diff --git a/src/renderer/src/components/terminal-pane/agent-completion-process-monitor.ts b/src/renderer/src/components/terminal-pane/agent-completion-process-monitor.ts index 78c859e4b4b..f34ceffa97e 100644 --- a/src/renderer/src/components/terminal-pane/agent-completion-process-monitor.ts +++ b/src/renderer/src/components/terminal-pane/agent-completion-process-monitor.ts @@ -101,6 +101,9 @@ export function createAgentCompletionProcessMonitor({ enqueueAgentProcessInspection({ priority, canRun: () => !state.disposed, + // Local reads all resolve out of one process-table capture; remote ones each cost their + // own execution-host round trip and stay admitted one at a time. + sharesHostObservation: options.isRemotePtyId?.(ptyId) !== true, run: async () => { let inspectedRecognizedAgent = false let inspectionSucceeded = false diff --git a/src/renderer/src/components/terminal-pane/agent-process-inspection-queue.ts b/src/renderer/src/components/terminal-pane/agent-process-inspection-queue.ts index 6dc1dc4ccf2..8e5aef6cade 100644 --- a/src/renderer/src/components/terminal-pane/agent-process-inspection-queue.ts +++ b/src/renderer/src/components/terminal-pane/agent-process-inspection-queue.ts @@ -4,29 +4,56 @@ type InspectionTask = { priority: InspectionPriority canRun: () => boolean run: () => Promise + /** + * Reads served by one shared host observation. Every local pane's inspection resolves out of + * the same TTL-and-in-flight-deduped process-table snapshot, so a whole round of them costs + * the host one capture however many panes ride it. + */ + sharesHostObservation?: boolean } const MAX_CONCURRENT_INSPECTIONS = 4 const MAX_INSPECTION_STARTS_PER_SECOND = 8 let activeInspections = 0 +let inspectionPumpQueued = false let inspectionPumpTimer: ReturnType | null = null const inspectionStarts: number[] = [] const inspectionQueue: InspectionTask[] = [] -function canStartInspection(now: number): boolean { +/** + * Host observations still admissible right now. A start is one observation, not one pane: a + * shared-observation round costs one however many panes ride it, an unshared task costs one each. + */ +function availableInspectionStarts(now: number): number { if (inspectionStarts.length > 0 && now < inspectionStarts[0]!) { inspectionStarts.length = 0 } while (inspectionStarts.length > 0 && now - inspectionStarts[0]! >= 1_000) { inspectionStarts.shift() } - return ( - activeInspections < MAX_CONCURRENT_INSPECTIONS && - inspectionStarts.length < MAX_INSPECTION_STARTS_PER_SECOND + return Math.min( + MAX_CONCURRENT_INSPECTIONS - activeInspections, + MAX_INSPECTION_STARTS_PER_SECOND - inspectionStarts.length ) } +/** + * Pump on a microtask, so a synchronous burst of enqueues forms one round. Pumping inline + * spent a start per pane until the concurrency slots filled and then parked the rest of the + * burst on the 100ms retry. + */ +function queueInspectionPump(): void { + if (inspectionPumpQueued) { + return + } + inspectionPumpQueued = true + queueMicrotask(() => { + inspectionPumpQueued = false + pumpInspectionQueue() + }) +} + function scheduleInspectionPump(delayMs = 0): void { if (inspectionPumpTimer !== null) { return @@ -37,44 +64,92 @@ function scheduleInspectionPump(delayMs = 0): void { }, delayMs) } -function pumpInspectionQueue(): void { - // Drop disposed tasks before slot/rate accounting. - for (let index = inspectionQueue.length - 1; index >= 0; index -= 1) { - const task = inspectionQueue[index] - if (task && !task.canRun()) { - inspectionQueue.splice(index, 1) +/** Compact disposed tasks out in one pass; a splice per drop is quadratic at pane scale. */ +function dropDisposedInspections(): void { + let write = 0 + for (let read = 0; read < inspectionQueue.length; read += 1) { + const task = inspectionQueue[read]! + if (task.canRun()) { + inspectionQueue[write] = task + write += 1 } } + inspectionQueue.length = write +} + +function startInspectionRound(tasks: InspectionTask[], now: number): void { + activeInspections += 1 + inspectionStarts.push(now) + let outstanding = tasks.length + const settleOne = (): void => { + outstanding -= 1 + if (outstanding > 0) { + return + } + activeInspections = Math.max(0, activeInspections - 1) + if (inspectionQueue.length > 0) { + scheduleInspectionPump() + } + } + for (const task of tasks) { + // Started synchronously so every read in the round lands in the same tick, hitting one + // process-table capture instead of serializing one capture window apart. + // Why the catch before finally: an unreachable runtime rejects the inspection on a cadence, and a + // bare `.finally()` chain re-raises it as a renderer-global unhandledrejection. Coordinators own + // their own failure/backoff state, so the queue only has to keep its accounting running. + void task + .run() + .catch(() => {}) + .finally(settleOne) + } +} + +/** Take every shared-observation task, in order, leaving the rest queued. */ +function takeSharedObservationRound(): InspectionTask[] { + const round: InspectionTask[] = [] + let write = 0 + for (let read = 0; read < inspectionQueue.length; read += 1) { + const task = inspectionQueue[read]! + if (task.sharesHostObservation === true) { + round.push(task) + } else { + inspectionQueue[write] = task + write += 1 + } + } + inspectionQueue.length = write + return round +} + +function pumpInspectionQueue(): void { + // Drop disposed tasks before slot/rate accounting. + dropDisposedInspections() if (inspectionQueue.length === 0) { return } const now = Date.now() - if (!canStartInspection(now)) { + let starts = availableInspectionStarts(now) + if (starts <= 0) { scheduleInspectionPump(100) return } - - const priorityIndex = inspectionQueue.findIndex((task) => task.priority === 'pending-title') - const next = - priorityIndex !== -1 ? inspectionQueue.splice(priorityIndex, 1)[0] : inspectionQueue.shift() - if (!next) { - return + // The whole shared-observation backlog goes on one start, so a pane's wait is bounded by the + // observation budget rather than by how many other panes are also due. + const sharedRound = takeSharedObservationRound() + if (sharedRound.length > 0) { + startInspectionRound(sharedRound, now) + starts -= 1 + } + while (starts > 0 && inspectionQueue.length > 0) { + const priorityIndex = inspectionQueue.findIndex((task) => task.priority === 'pending-title') + const next = + priorityIndex !== -1 ? inspectionQueue.splice(priorityIndex, 1)[0] : inspectionQueue.shift() + if (!next) { + break + } + startInspectionRound([next], now) + starts -= 1 } - - activeInspections += 1 - inspectionStarts.push(now) - // Why the catch before finally: an unreachable runtime rejects the inspection on a cadence, and a - // bare `.finally()` chain re-raises it as a renderer-global unhandledrejection. Coordinators own - // their own failure/backoff state, so the queue only has to keep its accounting running. - void next - .run() - .catch(() => {}) - .finally(() => { - activeInspections = Math.max(0, activeInspections - 1) - if (inspectionQueue.length > 0) { - scheduleInspectionPump() - } - }) if (inspectionQueue.length > 0) { scheduleInspectionPump() @@ -83,7 +158,7 @@ function pumpInspectionQueue(): void { export function enqueueAgentProcessInspection(task: InspectionTask): void { inspectionQueue.push(task) - pumpInspectionQueue() + queueInspectionPump() } export function resetAgentProcessInspectionQueueForTests(): void { @@ -91,6 +166,7 @@ export function resetAgentProcessInspectionQueueForTests(): void { clearTimeout(inspectionPumpTimer) inspectionPumpTimer = null } + inspectionPumpQueued = false activeInspections = 0 inspectionStarts.length = 0 inspectionQueue.length = 0 diff --git a/src/renderer/src/components/terminal-pane/agent-process-inspection-round.test.ts b/src/renderer/src/components/terminal-pane/agent-process-inspection-round.test.ts new file mode 100644 index 00000000000..216bf304a1c --- /dev/null +++ b/src/renderer/src/components/terminal-pane/agent-process-inspection-round.test.ts @@ -0,0 +1,149 @@ +// Regression guard for the inspection admission budget. The cadence tiers +// (active 750ms / idle 2000 / hidden 3000 / no-evidence 15000) were a per-pane +// promise the queue could not keep: the budget of 8 starts per second was spent +// one pane at a time, so N due panes meant roughly N/8 seconds between +// inspections for each of them and agent-completion latency degraded as the +// user added panes. Local inspections all resolve out of one TTL-and-in-flight- +// deduped process-table capture, so a whole round of them is one host +// observation and rides one start, launched in a single tick. The budget itself +// is unchanged — it just buys the whole round instead of one pane. +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' +import { + enqueueAgentProcessInspection, + resetAgentProcessInspectionQueueForTests +} from './agent-process-inspection-queue' + +const PANES = 300 + +beforeEach(() => { + vi.useFakeTimers() +}) + +afterEach(() => { + vi.useRealTimers() + resetAgentProcessInspectionQueueForTests() +}) + +describe('agent process inspection rounds', () => { + it('inspects every pane of a 300-pane round inside the same admission budget', async () => { + const inspected = new Set() + + for (let index = 0; index < PANES; index += 1) { + const ptyId = `pty-${index}` + enqueueAgentProcessInspection({ + priority: 'cadence', + canRun: () => true, + sharesHostObservation: true, + run: async () => { + await Promise.resolve() + inspected.add(ptyId) + } + }) + } + // Well inside one 1s rate-limiter window: pre-fix only the 8 starts that window + // allows are spent, so only 8 of the 300 panes are ever inspected. + await vi.advanceTimersByTimeAsync(200) + + expect(inspected.size).toBe(PANES) + }) + + it('launches the whole round in one tick on one start', async () => { + const launchesPerTick = new Map() + let unshared = 0 + + for (let index = 0; index < PANES; index += 1) { + enqueueAgentProcessInspection({ + priority: 'cadence', + canRun: () => true, + sharesHostObservation: true, + run: async () => { + // Fake timers freeze the clock inside a tick, so a shared timestamp is a shared burst. + const tick = Date.now() + launchesPerTick.set(tick, (launchesPerTick.get(tick) ?? 0) + 1) + } + }) + } + // Seven unshared panes still fit, which is what proves the round cost exactly one of + // the eight starts rather than one per pane until the concurrency slots filled. + for (let index = 0; index < 7; index += 1) { + enqueueAgentProcessInspection({ + priority: 'cadence', + canRun: () => true, + sharesHostObservation: false, + run: async () => { + unshared += 1 + } + }) + } + await vi.advanceTimersByTimeAsync(200) + + // One synchronous burst carries every pane, so they hit one process-table capture + // rather than serializing across the limiter. + expect([...launchesPerTick.values()]).toEqual([PANES]) + expect(unshared).toBe(7) + }) + + it('keeps a pane whose read is not shared admitted one round trip at a time', async () => { + const started: string[] = [] + + for (let index = 0; index < PANES; index += 1) { + enqueueAgentProcessInspection({ + priority: 'cadence', + canRun: () => true, + // Remote panes: each costs its own execution-host round trip, so no round shares them. + sharesHostObservation: false, + run: async () => { + started.push(`ssh-${index}`) + } + }) + } + await vi.advanceTimersByTimeAsync(200) + + expect(started.length).toBeLessThanOrEqual(8) + }) + + it('still serves a pending-title read ahead of the queued cadence backlog', async () => { + const order: string[] = [] + + for (let index = 0; index < 20; index += 1) { + enqueueAgentProcessInspection({ + priority: 'cadence', + canRun: () => true, + sharesHostObservation: false, + run: async () => { + order.push(`cadence-${index}`) + } + }) + } + enqueueAgentProcessInspection({ + priority: 'pending-title', + canRun: () => true, + sharesHostObservation: false, + run: async () => { + order.push('pending-title') + } + }) + await vi.advanceTimersByTimeAsync(200) + + expect(order.length).toBeLessThanOrEqual(8) + expect(order).toContain('pending-title') + }) + + it('drops disposed panes out of the round instead of inspecting them', async () => { + const inspected: number[] = [] + + for (let index = 0; index < PANES; index += 1) { + enqueueAgentProcessInspection({ + priority: 'cadence', + canRun: () => index % 2 === 0, + sharesHostObservation: true, + run: async () => { + inspected.push(index) + } + }) + } + await vi.advanceTimersByTimeAsync(200) + + expect(inspected).toEqual(Array.from({ length: PANES / 2 }, (_unused, index) => index * 2)) + }) +})