diff --git a/packages/shared/src/kernel/run-manager.test.ts b/packages/shared/src/kernel/run-manager.test.ts index c73cdfe..58e5e4d 100644 --- a/packages/shared/src/kernel/run-manager.test.ts +++ b/packages/shared/src/kernel/run-manager.test.ts @@ -1,7 +1,7 @@ import { mkdtemp, rm } from 'node:fs/promises'; import { tmpdir } from 'node:os'; import { join } from 'node:path'; -import { afterEach, beforeEach, describe, expect, it } from 'bun:test'; +import { afterEach, beforeEach, describe, expect, it, spyOn } from 'bun:test'; import type { AgentEvent, AgentEventPayload, @@ -153,6 +153,67 @@ function completedScript(answer: string): (input: AgentRunInput) => AsyncIterabl } describe('RunManager', () => { + for (const sameSession of [true, false]) { + it(`rejects concurrent starts during session lookup (${sameSession ? 'same' : 'different'} session)`, async () => { + const setup = Promise.withResolvers(); + const runtimeGate = Promise.withResolvers(); + const { sessions, runs, runtime } = makeKernel(completedScript('Done'), {}, runtimeGate.promise); + const first = await sessions.createSession('First'); + const second = sameSession ? first : await sessions.createSession('Second'); + const getSession = sessions.getSession.bind(sessions); + const lookup = spyOn(sessions, 'getSession').mockImplementation(async (id) => { + await setup.promise; + return getSession(id); + }); + const pending = Promise.allSettled([ + runs.startRun(first.id, 'first question'), + runs.startRun(second.id, 'second question'), + ]); + setup.resolve(); + const results = await pending; + lookup.mockRestore(); + runtimeGate.resolve(); + // Drain even the buggy two-run baseline before asserting or removing storage. + await waitFor(async () => { + const finished = await Promise.all(results.map(async (result) => + result.status === 'rejected' || + (await sessions.getRun(result.value.sessionId, result.value.id))?.status === 'completed' + )); + return finished.every(Boolean) && !runs.isRunning(); + }); + expect(results[0].status).toBe('fulfilled'); + expect(results[1]).toMatchObject({ status: 'rejected', reason: { code: 'RUN_IN_PROGRESS' } }); + expect(runtime.ensureSessionCalls).toHaveLength(1); + expect((await sessions.listMessages(first.id)).map((message) => message.content)).toEqual(['first question', 'Done']); + if (!sameSession) expect(await sessions.listMessages(second.id)).toEqual([]); + }); + } + + it('reports a starting session as busy and releases it after a lookup failure', async () => { + const setup = Promise.withResolvers(); + const { sessions, runs } = makeKernel(completedScript('Recovered')); + const session = await sessions.createSession('A'); + const lookup = spyOn(sessions, 'getSession').mockImplementationOnce(async () => { + await setup.promise; + throw new Error('storage unavailable'); + }); + const pending = runs.startRun(session.id, 'failed attempt'); + const busy = runs.isRunning(); + const sessionBusy = runs.hasActiveRun(session.id); + const otherBusy = runs.hasActiveRun('other'); + setup.resolve(); + await expect(pending).rejects.toThrow('storage unavailable'); + lookup.mockRestore(); + expect(runs.isRunning()).toBe(false); + expect(runs.hasActiveRun(session.id)).toBe(false); + await runs.startRun(session.id, 'retry'); + await waitFor(async () => !runs.isRunning()); + expect((await sessions.listMessages(session.id)).map((message) => message.content)).toEqual(['retry', 'Recovered']); + expect(busy).toBe(true); + expect(sessionBusy).toBe(true); + expect(otherBusy).toBe(false); + }); + it('runs a full loop: run_started → tools → deltas → run_completed, persisted', async () => { const { sessions, runs, runtime, store } = makeKernel(completedScript('Portfolio risk is moderate.')); const session = await sessions.createSession('Portfolio Review'); diff --git a/packages/shared/src/kernel/run-manager.ts b/packages/shared/src/kernel/run-manager.ts index 0ac0b0d..c07e4fc 100644 --- a/packages/shared/src/kernel/run-manager.ts +++ b/packages/shared/src/kernel/run-manager.ts @@ -126,6 +126,7 @@ export class RunManager { private readonly budgetInput: ResolveBudgetInput; private readonly searchToolPatterns: readonly string[]; private readonly runawayPolicy: Partial; + private startingSessionId: string | null = null; private readonly listeners = new Set<(event: AgentEvent) => void>(); private readonly streamListeners = new Set<(sessionId: string, event: StreamEvent) => void>(); private activeRun: ActiveRun | null = null; @@ -175,12 +176,12 @@ export class RunManager { /** Whether a run is currently executing (Pi runtime executes one at a time). */ isRunning(): boolean { - return this.activeRun !== null; + return this.startingSessionId !== null || this.activeRun !== null; } /** Whether a run is currently executing for the given session. */ hasActiveRun(sessionId: string): boolean { - return this.activeRun?.sessionId === sessionId; + return this.startingSessionId === sessionId || this.activeRun?.sessionId === sessionId; } /** @@ -203,13 +204,30 @@ export class RunManager { if (!text) { throw createCodeError('INVALID_ARGUMENT', 'Message content is required.'); } - if (this.activeRun) { + if (this.isRunning()) { throw createCodeError( 'RUN_IN_PROGRESS', 'Another run is still in progress. Stop it before sending a new message.' ); } + // Reserve the single runtime before any asynchronous lookup or persistence. + // A failed start must release the reservation so the caller can retry. + this.startingSessionId = sessionId; + try { + return await this.prepareRun(sessionId, text, workspaceContext, locale, budgetOverrides); + } finally { + this.startingSessionId = null; + } + } + + private async prepareRun( + sessionId: string, + text: string, + workspaceContext?: WorkspaceContext, + locale?: SupportedLocale, + budgetOverrides?: RunBudgetLimits + ): Promise { const session = await this.sessions.getSession(sessionId); if (!session) { throw createCodeError('SESSION_NOT_FOUND', `Session ${sessionId} was not found.`);