mirror of
https://github.com/stablyai/orca.git
synced 2026-10-07 08:02:21 +00:00
perf: bound daemon checkpoint fanout (#4161)
perf: bound daemon checkpoint fanout
This commit is contained in:
@@ -7,6 +7,7 @@ import { DaemonPtyAdapter } from './daemon-pty-adapter'
|
||||
import { DaemonServer } from './daemon-server'
|
||||
import { getHistorySessionDirName } from './history-paths'
|
||||
import type { SubprocessHandle } from './session'
|
||||
import type { GetSnapshotResult, TerminalSnapshot } from './types'
|
||||
import type * as DaemonHealthModule from './daemon-health'
|
||||
|
||||
const { getMacDaemonSystemResolverHealthMock } = vi.hoisted(() => ({
|
||||
@@ -55,6 +56,27 @@ function createMockSubprocess(): SubprocessHandle & {
|
||||
}
|
||||
}
|
||||
|
||||
function createTestSnapshot(label: string): TerminalSnapshot {
|
||||
return {
|
||||
snapshotAnsi: label,
|
||||
scrollbackAnsi: '',
|
||||
rehydrateSequences: '',
|
||||
cwd: null,
|
||||
modes: {
|
||||
bracketedPaste: false,
|
||||
mouseTracking: false,
|
||||
mouseTrackingMode: 'none',
|
||||
sgrMouseMode: false,
|
||||
sgrMousePixelsMode: false,
|
||||
applicationCursor: false,
|
||||
alternateScreen: false
|
||||
},
|
||||
cols: 80,
|
||||
rows: 24,
|
||||
scrollbackLines: 0
|
||||
}
|
||||
}
|
||||
|
||||
async function waitFor(predicate: () => boolean, timeoutMs = 2000): Promise<void> {
|
||||
const start = Date.now()
|
||||
while (!predicate()) {
|
||||
@@ -539,6 +561,57 @@ describe('DaemonPtyAdapter (IPtyProvider)', () => {
|
||||
}
|
||||
})
|
||||
|
||||
it('limits concurrent checkpoint snapshot and disk work', async () => {
|
||||
historyAdapter = new DaemonPtyAdapter({ socketPath, tokenPath, historyPath: historyDir })
|
||||
const releaseSnapshotRequests: (() => void)[] = []
|
||||
const requestedSessionIds: string[] = []
|
||||
let inFlight = 0
|
||||
let maxInFlight = 0
|
||||
const request = vi.fn(
|
||||
async (_type: string, payload: { sessionId: string }): Promise<GetSnapshotResult> => {
|
||||
requestedSessionIds.push(payload.sessionId)
|
||||
inFlight++
|
||||
maxInFlight = Math.max(maxInFlight, inFlight)
|
||||
await new Promise<void>((resolve) => {
|
||||
releaseSnapshotRequests.push(() => {
|
||||
inFlight--
|
||||
resolve()
|
||||
})
|
||||
})
|
||||
return { snapshot: createTestSnapshot(payload.sessionId) }
|
||||
}
|
||||
)
|
||||
const checkpoint = vi.fn(async () => {})
|
||||
const dispose = vi.fn(async () => {})
|
||||
const disconnect = vi.fn()
|
||||
const internals = historyAdapter as unknown as {
|
||||
client: { request: typeof request; disconnect: typeof disconnect }
|
||||
historyManager: { checkpoint: typeof checkpoint; dispose: typeof dispose }
|
||||
checkpointSessions(sessionIds: Iterable<string>): Promise<Set<string>>
|
||||
}
|
||||
internals.client = { request, disconnect }
|
||||
internals.historyManager = { checkpoint, dispose }
|
||||
|
||||
const checkpointing = internals.checkpointSessions(['a', 'b', 'c', 'd', 'e', 'f'])
|
||||
await waitFor(() => requestedSessionIds.length === 4)
|
||||
|
||||
expect(maxInFlight).toBe(4)
|
||||
expect(requestedSessionIds).toEqual(['a', 'b', 'c', 'd'])
|
||||
|
||||
for (const release of releaseSnapshotRequests.splice(0)) {
|
||||
release()
|
||||
}
|
||||
await waitFor(() => requestedSessionIds.length === 6)
|
||||
|
||||
expect(maxInFlight).toBe(4)
|
||||
|
||||
for (const release of releaseSnapshotRequests.splice(0)) {
|
||||
release()
|
||||
}
|
||||
await expect(checkpointing).resolves.toEqual(new Set(['a', 'b', 'c', 'd', 'e', 'f']))
|
||||
expect(checkpoint).toHaveBeenCalledTimes(6)
|
||||
})
|
||||
|
||||
it('does not schedule a checkpoint timer until a session is dirty', async () => {
|
||||
const adapterClass = DaemonPtyAdapter as unknown as { CHECKPOINT_INTERVAL_MS: number }
|
||||
const previousInterval = adapterClass.CHECKPOINT_INTERVAL_MS
|
||||
|
||||
@@ -32,6 +32,7 @@ export type DaemonPtyAdapterOptions = {
|
||||
}
|
||||
|
||||
const MAX_TOMBSTONES = 1000
|
||||
const MAX_CONCURRENT_CHECKPOINTS = 4
|
||||
|
||||
export class TerminalKilledError extends Error {
|
||||
constructor(sessionId: string) {
|
||||
@@ -625,10 +626,18 @@ export class DaemonPtyAdapter implements IPtyProvider {
|
||||
if (!this.historyManager) {
|
||||
return completed
|
||||
}
|
||||
const promises: Promise<void>[] = []
|
||||
for (const sessionId of sessionIds) {
|
||||
promises.push(
|
||||
this.client
|
||||
const ids = Array.from(sessionIds)
|
||||
let nextIndex = 0
|
||||
|
||||
const checkpointNext = async (): Promise<void> => {
|
||||
for (;;) {
|
||||
const index = nextIndex
|
||||
nextIndex++
|
||||
if (index >= ids.length) {
|
||||
return
|
||||
}
|
||||
const sessionId = ids[index]
|
||||
await this.client
|
||||
.request<GetSnapshotResult>('getSnapshot', { sessionId })
|
||||
.then((result) => {
|
||||
if (result.snapshot && this.historyManager) {
|
||||
@@ -640,9 +649,15 @@ export class DaemonPtyAdapter implements IPtyProvider {
|
||||
return undefined
|
||||
})
|
||||
.catch((err) => console.warn('[history] checkpoint failed:', sessionId, err))
|
||||
)
|
||||
}
|
||||
}
|
||||
await Promise.all(promises)
|
||||
// Why: snapshot serialization and checkpoint writes are CPU/disk heavy.
|
||||
// Dirty-session filtering keeps idle terminals out; this cap prevents one
|
||||
// tick from snapshotting every active dirty terminal at once.
|
||||
const workers = Array.from({ length: Math.min(MAX_CONCURRENT_CHECKPOINTS, ids.length) }, () =>
|
||||
checkpointNext()
|
||||
)
|
||||
await Promise.all(workers)
|
||||
return completed
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user