diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 0d887c98fa78..f278d5ff6b98 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -195,6 +195,7 @@ it.effect.each( ); const fallbackProviderInstanceId = ProviderInstanceId.make("claudeAgent"); const continuationSent = yield* Deferred.make(); + const continuationRunning = yield* Deferred.make(); const continuationCleared = yield* Deferred.make(); const sends: ProviderSendTurnInput[] = []; const dispatched: OrchestrationCommand[] = []; @@ -288,16 +289,20 @@ it.effect.each( }, dispatch: (command) => Effect.sync(() => dispatched.push(command)).pipe( + Effect.flatMap(() => { + const runningCount = dispatched.filter( + (entry) => + entry.type === "thread.session.set" && entry.session.status === "running", + ).length; + return runningCount === 2 + ? Deferred.succeed(continuationRunning, undefined) + : Effect.void; + }), Effect.as({ sequence: dispatched.length }), ), }); - assert.isTrue( - dispatched.every( - (command) => - command.type === "thread.session.set" && command.session.status === "starting", - ), - ); yield* Deferred.await(continuationSent); + yield* Deferred.await(continuationRunning); yield* Deferred.await(continuationCleared); assert.deepStrictEqual( @@ -313,28 +318,50 @@ it.effect.each( }, ], ); - assert.deepStrictEqual( - dispatched.map((command) => - command.type === "thread.session.set" - ? { + const sessionSets = dispatched.flatMap((command) => + command.type === "thread.session.set" + ? [ + { threadId: command.threadId, status: command.session.status, activeTurnId: command.session.activeTurnId, - } - : null, - ), + }, + ] + : [], + ); + const byThreadId = (left: { threadId: ThreadId }, right: { threadId: ThreadId }) => + String(left.threadId).localeCompare(String(right.threadId)); + assert.deepStrictEqual( + sessionSets + .filter((command) => command.status === "starting") + .slice() + .sort(byThreadId), + [ + { threadId: codex.id, status: "starting" as const, activeTurnId: null }, + { threadId: fallback.id, status: "starting" as const, activeTurnId: null }, + ] + .slice() + .sort(byThreadId), + ); + assert.deepStrictEqual( + sessionSets + .filter((command) => command.status === "running") + .slice() + .sort(byThreadId), [ { threadId: codex.id, - status: "starting", - activeTurnId: null, + status: "running" as const, + activeTurnId: TurnId.make(`continued-${String(codex.id)}`), }, { threadId: fallback.id, - status: "starting", - activeTurnId: null, + status: "running" as const, + activeTurnId: TurnId.make(`continued-${String(fallback.id)}`), }, - ], + ] + .slice() + .sort(byThreadId), ); for (const [thread, continuationTurnId] of [ [codex, codex.session.activeTurnId], @@ -909,6 +936,8 @@ for (const preparedStatus of [ assert.deepStrictEqual(sends, [ { threadId: thread.id, continuation: true, interactionMode: "default" }, ]); + assert.equal(thread.session.status, "running"); + assert.equal(thread.session.activeTurnId, TurnId.make("turn-recovered")); assert.deepStrictEqual(binding.runtimePayload, { activeTurnId: null, continueAfterServerUpdate: null, diff --git a/apps/server/src/serverRuntimeStartup.ts b/apps/server/src/serverRuntimeStartup.ts index 148d6aa75110..7d9dbcb1095a 100644 --- a/apps/server/src/serverRuntimeStartup.ts +++ b/apps/server/src/serverRuntimeStartup.ts @@ -705,13 +705,30 @@ export const reconcileProviderSessions = Effect.gen(function* () { }); } const capabilities = yield* providerService.getCapabilities(providerInstanceId); - yield* providerService.sendTurn({ + const result = yield* providerService.sendTurn({ threadId: thread.id, ...(capabilities.promptlessTurnContinuation === true ? { continuation: true } : { input: SERVER_UPDATE_CONTINUATION_PROMPT }), interactionMode: thread.interactionMode, }); + // Resume/continuation often does not emit turn.started until the + // replacement turn finishes. Leave the session running with the + // admitted turn id so queue.steer is not blocked for the whole turn. + const continuedAt = DateTime.formatIso(yield* DateTime.now); + yield* orchestrationEngine.dispatch({ + type: "thread.session.set", + commandId: CommandId.make(yield* crypto.randomUUIDv4), + threadId: thread.id, + session: { + ...session, + status: "running", + activeTurnId: result.turnId, + lastError: null, + updatedAt: continuedAt, + }, + createdAt: continuedAt, + }); }); const continuationExit = yield* Effect.exit(continuation); if (Exit.isSuccess(continuationExit) || Cause.hasInterrupts(continuationExit.cause)) { diff --git a/apps/web/src/components/files/fileEditorLanguageReadiness.test.ts b/apps/web/src/components/files/fileEditorLanguageReadiness.test.ts index 520c0fa82d0e..52253d6c27d6 100644 --- a/apps/web/src/components/files/fileEditorLanguageReadiness.test.ts +++ b/apps/web/src/components/files/fileEditorLanguageReadiness.test.ts @@ -40,6 +40,7 @@ const source = "export const View = () =>
Ready
;"; let pool: WorkerPoolManager; let renderer: FileRenderer; let terminationPromises: Promise[]; +let animationFrames: Set>; class WorkerTransport { private readonly worker = new NodeWorkerThreads.Worker( @@ -96,10 +97,19 @@ function firstEnter(highlighter: DiffsHighlighter, file: FileContents, language: beforeEach(async () => { terminationPromises = []; - vi.stubGlobal("requestAnimationFrame", (callback: FrameRequestCallback) => - setImmediate(() => callback(0)), - ); - vi.stubGlobal("cancelAnimationFrame", clearImmediate); + animationFrames = new Set(); + vi.stubGlobal("requestAnimationFrame", (callback: FrameRequestCallback) => { + const frame = setImmediate(() => { + animationFrames.delete(frame); + callback(0); + }); + animationFrames.add(frame); + return frame; + }); + vi.stubGlobal("cancelAnimationFrame", (frame: ReturnType) => { + animationFrames.delete(frame); + clearImmediate(frame); + }); vi.stubGlobal("window", { matchMedia: () => ({ matches: true }) }); await disposeHighlighter(); pool = new WorkerPoolManager( @@ -116,6 +126,9 @@ afterEach(async () => { pool?.terminate(); await Promise.all(terminationPromises); await disposeHighlighter(); + // Pool termination can queue a final broadcast after its workers have exited. + for (const frame of animationFrames) clearImmediate(frame); + animationFrames.clear(); vi.unstubAllGlobals(); });