diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 0a913831890d..abdfc29579cd 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -4382,8 +4382,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const source = projection.runs.find((run) => run.id === command.restartContinuationOfRunId); if ( !source || - (source.status !== "cancelled" && - !isRestartNoteSource(source, projection.providerTurns)) || + source.status !== "cancelled" || + isRestartNoteSource(source, projection.providerTurns) || projection.thread.archivedAt !== null || projection.thread.deletedAt !== null || projection.thread.providerInstanceId !== source.providerInstanceId || diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index 98287944dee1..6eb0ef12e605 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -149,33 +149,6 @@ function resolveStaleBackgroundItemProviderInstanceId( return projection.providerThreads[0]?.providerInstanceId ?? projection.thread.providerInstanceId; } -/** - * Provider threads with background work reconciliation would cancel and - * record for their next turn. - */ -function providerThreadsWithOpenBackgroundWork( - projection: ProjectionStore.ProjectionRuntimeRecoveryState, -): ReadonlySet { - const ids = new Set(); - for (const item of projection.turnItems ?? []) { - if ( - !isBackgroundCapableTurnItemType(item.type) || - !isNonterminalTurnItemStatus(item.status) || - isAppOwnedDelegationItem(item) - ) - continue; - const providerThreadId = - item.providerThreadId ?? - projection.runs.find((run) => run.id === item.runId)?.providerThreadId; - if (providerThreadId != null) ids.add(providerThreadId); - } - for (const thread of projection.providerThreads ?? []) { - if (thread.ownerNodeId === null && providerThreadHasPendingBackgroundTasks(thread)) - ids.add(thread.id); - } - return ids; -} - /** * A provider thread's latest started run: the last turn that provider saw. * Restart recovery records the thread's cancelled background work on it, and @@ -668,7 +641,7 @@ export const make = Effect.gen(function* () { } const continuationRun = continueAfterRestart && trigger === "startup" - ? restartContinuationRun(projection, new Set(cancelledBackgroundWork.keys())) + ? restartContinuationRun(projection) : undefined; const effects: Array = continuationRun ? [ @@ -804,12 +777,7 @@ export const make = Effect.gen(function* () { .continueThreadsAfterServerUpdate ) return; - // Shutdown reconciliation cancels the background work below, so a - // settled thread's continuation must be captured while it is still open. - const run = restartContinuationRun( - projection, - providerThreadsWithOpenBackgroundWork(projection), - ); + const run = restartContinuationRun(projection); if (!run) return; const commandId = CommandId.make(`command:restart-prepare:${run.id}`); yield* eventSink.writeWithEffects({ diff --git a/apps/server/src/orchestration-v2/RestartContinuation.test.ts b/apps/server/src/orchestration-v2/RestartContinuation.test.ts index e7cd4b9476a7..60704e23ec83 100644 --- a/apps/server/src/orchestration-v2/RestartContinuation.test.ts +++ b/apps/server/src/orchestration-v2/RestartContinuation.test.ts @@ -32,6 +32,8 @@ const driver = ProviderDriverKind.make("codex"); const providerThreadId = ProviderThreadId.make("provider-thread:restart"); const sessionId = ProviderSessionId.make("session:restart"); const attemptId = RunAttemptId.make("attempt:restart"); +// "No project" threads belong to the environment's Scratch project. +const scratchProjectId = ProjectId.make("project:scratch"); function makeProjection() { return { @@ -155,76 +157,63 @@ it("recovers an admitted continuation after another crash before provider start" assert.equal(restartContinuationRun(starting)?.id, runId); }); -it("continues a settled root run only when the restart cancelled its background work", () => { +it("does not continue settled root runs with restart-cancelled background work", () => { const projection = makeProjection(); const settled = { ...projection, - runs: [{ ...projection.runs[0]!, status: "completed" as const }], + runs: [ + { + ...projection.runs[0]!, + status: "completed" as const, + restartCancelledBackgroundWork: [{ kind: "shell" as const, label: "sleep 25" }], + }, + ], providerThreads: [{ ...projection.providerThreads[0]!, status: "idle" as const }], providerSessions: [{ ...projection.providerSessions[0]!, status: "stopped" as const }], providerTurns: [], }; - const lostWork = new Set([providerThreadId]); - assert.isUndefined(restartContinuationRun(settled)); - assert.equal(restartContinuationRun(settled, lostWork)?.id, runId); - // Work an older provider thread launched (before a provider switch) is not - // this run's: its provider was never told about it and cannot continue it. - assert.isUndefined( - restartContinuationRun(settled, new Set([ProviderThreadId.make("provider-thread:claude")])), - ); - for (const invalid of [ - { ...settled, thread: { ...settled.thread, archivedAt: {} } }, - { ...settled, thread: { ...settled.thread, deletedAt: {} } }, - { ...settled, runs: [{ ...settled.runs[0]!, status: "failed" as const }] }, - { - ...settled, - providerThreads: [{ ...settled.providerThreads[0]!, nativeThreadRef: null }], - }, - ]) - assert.isUndefined( - restartContinuationRun(invalid as OrchestrationV2ThreadProjection, lostWork), - ); + for (const projectId of [scratchProjectId, projection.thread.projectId]) + for (const status of ["completed", "waiting"] as const) + assert.isUndefined( + restartContinuationRun({ + ...settled, + thread: { ...settled.thread, projectId }, + runs: [{ ...settled.runs[0]!, status }], + }), + ); }); -it.effect("prompts a settled thread's continuation with the note of its lost work", () => - Effect.gen(function* () { - const base = makeProjection(); - const work = [{ kind: "shell" as const, label: "sleep 25 && echo DONE" }]; - const projection = { - ...base, - runs: [{ ...base.runs[0]!, status: "completed", restartCancelledBackgroundWork: work }], - providerTurns: [{ ...base.providerTurns[0]!, status: "completed" }], - } as unknown as OrchestrationV2ThreadProjection; - const commands: Parameters< - ThreadManagementService.ThreadManagementService["Service"]["dispatch"] - >[0][] = []; - yield* continueRestartedRun({ threadId, sourceRunId: runId }).pipe( - Effect.provide( - Layer.merge( - Layer.mock(ThreadManagementService.ThreadManagementService)({ - getThreadRecords: () => Effect.succeed(projection), - recoverDelegatedTask: () => Effect.void, - dispatch: (command) => { - commands.push(command); - return Effect.succeed({} as never); - }, - }), - ServerSettings.layerTest({ continueThreadsAfterServerUpdate: true }), +it.effect.each(["completed", "waiting", "cancelled"] as const)( + "ignores a pending restart continuation for a %s run whose provider turn settled", + (status) => + Effect.gen(function* () { + const base = makeProjection(); + const work = [{ kind: "shell" as const, label: "sleep 25 && echo DONE" }]; + const projection = { + ...base, + runs: [{ ...base.runs[0]!, status, restartCancelledBackgroundWork: work }], + providerTurns: [{ ...base.providerTurns[0]!, status: "completed" }], + } as unknown as OrchestrationV2ThreadProjection; + const commands: Parameters< + ThreadManagementService.ThreadManagementService["Service"]["dispatch"] + >[0][] = []; + yield* continueRestartedRun({ threadId, sourceRunId: runId }).pipe( + Effect.provide( + Layer.merge( + Layer.mock(ThreadManagementService.ThreadManagementService)({ + getThreadRecords: () => Effect.succeed(projection), + recoverDelegatedTask: () => Effect.void, + dispatch: (command) => { + commands.push(command); + return Effect.succeed({} as never); + }, + }), + ServerSettings.layerTest({ continueThreadsAfterServerUpdate: true }), + ), ), - ), - ); - assert.lengthOf(commands, 1); - const command = commands[0]!; - assert.equal( - command.type === "message.dispatch" ? command.restartContinuationOfRunId : null, - runId, - ); - assert.include( - command.type === "message.dispatch" ? command.text : "", - "sleep 25 && echo DONE", - ); - assert.notInclude(command.type === "message.dispatch" ? command.text : "", "Continue where"); - }), + ); + assert.lengthOf(commands, 0); + }), ); it.effect("does not continue a failed run that lost background work", () => @@ -419,7 +408,7 @@ it.effect("prepares no continuation for background work another provider thread ); yield* recovery.prepareForShutdown; assert.lengthOf(writes, 0); - // The same work on the Codex run's own provider thread is continued. + // Work on the settled run's own provider thread must not wake it either. const ownWork = { ...projection, turnItems: [{ ...projection.turnItems[0]!, providerThreadId }], @@ -446,13 +435,103 @@ it.effect("prepares no continuation for background work another provider thread ), ); yield* ownRecovery.prepareForShutdown; - assert.deepEqual( - writes.map((write) => write.effects[0]?.request), - [{ type: "provider-runtime.continue", sourceRunId: runId }], - ); + assert.lengthOf(writes, 0); }), ); +it.effect.each([ + ["completed", scratchProjectId], + ["completed", ProjectId.make("restart-project")], + ["waiting", scratchProjectId], + ["waiting", ProjectId.make("restart-project")], +] as const)( + "cleans up background work without waking a %s run in project %s", + ([status, projectId]) => + Effect.gen(function* () { + const base = makeProjection(); + const projection = { + ...base, + thread: { ...base.thread, projectId }, + runs: [{ ...base.runs[0]!, status }], + providerTurns: [{ ...base.providerTurns[0]!, status: "completed" }], + turnItems: [ + { + id: "turn-item:background-subagent", + runId, + nodeId: null, + providerThreadId, + providerInstanceId: instanceId, + type: "subagent", + subagentId: "subagent:background", + title: "Background reviewer", + status: "running", + }, + ], + subagents: [ + { + id: "subagent:background", + runId, + driver, + providerInstanceId: instanceId, + status: "running", + }, + ], + } as unknown as OrchestrationV2ThreadProjection; + for (const trigger of ["startup", "shutdown"] as const) { + const commits: Parameters[0][] = []; + const writes: Parameters[0][] = []; + const recovery = yield* ProviderRuntimeRecovery.make.pipe( + Effect.provide( + Layer.mergeAll( + ServerSettings.layerTest({ continueThreadsAfterServerUpdate: true }), + Layer.mock(ProjectionStore.ProjectionStoreV2)({ + getRecoveryThreadIds: () => Effect.succeed([threadId]), + getRuntimeRecoveryProjection: () => Effect.succeed(projection), + }), + Layer.mock(EventSink.EventSinkV2)({ + writeWithEffects: (input) => + Effect.sync(() => { + writes.push(input); + return []; + }), + commitCommand: (input) => + Effect.sync(() => { + commits.push(input); + return { committed: true, cancelledEffectCount: 0 } as never; + }), + }), + IdAllocator.layer, + Layer.mock(EffectWorker.OrchestrationEffectWorkerV2)({ + runRecoveryOnce: Effect.succeed(false), + }), + Layer.mock(EffectOutbox.EffectOutboxV2)({ + listByCommandId: () => Effect.succeed([]), + reconcileAfterProcessLoss: Effect.succeed({ requeued: 0, cancelled: 0 }), + }), + ), + ), + ); + if (trigger === "shutdown") yield* recovery.prepareForShutdown; + yield* recovery.reconcile(trigger); + assert.lengthOf(writes, 0); + assert.lengthOf(commits, 1); + assert.lengthOf(commits[0]!.effects, 0); + const events = commits[0]!.events; + for (const type of ["turn-item.updated", "subagent.updated"] as const) + assert.isTrue( + events.some((event) => event.type === type && event.payload.status === "cancelled"), + ); + const note = events.find((event) => event.type === "run.background-work-cancelled"); + assert.equal(note?.runId, runId); + assert.deepEqual(note?.payload.restartCancelledBackgroundWork, [ + { kind: "subagent", label: "Background reviewer", id: "turn-item:background-subagent" }, + ]); + if (status === "completed") + assert.isFalse(events.some((event) => event.type === "run.updated")); + } + }), +); + it.effect("does not cancel or resume a run that completes while shutdown intent commits", () => Effect.gen(function* () { let projection = makeProjection(); diff --git a/apps/server/src/orchestration-v2/RestartContinuation.ts b/apps/server/src/orchestration-v2/RestartContinuation.ts index f91d1a1e5f73..013e7a07fbc4 100644 --- a/apps/server/src/orchestration-v2/RestartContinuation.ts +++ b/apps/server/src/orchestration-v2/RestartContinuation.ts @@ -4,7 +4,6 @@ import { CommandId, MessageId, type OrchestrationV2Run, - type ProviderThreadId, type RunId, type ThreadId, } from "@t3tools/contracts"; @@ -23,16 +22,14 @@ import { const CONTINUE_PROMPT = "Continue where you left off."; /** - * The run a restart continuation resumes, if any: an unfinished root run, or a - * settled one whose own provider thread lost background work in the restart - * (`cancelledWorkProviderThreadIds`, which recovery records on that thread). + * Resume only an unfinished root turn. Leftover background work is cleaned up + * separately and reported on the next user turn; it must not wake a settled run. */ export function restartContinuationRun( projection: Pick< ProjectionRuntimeRecoveryState, "thread" | "runs" | "providerThreads" | "providerSessions" | "providerTurns" >, - cancelledWorkProviderThreadIds: ReadonlySet = new Set(), ): OrchestrationV2Run | undefined { if (projection.thread.archivedAt !== null || projection.thread.deletedAt !== null) return; // Queued runs never started; recovery holds them behind the cut run. @@ -46,13 +43,8 @@ export function restartContinuationRun( if (!run) return; const preparedContinuation = run.status === "starting" && run.restartContinuationOfRunId !== undefined; - // Background work outlived this settled turn; the provider has no live turn. - const settledWithCancelledWork = - (run.status === "completed" || run.status === "waiting") && - run.providerThreadId !== null && - cancelledWorkProviderThreadIds.has(run.providerThreadId); - if (run.status !== "running" && !preparedContinuation && !settledWithCancelledWork) return; - const liveTurnRequired = !preparedContinuation && !settledWithCancelledWork; + if (run.status !== "running" && !preparedContinuation) return; + const liveTurnRequired = !preparedContinuation; if (projection.thread.providerInstanceId !== run.providerInstanceId) return; const providerThread = projection.providerThreads.find( (thread) => thread.id === run.providerThreadId, @@ -73,16 +65,13 @@ export function restartContinuationRun( const session = projection.providerSessions.find( (candidate) => candidate.id === providerThread.providerSessionId, ); - // A settled thread's session may already be stopped and out of the recovery - // read; the continuation reopens it from the provider thread's native ref. // Most adapters keep a live session "ready" through its turns, so only a // stopped or failed session rules out a live turn. if ( - session === undefined - ? !settledWithCancelledWork - : session.providerInstanceId !== run.providerInstanceId || - session.driver !== providerThread.driver || - (liveTurnRequired && (session.status === "stopped" || session.status === "error")) + session === undefined || + session.providerInstanceId !== run.providerInstanceId || + session.driver !== providerThread.driver || + (liveTurnRequired && (session.status === "stopped" || session.status === "error")) ) return; if ( @@ -119,10 +108,14 @@ export const continueRestartedRun = Effect.fn("RestartContinuation.continueResta if (projection.messages.some((message) => message.id === messageId)) return; const source = projection.runs.find((run) => run.id === input.sourceRunId); - // A settled source prompts with the note of the background work it lost. - const noteSource = - source !== undefined && isRestartNoteSource(source, projection.providerTurns); - if (!source || (source.status !== "cancelled" && !noteSource)) return; + // Pending effects from older versions may target settled background work, + // including waiting runs that reconciliation subsequently cancelled. + if ( + !source || + source.status !== "cancelled" || + isRestartNoteSource(source, projection.providerTurns) + ) + return; // A user submission after reconciliation takes precedence over an automatic // prompt. Queued runs never started and stay held behind this one. if ( diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 84d6df90160a..c2f0024b0542 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -2645,84 +2645,86 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }), ); - it.effect("does not admit a restart continuation of a failed run that lost background work", () => - Effect.gen(function* () { - const orchestrator = yield* Orchestrator.OrchestratorV2; - const eventSink = yield* EventSink.EventSinkV2; - const threadId = ThreadId.make("runtime-layer-restart-failed-source"); - yield* orchestrator.dispatch({ - type: "thread.create", - createdBy: "user", - creationSource: "web", - commandId: CommandId.make("restart-failed-create"), - threadId, - projectId: ProjectId.make("restart-project"), - title: "Restart", - modelSelection, - runtimeMode: "full-access", - interactionMode: "default", - branch: null, - worktreePath: "/tmp/runtime-layer-restart-failed", - }); - yield* orchestrator.dispatch({ - type: "message.dispatch", - createdBy: "user", - creationSource: "web", - commandId: CommandId.make("restart-failed-user-message"), - threadId, - messageId: MessageId.make("restart-failed-user-message"), - text: "Original work", - attachments: [], - modelSelection, - dispatchMode: { type: "start_immediately" }, - }); - const original = (yield* orchestrator.getThreadProjection(threadId)).runs[0]!; - const now = yield* DateTime.now; - yield* eventSink.commitCommand({ - commandId: CommandId.make("restart-failed-reconcile"), - threadId, - commandType: "provider-runtime.reconcile", - acceptedAt: now, - events: [ - { - id: EventId.make("restart-failed-run"), - type: "run.updated", - threadId, - runId: original.id, - occurredAt: now, - payload: { ...original, status: "failed", completedAt: now }, - }, - { - id: EventId.make("restart-failed-work"), - type: "run.background-work-cancelled", - threadId, - runId: original.id, - occurredAt: now, - payload: { + it.effect.each(["failed", "completed", "waiting", "cancelled"] as const)( + "does not admit a restart continuation of a %s run with only lost background work", + (status) => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const eventSink = yield* EventSink.EventSinkV2; + const threadId = ThreadId.make(`runtime-layer-restart-${status}-source`); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`restart-${status}-create`), + threadId, + projectId: ProjectId.make("restart-project"), + title: "Restart", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: "/tmp/runtime-layer-restart-failed", + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`restart-${status}-user-message`), + threadId, + messageId: MessageId.make(`restart-${status}-user-message`), + text: "Original work", + attachments: [], + modelSelection, + dispatchMode: { type: "start_immediately" }, + }); + const original = (yield* orchestrator.getThreadProjection(threadId)).runs[0]!; + const now = yield* DateTime.now; + yield* eventSink.commitCommand({ + commandId: CommandId.make(`restart-${status}-reconcile`), + threadId, + commandType: "provider-runtime.reconcile", + acceptedAt: now, + events: [ + { + id: EventId.make(`restart-${status}-run`), + type: "run.updated", + threadId, runId: original.id, - restartCancelledBackgroundWork: [{ kind: "shell", label: "sleep 25" }], + occurredAt: now, + payload: { ...original, status, completedAt: now }, }, - }, - ], - effects: [], - }); - yield* orchestrator.dispatch({ - type: "message.dispatch", - createdBy: "agent", - creationSource: "server", - commandId: CommandId.make("restart-failed-continuation"), - threadId, - messageId: MessageId.make("restart-failed-continuation"), - text: "Note: the T3 server restarted.", - attachments: [], - modelSelection, - dispatchMode: { type: "start_immediately" }, - restartContinuationOfRunId: original.id, - }); - const projection = yield* orchestrator.getThreadProjection(threadId); - assert.lengthOf(projection.runs, 1); - assert.equal(projection.runs[0]?.status, "failed"); - }), + { + id: EventId.make(`restart-${status}-work`), + type: "run.background-work-cancelled", + threadId, + runId: original.id, + occurredAt: now, + payload: { + runId: original.id, + restartCancelledBackgroundWork: [{ kind: "shell", label: "sleep 25" }], + }, + }, + ], + effects: [], + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "agent", + creationSource: "server", + commandId: CommandId.make(`restart-${status}-continuation`), + threadId, + messageId: MessageId.make(`restart-${status}-continuation`), + text: "Note: the T3 server restarted.", + attachments: [], + modelSelection, + dispatchMode: { type: "start_immediately" }, + restartContinuationOfRunId: original.id, + }); + const projection = yield* orchestrator.getThreadProjection(threadId); + assert.lengthOf(projection.runs, 1); + assert.equal(projection.runs[0]?.status, status); + }), ); it.effect("rejects settling a thread while a run is active", () => diff --git a/apps/server/src/orchestration-v2/testkit/OrchestratorReplayRestartBackgroundNote.integration.test.ts b/apps/server/src/orchestration-v2/testkit/OrchestratorReplayRestartBackgroundNote.integration.test.ts index fc61c61646bc..d9e6a7f13161 100644 --- a/apps/server/src/orchestration-v2/testkit/OrchestratorReplayRestartBackgroundNote.integration.test.ts +++ b/apps/server/src/orchestration-v2/testkit/OrchestratorReplayRestartBackgroundNote.integration.test.ts @@ -11,6 +11,7 @@ import { makeClaudeRestartReplayHarness, } from "../Adapters/ClaudeAdapterV2.testkit.ts"; import * as IdAllocator from "../IdAllocator.ts"; +import * as EffectWorker from "../EffectWorker.ts"; import { makeSqlitePersistenceLive } from "../../persistence/Layers/Sqlite.ts"; import { provideDeterministicTestRuntime } from "./DeterministicRuntime.ts"; import { CLAUDE_BACKGROUND_SUBAGENT_AFTER_ROOT_PROMPT } from "./fixtures/claude_background_subagent_after_root/input.ts"; @@ -19,7 +20,11 @@ import { materializeFixtureInput, projectionFor, } from "./fixtures/shared.ts"; -import { runOrchestratorV2ProviderReplayScenario } from "./ProviderReplayHarness.ts"; +import { + makeOrchestratorV2ProviderReplayLayer, + runOrchestratorV2ProviderReplayScenario, +} from "./ProviderReplayHarness.ts"; +import { runOrchestratorV2Scenario } from "./OrchestratorScenario.ts"; import { checkpointWorkspace } from "./ReplayFixtureWorkspace.ts"; import { readProviderReplayTranscript } from "./ReplayTranscriptNdjson.ts"; @@ -155,21 +160,7 @@ const runRestart = Effect.fn("runRestart")(function* (input: { // Phase 1 ends once the root run settles; the subagent is still open. const firstIdle = materialized.steps.findIndex((step) => step.type === "await_thread_idle"); const phase1Steps = materialized.steps.slice(0, firstIdle + 1); - const threadId = materialized.projectionThreadIds[0]!; - const phase2Steps = [ - // Let recovery's continuation (run 2) finish before the next user message. - ...(input.continueThreadsAfterServerUpdate - ? [ - { - type: "await_run_status" as const, - threadId, - runId: (yield* IdAllocator.IdAllocatorV2).derive.run({ threadId, ordinal: 2 }), - status: "completed" as const, - }, - ] - : []), - ...materialized.steps.slice(firstIdle + 1), - ]; + const phase2Steps = materialized.steps.slice(firstIdle + 1); const { harness, assertComplete } = makeClaudeRestartReplayHarness(transcript); const databaseLayer = makeSqlitePersistenceLive(path.join(tempDir, "state.sqlite")).pipe( Layer.provide(NodeServices.layer), @@ -192,10 +183,35 @@ const runRestart = Effect.fn("runRestart")(function* (input: { assert.equal(settled.runs[0]?.status, "completed"); assert.equal(settled.subagents[0]?.status, "running"); + // Drain recovery before submitting any user work, which would otherwise + // take precedence over an incorrectly queued automatic continuation. + const restartScenario = scenario("restarted", []); + const restarted = yield* Effect.scoped( + Effect.gen(function* () { + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + yield* worker.drain(); + return yield* runOrchestratorV2Scenario(restartScenario); + }).pipe( + Effect.provide( + makeOrchestratorV2ProviderReplayLayer(restartScenario, harness, { + databaseLayer, + recoverOnStartup: true, + continueThreadsAfterServerUpdate: input.continueThreadsAfterServerUpdate, + }), + ), + ), + ); + const recovered = projectionFor(restarted, SCENARIO); + assert.deepEqual( + recovered.runs.map((run) => run.status), + ["completed"], + ); + assert.equal(recovered.subagents[0]?.status, "cancelled"); + assert.equal(recovered.runs[0]?.restartCancelledBackgroundWork?.[0]?.kind, "subagent"); + const after = yield* Effect.scoped( runOrchestratorV2ProviderReplayScenario(scenario("after-restart", phase2Steps), harness, { databaseLayer, - recoverOnStartup: true, continueThreadsAfterServerUpdate: input.continueThreadsAfterServerUpdate, }), ); @@ -210,61 +226,37 @@ const userTexts = (projection: ReturnType) => projection.turnItems.flatMap((item) => (item.type === "user_message" ? [item.text] : [])); describe("restart-cancelled background work", () => { - it.effect("tells the next provider turn once that its background subagent died", () => - Effect.scoped( - Effect.gen(function* () { - const projection = yield* runRestart({ - // The first turn after the restart carries the note; the second does not. - resumedPrompts: [ - `${NOTE}\n\nUser message:\n${FIRST_AFTER_RESTART}`, + it.effect.each([false, true])( + "keeps the thread settled and tells the next user turn once when continuation is %s", + (continueThreadsAfterServerUpdate) => + Effect.scoped( + Effect.gen(function* () { + const projection = yield* runRestart({ + // The first turn after the restart carries the note; the second does not. + resumedPrompts: [ + `${NOTE}\n\nUser message:\n${FIRST_AFTER_RESTART}`, + SECOND_AFTER_RESTART, + ], + userMessagesAfterRestart: [FIRST_AFTER_RESTART, SECOND_AFTER_RESTART], + continueThreadsAfterServerUpdate, + }); + assert.deepEqual( + projection.runs.map((run) => run.status), + ["completed", "completed", "completed"], + ); + assert.isFalse( + projection.runs.some((run) => run.restartContinuationOfRunId !== undefined), + ); + // The note reaches the provider only; the timeline keeps what the user sent. + assert.deepEqual(userTexts(projection), [ + CLAUDE_BACKGROUND_SUBAGENT_AFTER_ROOT_PROMPT, + FIRST_AFTER_RESTART, SECOND_AFTER_RESTART, - ], - userMessagesAfterRestart: [FIRST_AFTER_RESTART, SECOND_AFTER_RESTART], - continueThreadsAfterServerUpdate: false, - }); - assert.deepEqual( - projection.runs.map((run) => run.status), - ["completed", "completed", "completed"], - ); - assert.isFalse(projection.runs.some((run) => run.restartContinuationOfRunId !== undefined)); - // The note reaches the provider only; the timeline keeps what the user sent. - assert.deepEqual(userTexts(projection), [ - CLAUDE_BACKGROUND_SUBAGENT_AFTER_ROOT_PROMPT, - FIRST_AFTER_RESTART, - SECOND_AFTER_RESTART, - ]); - }).pipe( - provideDeterministicTestRuntime, - Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer)), - ), - ), - ); - - it.effect("continues a settled thread with the note when restart continuation is on", () => - Effect.scoped( - Effect.gen(function* () { - const projection = yield* runRestart({ - // Recovery's continuation prompts with the note; the user turn after it does not. - resumedPrompts: [NOTE, FIRST_AFTER_RESTART], - userMessagesAfterRestart: [FIRST_AFTER_RESTART], - continueThreadsAfterServerUpdate: true, - }); - const [root, continuation, user] = projection.runs; - assert.deepEqual( - projection.runs.map((run) => run.status), - ["completed", "completed", "completed"], - ); - assert.equal(continuation?.restartContinuationOfRunId, root?.id); - assert.isUndefined(user?.restartContinuationOfRunId); - assert.deepEqual(userTexts(projection), [ - CLAUDE_BACKGROUND_SUBAGENT_AFTER_ROOT_PROMPT, - NOTE, - FIRST_AFTER_RESTART, - ]); - }).pipe( - provideDeterministicTestRuntime, - Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer)), + ]); + }).pipe( + provideDeterministicTestRuntime, + Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer)), + ), ), - ), ); });