diff --git a/packages/shared/src/kernel/run-manager.test.ts b/packages/shared/src/kernel/run-manager.test.ts index 8f41710..7425533 100644 --- a/packages/shared/src/kernel/run-manager.test.ts +++ b/packages/shared/src/kernel/run-manager.test.ts @@ -21,6 +21,30 @@ import { RunManager, type RunManagerOptions } from './run-manager.ts'; let dir = ''; let clock = 1000; +interface ManualTimerEntry { + callback: () => void; + cancelled: boolean; + delayMs: number; +} + +class ManualScheduler { + readonly entries: ManualTimerEntry[] = []; + + setTimeout(callback: () => void, delayMs: number) { + const entry: ManualTimerEntry = { callback, cancelled: false, delayMs }; + this.entries.push(entry); + return { + cancel: () => { + entry.cancelled = true; + }, + }; + } + + fire(index: number): void { + this.entries[index]?.callback(); + } +} + beforeEach(async () => { dir = await mkdtemp(join(tmpdir(), 'finagent-kernel-')); clock = 1000; @@ -32,9 +56,10 @@ afterEach(async () => { function makeKernel( script: (input: AgentRunInput) => AsyncIterable, - extra: Partial> = {} + extra: Partial> = {}, + ensureSessionGate?: Promise ) { - const runtime = new ScriptedRuntime(script); + const runtime = new ScriptedRuntime(script, ensureSessionGate); const store = new JsonFileStore(dir); const sessions = new SessionManager({ sessions: new SessionRepository(store), @@ -57,7 +82,10 @@ class ScriptedRuntime implements AgentRuntime { ensureSessionCalls: Array<{ id: string; sessionPath?: string }> = []; cancelCalls: Array<{ sessionId: string; runId: string }> = []; - constructor(private readonly script: (input: AgentRunInput) => AsyncIterable) {} + constructor( + private readonly script: (input: AgentRunInput) => AsyncIterable, + private readonly ensureSessionGate?: Promise + ) {} async getTools(): Promise> { return { ok: true, data: [] }; @@ -65,6 +93,7 @@ class ScriptedRuntime implements AgentRuntime { async ensureSession(session: { id: string; title?: string; sessionPath?: string }): Promise { this.ensureSessionCalls.push(session); + await this.ensureSessionGate; return { sessionId: session.id, status: 'active' }; } @@ -445,6 +474,177 @@ describe('RunManager budgets and runaway detection (#17)', () => { }); }); + it('stops a stream that stalls before its first event', async () => { + let release!: () => void; + const stalled = new Promise((resolve) => { + release = resolve; + }); + const scheduler = new ManualScheduler(); + const { sessions, runs, runtime } = makeKernel( + async function* () { + await stalled; + }, + { budgets: { defaults: { wallClockMs: 1_000 } }, scheduler } + ); + const session = await sessions.createSession('Silent stream'); + + const run = await runs.startRun(session.id, 'wait forever'); + expect(scheduler.entries).toHaveLength(1); + expect(scheduler.entries[0]?.delayMs).toBe(1_000); + + clock += 1_000; + scheduler.fire(0); + await waitFor(async () => runtime.cancelCalls.length === 1); + + release(); + await waitFor(async () => !runs.isRunning()); + + expect(await sessions.getRun(session.id, run.id)).toMatchObject({ + status: 'cancelled', + stopReason: 'budget_exhausted', + stopDetail: { key: 'wallClockMs', limit: 1_000, used: 1_000 }, + }); + }); + + it('keeps partial output when the stream stalls after a message', async () => { + let release!: () => void; + const stalled = new Promise((resolve) => { + release = resolve; + }); + const scheduler = new ManualScheduler(); + const { sessions, runs, runtime } = makeKernel( + async function* (input) { + yield event(input.sessionId, input.runId, 'message_completed', { answer: 'partial' }, 1); + await stalled; + }, + { budgets: { defaults: { wallClockMs: 1_000 } }, scheduler } + ); + const session = await sessions.createSession('Partial'); + let sawPartial = false; + runs.subscribe((agentEvent) => { + if (agentEvent.type === 'message_completed') sawPartial = true; + }); + + const run = await runs.startRun(session.id, 'slow task'); + await waitFor(async () => sawPartial); + + clock += 1_000; + scheduler.fire(0); + await waitFor(async () => runtime.cancelCalls.length === 1); + + release(); + await waitFor(async () => !runs.isRunning()); + + expect(await sessions.getRun(session.id, run.id)).toMatchObject({ + status: 'cancelled', + answer: 'partial', + stopReason: 'budget_exhausted', + stopDetail: { key: 'wallClockMs', limit: 1_000, used: 1_000 }, + }); + const messages = await sessions.listMessages(session.id); + expect(messages[1]).toMatchObject({ role: 'assistant', content: 'partial' }); + }); + + it('does not start the runtime when the deadline expires during session setup', async () => { + let releaseEnsure!: () => void; + const ensureSessionGate = new Promise((resolve) => { + releaseEnsure = resolve; + }); + let runStarted = false; + const scheduler = new ManualScheduler(); + const { sessions, runs, runtime } = makeKernel( + async function* () { + runStarted = true; + }, + { budgets: { defaults: { wallClockMs: 1_000 } }, scheduler }, + ensureSessionGate + ); + const session = await sessions.createSession('Setup race'); + + const run = await runs.startRun(session.id, 'slow setup'); + expect(runtime.ensureSessionCalls).toHaveLength(1); + + clock += 1_000; + scheduler.fire(0); + await waitFor(async () => runtime.cancelCalls.length === 1); + + releaseEnsure(); + await waitFor(async () => !runs.isRunning()); + + expect(runStarted).toBe(false); + expect(await sessions.getRun(session.id, run.id)).toMatchObject({ + status: 'cancelled', + stopReason: 'budget_exhausted', + stopDetail: { key: 'wallClockMs', limit: 1_000, used: 1_000 }, + }); + }); + + it('cancels a wall-clock expiry only once even if the timer callback repeats', async () => { + let release!: () => void; + const stalled = new Promise((resolve) => { + release = resolve; + }); + const scheduler = new ManualScheduler(); + const { sessions, runs, runtime } = makeKernel( + async function* () { + await stalled; + }, + { budgets: { defaults: { wallClockMs: 1_000 } }, scheduler } + ); + const session = await sessions.createSession('Once'); + const run = await runs.startRun(session.id, 'loop guard'); + + clock += 1_000; + scheduler.fire(0); + scheduler.fire(0); + await waitFor(async () => runtime.cancelCalls.length === 1); + + release(); + await waitFor(async () => !runs.isRunning()); + + expect(runtime.cancelCalls).toEqual([{ sessionId: session.id, runId: run.id }]); + }); + + it('clears a completed run deadline so a stale callback cannot cancel a later run', async () => { + let releaseLater!: () => void; + const laterStalled = new Promise((resolve) => { + releaseLater = resolve; + }); + const scheduler = new ManualScheduler(); + const { sessions, runs, runtime } = makeKernel( + async function* (input) { + if (input.content === 'first run') { + yield event(input.sessionId, input.runId, 'run_completed', { + answer: 'done', + toolCalls: [], + }); + return; + } + await laterStalled; + }, + { budgets: { defaults: { wallClockMs: 1_000 } }, scheduler } + ); + const session = await sessions.createSession('Isolation'); + + const first = await runs.startRun(session.id, 'first run'); + await waitFor(async () => !runs.isRunning()); + expect(await sessions.getRun(session.id, first.id)).toMatchObject({ status: 'completed' }); + expect(scheduler.entries[0]?.cancelled).toBe(true); + + const later = await runs.startRun(session.id, 'second run'); + expect(scheduler.entries).toHaveLength(2); + + // Simulate an already-dispatched stale callback from the completed run. + scheduler.fire(0); + await new Promise((resolve) => setTimeout(resolve, 0)); + expect(runtime.cancelCalls).toEqual([]); + expect(runs.hasActiveRun(session.id)).toBe(true); + + releaseLater(); + await waitFor(async () => !runs.isRunning()); + expect(await sessions.getRun(session.id, later.id)).toMatchObject({ status: 'completed' }); + }); + it('leaves a run untouched when neither budgets nor detectors are configured', async () => { const { sessions, runs, runtime } = makeKernel(completedScript('Answer')); const session = await sessions.createSession('Plain'); diff --git a/packages/shared/src/kernel/run-manager.ts b/packages/shared/src/kernel/run-manager.ts index 00cc0c9..fc6e191 100644 --- a/packages/shared/src/kernel/run-manager.ts +++ b/packages/shared/src/kernel/run-manager.ts @@ -37,11 +37,28 @@ import { } from './runaway-detector.ts'; import { buildFinancialEvidence } from '../evidence/financial-evidence.ts'; +export interface RunManagerTimer { + cancel(): void; +} + +export interface RunManagerScheduler { + setTimeout(callback: () => void, delayMs: number): RunManagerTimer; +} + +const DEFAULT_SCHEDULER: RunManagerScheduler = { + setTimeout(callback, delayMs) { + const handle = setTimeout(callback, delayMs); + return { cancel: () => clearTimeout(handle) }; + }, +}; + export interface RunManagerOptions { sessions: SessionManager; runs: RunRepository; runtime: AgentRuntime; now?: () => number; + /** Timer implementation for deterministic tests; production uses setTimeout. */ + scheduler?: RunManagerScheduler; /** * Budget defaults and the system ceiling applied to every run (#17). A run * may override the defaults but never the ceiling; with no `budgets` option a @@ -66,6 +83,8 @@ interface ActiveRun { runaway: RunawayState; /** Set when a budget or a runaway detector stopped the run. */ stop?: RunStop; + /** Active wall-clock deadline; absent when the limit is unset or already cleared. */ + wallClockTimer?: RunManagerTimer; } /** @@ -82,6 +101,7 @@ export class RunManager { private readonly runs: RunRepository; private readonly runtime: AgentRuntime; private readonly now: () => number; + private readonly scheduler: RunManagerScheduler; private readonly budgetInput: ResolveBudgetInput; private readonly searchToolPatterns: readonly string[]; private readonly runawayPolicy: Partial; @@ -93,6 +113,7 @@ export class RunManager { this.runs = options.runs; this.runtime = options.runtime; this.now = options.now ?? Date.now; + this.scheduler = options.scheduler ?? DEFAULT_SCHEDULER; this.budgetInput = options.budgets ?? {}; this.searchToolPatterns = options.searchTools ?? []; this.runawayPolicy = options.runaway ?? {}; @@ -169,7 +190,7 @@ export class RunManager { ceiling: this.budgetInput.ceiling, overrides: budgetOverrides, }); - this.activeRun = { + const active: ActiveRun = { sessionId, runId: run.id, cancelRequested: false, @@ -178,6 +199,8 @@ export class RunManager { usage: createUsage(), runaway: createRunawayState(), }; + this.activeRun = active; + this.armWallClockBudget(active); this.emit({ id: randomUUID(), sessionId, @@ -199,6 +222,7 @@ export class RunManager { return; } active.cancelRequested = true; + this.clearWallClockTimer(active); await this.runtime.cancel({ sessionId, runId }); } @@ -221,29 +245,33 @@ export class RunManager { recentSymbols: session.recentSymbols, }); - for await (const event of this.runtime.run({ - sessionId: run.sessionId, - runId: run.id, - content: run.input, - workspaceContext, - locale, - })) { - this.emit(event); - if (event.type === 'message_delta' || event.type === 'message_completed') { - answer = event.payload.answer; - } else if (event.type === 'tool_completed') { - toolCalls.push(event.payload.toolCall); - } else if (event.type === 'run_failed') { - failure = event.payload.error; - sawTerminal = true; - } else if (event.type === 'run_completed') { - answer = event.payload.answer; - sawTerminal = true; + // A wall-clock deadline can expire while ensureSession is still pending. + // Never start a runtime run that was already cancelled before it began. + if (!this.activeRun?.cancelRequested) { + for await (const event of this.runtime.run({ + sessionId: run.sessionId, + runId: run.id, + content: run.input, + workspaceContext, + locale, + })) { + this.emit(event); + if (event.type === 'message_delta' || event.type === 'message_completed') { + answer = event.payload.answer; + } else if (event.type === 'tool_completed') { + toolCalls.push(event.payload.toolCall); + } else if (event.type === 'run_failed') { + failure = event.payload.error; + sawTerminal = true; + } else if (event.type === 'run_completed') { + answer = event.payload.answer; + sawTerminal = true; + } + + // Budgets and detectors are evaluated after the event is accounted for, + // so a run that stops keeps the evidence it had already produced. + if (await this.applyBudget(event)) break; } - - // Budgets and detectors are evaluated after the event is accounted for, - // so a run that stops keeps the evidence it had already produced. - if (await this.applyBudget(event)) break; } } catch (error) { failure = toApiError(error); @@ -300,6 +328,8 @@ export class RunManager { recentSymbols: collectSymbols(toolCalls), }); + if (active) this.clearWallClockTimer(active); + // The run is fully settled (persisted) only now; only then allow the next run. this.activeRun = null; @@ -365,13 +395,52 @@ export class RunManager { } if (stop === undefined) return false; + return this.requestSafeguardStop(stop); + } + + /** Schedule the wall-clock safeguard for one active run. */ + private armWallClockBudget(active: ActiveRun): void { + const limit = active.limits.wallClockMs; + if (limit === undefined) return; + + active.wallClockTimer = this.scheduler.setTimeout(() => { + void this.handleWallClockExpiry(active.runId).catch(() => undefined); + }, limit); + } + + /** Stop once from a timer callback, even if the callback is accidentally fired twice. */ + private async handleWallClockExpiry(runId: string): Promise { + const active = this.activeRun; + if (!active || active.runId !== runId) return; + + const limit = active.limits.wallClockMs; + if (limit === undefined) return; + + const used = Math.max(this.now() - active.startedAt, limit); + active.usage = { ...active.usage, wallClockMs: used }; + await this.requestSafeguardStop(budgetStop({ key: 'wallClockMs', limit, used })); + } + + /** + * Shared idempotent stop path for event-driven budgets and timer deadlines. + * The first caller owns cancellation; later callers are no-ops. + */ + private async requestSafeguardStop(stop: RunStop): Promise { + const active = this.activeRun; + if (!active || active.stop !== undefined || active.cancelRequested) return false; active.stop = stop; active.cancelRequested = true; + this.clearWallClockTimer(active); await this.runtime.cancel({ sessionId: active.sessionId, runId: active.runId }); return true; } + private clearWallClockTimer(active: ActiveRun): void { + active.wallClockTimer?.cancel(); + active.wallClockTimer = undefined; + } + /** * The query a search-tool call carries, when this run tracks search loops. * @param toolCall - the completed tool call.