mirror of
https://github.com/stablyai/orca.git
synced 2026-09-29 16:02:50 +00:00
perf(terminals): spend one inspection start on a whole cadence round
The inspection rate limiter counted panes when it should have counted host observations. `MAX_INSPECTION_STARTS_PER_SECOND = 8` is global, and it was spent one pane at a time, so N due panes meant an effective per-pane period of max(tier, N/8 seconds) — ~37.5s at 300 panes for a pane the code polls at 750ms. Agent-completion latency degraded monotonically as panes were added. Every local pane's inspection resolves out of the same TTL-and-in-flight- deduped process-table capture, so a whole round of them is one host observation. The queue now drains all shared-observation tasks as one round on one start, launched in a single tick. Remote panes each cost their own execution-host round trip and stay admitted one at a time. Both the budget and the cadence tiers are numerically unchanged. Disposed tasks are also compacted out in one pass instead of a splice per drop, so the per-round predicate cost is linear rather than quadratic at pane scale. No IPC, preload, wire, or main-process change: each pane keeps its existing per-pane `pty:inspectProcess` invoke.
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -4,29 +4,56 @@ type InspectionTask = {
|
||||
priority: InspectionPriority
|
||||
canRun: () => boolean
|
||||
run: () => Promise<void>
|
||||
/**
|
||||
* 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<typeof setTimeout> | 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
|
||||
|
||||
@@ -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<string>()
|
||||
|
||||
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<number, number>()
|
||||
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))
|
||||
})
|
||||
})
|
||||
Reference in New Issue
Block a user