diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts index 29c96447e9f4..877b02049e79 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.test.ts @@ -3492,6 +3492,51 @@ describe("ClaudeAdapterV2 background wake turns", () => { ), ); + // The CLI process exits on its own after the turn settled (idle, crash), + // leaving its background task on the roster. Stop must succeed so the + // orchestrator goes on to settle what the thread still shows. + it.effect("a settled Stop with no CLI process left succeeds and clears the roster", () => + Effect.scoped( + Effect.gen(function* () { + const harness = yield* makeWakeHarnessWithOptions(); + const now = yield* DateTime.now; + yield* harness.runtime.startTurn( + makeClaudeTestTurnInput({ + threadId: harness.threadId, + providerThread: harness.providerThread, + now, + attemptId: RunAttemptId.make("attempt-claude-settled-stop-no-process"), + text: "Run the build in the background.", + attachments: [], + }), + ); + yield* Queue.offer(harness.sdkMessages, wakeTaskStarted); + yield* Queue.offer(harness.sdkMessages, turnOneResult); + yield* awaitUntil(() => harness.terminalEvents().length === 1, "first turn terminal"); + const settledThread = providerThreadRosterEvents(harness.events).at(-1)?.providerThread; + assert.equal(settledThread?.pendingBackgroundTasks?.[0]?.taskId, WAKE_TASK_ID); + + yield* Queue.shutdown(harness.sdkMessages); + let quietYields = 0; + yield* awaitUntil(() => quietYields++ >= 50, "query exit"); + + yield* harness.runtime.interruptTurn({ + providerThread: settledThread ?? harness.providerThread, + providerTurnId: harness.terminalEvents()[0]!.providerTurnId, + requestRuntimeRestart: true, + }); + yield* awaitUntil( + () => + (providerThreadRosterEvents(harness.events).at(-1)?.providerThread + .pendingBackgroundTasks?.length ?? 0) === 0, + "roster clear after Stop", + ); + assert.isFalse(yield* harness.hasPendingBackgroundWork); + assert.lengthOf(harness.continuationRequests, 0); + }).pipe(Effect.provide(Layer.merge(IdAllocator.layer, NodeServices.layer))), + ), + ); + it.effect("a settled Stop leaves a turn that replaced the closing process alone", () => Effect.scoped( Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts index eeaf367af70a..6f08089cca46 100644 --- a/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/ClaudeAdapterV2.ts @@ -6940,22 +6940,23 @@ export function makeClaudeAdapterV2( const interruptTurn = Effect.fn("ClaudeAdapterV2.interruptTurn")( function* (turnInput: ProviderAdapter.ProviderAdapterV2InterruptInput) { const existing = yield* Ref.get(queryContext); - if (existing === null) { - return yield* new ProviderAdapter.ProviderAdapterProtocolError({ - driver: CLAUDE_PROVIDER, - detail: `Claude provider thread ${turnInput.providerThread.id} has no live query.`, - }); - } const currentTurn = yield* Ref.get(activeTurn); const nativeThreadId = turnInput.providerThread.nativeThreadRef?.nativeId ?? null; - if ( - currentTurn === null && - turnInput.requestRuntimeRestart === true && - nativeThreadId !== null && - existing.nativeThreadId === nativeThreadId - ) { - // Stop after the turn settled: the background shells belong to - // the CLI process, so closing its query is what stops them. + if (currentTurn === null && turnInput.requestRuntimeRestart === true) { + // Stop after the turn settled. With no CLI process of this + // native thread left, nothing it started is still running: its + // roster is not authoritative any more, and the orchestrator + // settles the items the thread still shows. + if (nativeThreadId === null) return; + if (existing === null || existing.nativeThreadId !== nativeThreadId) { + yield* clearWakeStateForNativeThread(nativeThreadId); + yield* resetBackgroundTaskStateForNativeThreadProcess(nativeThreadId, { + status: "idle", + }); + return; + } + // The background shells belong to the CLI process, so closing + // its query is what stops them. yield* closeLiveQueryForNativeThread(nativeThreadId); // A turn started while the close was pending may have opened a // replacement process. Its Waiting and wake state are its own. @@ -6968,6 +6969,12 @@ export function makeClaudeAdapterV2( } return; } + if (existing === null) { + return yield* new ProviderAdapter.ProviderAdapterProtocolError({ + driver: CLAUDE_PROVIDER, + detail: `Claude provider thread ${turnInput.providerThread.id} has no live query.`, + }); + } if (currentTurn?.providerTurnId !== turnInput.providerTurnId) { return yield* new ProviderAdapter.ProviderAdapterProtocolError({ driver: CLAUDE_PROVIDER, diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts index b27309f2e81e..27fadf2c4381 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.test.ts @@ -3987,6 +3987,126 @@ describe("CodexAdapterV2 post-settle continuation", () => { } } + // The app-server exits after the root turn, before the command's own + // item/completed (Codex always sends one, so only a lost notification or a + // gone process leaves it running). Nothing tracks the command any more, yet + // the thread still shows it, and Stop is the only way to clear it. + it.effect("Stop ends a background command no Codex process tracks any more", () => + Effect.scoped( + Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const cwd = yield* fs.makeTempDirectoryScoped({ prefix: "t3-bg-stale-workspace-" }); + const staleTranscript = makeCodexReplayTranscript({ + scenario: "codex-bg-stop-untracked", + entries: [ + ...backgroundExecTranscript.entries.slice(0, -1), + { type: "runtime_exit", status: "success" }, + ], + }); + const localTranscript = yield* decodeReplayTranscriptJson( + (yield* encodeReplayTranscriptJson(staleTranscript)).replaceAll( + yield* encodeStringJson("/workspace"), + yield* encodeStringJson(cwd), + ), + ); + const spawner = yield* ChildProcessSpawner.ChildProcessSpawner; + for (const args of [ + ["init", "--quiet"], + [ + "-c", + "user.name=Test", + "-c", + "user.email=test@example.com", + "commit", + "--allow-empty", + "--quiet", + "-m", + "Initial commit", + ], + ]) { + assert.equal(Number(yield* spawner.exitCode(ChildProcess.make("git", args, { cwd }))), 0); + } + yield* Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const worker = yield* EffectWorker.OrchestrationEffectWorkerV2; + const threadId = ThreadId.make("thread:background-stop-untracked"); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("create-background-stop-untracked"), + threadId, + projectId: ProjectId.make("project:background-stop-untracked"), + title: "Background stop untracked", + modelSelection: CODEX_TEST_MODEL_SELECTION, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: cwd, + createdBy: "user", + creationSource: "web", + }); + const settled = yield* orchestrator.streamDomainEvents.pipe( + Stream.filter( + (event) => + event.type === "run.updated" && + (event.payload.status === "waiting" || event.payload.status === "completed"), + ), + Stream.runHead, + Effect.forkChild({ startImmediately: true }), + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make("start-background-stop-untracked"), + threadId, + messageId: MessageId.make("message:background-stop-untracked"), + text: BG_PROMPT, + attachments: [], + createdBy: "user", + creationSource: "web", + dispatchMode: { type: "start_immediately" }, + }); + yield* worker.drain(); + yield* Fiber.join(settled); + yield* worker.drain(); + const before = yield* orchestrator.getThreadShell(threadId); + assert.deepEqual( + before?.pendingBackgroundTasks?.map((task) => task.kind), + ["command"], + "the thread still shows the command the gone process never finished", + ); + const run = (yield* orchestrator.getThreadProjection(threadId)).runs.at(-1)!; + yield* orchestrator.dispatch({ + type: "run.interrupt", + commandId: CommandId.make("stop-background-untracked"), + threadId, + runId: run.id, + holdQueue: true, + }); + yield* worker.drain(); + const projection = yield* orchestrator.getThreadProjection(threadId); + assert.deepEqual( + projection.turnItems.flatMap((item) => + item.type === "command_execution" ? [item.status] : [], + ), + ["interrupted"], + ); + assert.equal(projection.runs.at(-1)?.status, "completed"); + assert.deepEqual( + (yield* orchestrator.getThreadShell(threadId))?.pendingBackgroundTasks, + [], + ); + }).pipe( + Effect.provide( + makeOrchestratorV2ReplayLayerWithRegistry( + { name: "codex-background-stop-untracked", runtimePolicyOverride: { cwd } }, + makeCodexProviderAdapterRegistryReplayLayer({ transcript: localTranscript }), + { runEffectWorker: false }, + ), + ), + ); + }).pipe(Effect.provide(NodeServices.layer)), + ), + ); + const PRE_SETTLE_SCENARIO = "codex-bg-exec-pre-settle"; const PRE_SETTLE_NATIVE_THREAD = "native-codex-pre-settle-thread"; const PRE_SETTLE_NATIVE_TURN = "native-codex-pre-settle-turn"; diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 7ceecc4f487f..88cd870f53c2 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -2042,10 +2042,19 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi const terminateBackgroundTerminal = Effect.fn("CodexAdapterV2.terminateBackgroundTerminal")( function* (nativeThreadId: string, processId: string) { - const response = yield* client.raw.request("thread/backgroundTerminals/terminate", { - threadId: nativeThreadId, - processId, - }); + const response = yield* client.raw + .request("thread/backgroundTerminals/terminate", { + threadId: nativeThreadId, + processId, + }) + .pipe( + // The app-server that ran the terminal is gone, and with it + // the only handle to the terminal: nothing is left to stop. + Effect.catchTags({ + CodexAppServerProcessExitedError: () => Effect.succeed({ terminated: true }), + CodexAppServerInputStreamEndedError: () => Effect.succeed({ terminated: true }), + }), + ); const result = yield* decodeCodexBackgroundTerminalTerminateResponse(response); if (result.terminated) return; let cursor: string | null = null; @@ -5653,6 +5662,9 @@ export function makeCodexAdapterV2(adapterOptions: CodexAdapterV2Options): Provi ) : undefined); if (activeTurn === undefined) { + // Stop on a settled turn this process retains nothing for + // (released, restarted, or every command already reported). + if (turnInput.requestRuntimeRestart === true) return; return yield* toProtocolError( `Provider turn ${turnInput.providerTurnId} is not active and cannot be interrupted.`, ); diff --git a/apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.ts index 97ddcc7c26ab..bb9de20da64f 100644 --- a/apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CursorAdapterV2.ts @@ -2429,6 +2429,9 @@ export function makeCursorAdapterV2( interruptTurn: Effect.fn("CursorAdapterV2.interruptTurn")( function* (turnInput: ProviderAdapter.ProviderAdapterV2InterruptInput) { const context = yield* Ref.get(activeTurn); + // Stop on a settled turn: finalization already ended its tools + // and subagents, so nothing of it is left running to stop. + if (context === null && turnInput.requestRuntimeRestart === true) return; if (context?.providerTurnId !== turnInput.providerTurnId) { return yield* new ProviderAdapter.ProviderAdapterProtocolError({ driver: CursorAgentSdk.CURSOR_PROVIDER, diff --git a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts index 98203a7f91af..11d6da13c7d5 100644 --- a/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/OpenCodeAdapterV2.ts @@ -3385,6 +3385,17 @@ export function makeOpenCodeAdapterV2( const sessionId = nativeThreadId(interruptInput.providerThread); const state = threads.get(sessionId); const turn = state?.activeTurn; + // Stop on a settled turn: a background task child can outlive + // it (busySessionIds), and aborting the descendants ends it. + if ( + (turn === undefined || turn === null) && + interruptInput.requestRuntimeRestart === true + ) { + for (const controller of commandControllers.get(sessionId) ?? []) + controller.abort(); + yield* abortDescendants(sessionId); + return; + } if ( turn === undefined || turn === null || diff --git a/apps/server/src/orchestration-v2/Adapters/PiAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/PiAdapterV2.ts index 88489a1f9040..6f0bcf258acb 100644 --- a/apps/server/src/orchestration-v2/Adapters/PiAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/PiAdapterV2.ts @@ -2466,6 +2466,9 @@ export function makePiAdapterV2( interruptTurn: (interruptInput) => Effect.gen(function* () { const turn = threadState?.activeTurn ?? null; + // Stop on a settled turn: Pi runs nothing between prompts, so + // nothing of that turn is left to stop. + if (turn === null && interruptInput.requestRuntimeRestart === true) return; if (turn === null || turn.providerTurn.id !== interruptInput.providerTurnId) { return yield* protocolError(`Pi turn ${interruptInput.providerTurnId} is not active`); } diff --git a/apps/server/src/orchestration-v2/EffectWorker.ts b/apps/server/src/orchestration-v2/EffectWorker.ts index 3d94f1baa8df..48dcfa78e96b 100644 --- a/apps/server/src/orchestration-v2/EffectWorker.ts +++ b/apps/server/src/orchestration-v2/EffectWorker.ts @@ -166,6 +166,18 @@ export const executorLayer: Layer.Layer< providerTurnId: effect.request.providerTurnId, }) .pipe( + // The provider has stopped what it still ran and reported it. + // Whatever the thread still shows on that provider thread is + // work no process will report on, so the Stop ends it too. + Effect.andThen( + threads.dispatch({ + type: "thread.background-work.settle", + commandId: CommandId.make(`${effect.commandId}:background-work-settled`), + threadId: effect.threadId, + providerThreadId: effect.request.providerThreadId, + providerTurnId: effect.request.providerTurnId, + }), + ), Effect.mapError( (cause) => new OrchestrationEffectExecutionError({ diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index fffacab5a065..ad41ae6ffde6 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -9,6 +9,10 @@ import { ProviderDriverKind, ProviderInstanceId, ProviderSessionId, + ProviderThreadId, + ProviderTurnId, + RunAttemptId, + RunId, RuntimeRequestId, ThreadId, TurnItemId, @@ -329,3 +333,204 @@ it.effect("implements a proposed plan that the command projection leaves out", ( assert.equal((yield* projections.getPlan(threadId, planId))?.status, "completed"); }).pipe(Effect.provide(testLayer)), ); + +// Stop's settle follow-up runs after the provider interrupt returns, possibly +// long after the Stop (retries) or again (an effect replayed after a crash). +// A later run's background work is not that Stop's to end. +it.effect("settles only the stopped run's background work, once", () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const threadId = ThreadId.make("thread:settle-binding"); + const providerThreadId = ProviderThreadId.make("provider-thread:settle-binding"); + const now = yield* DateTime.now; + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("create-settle-binding"), + threadId, + projectId: ProjectId.make("project:settle-binding"), + title: "Settle binding", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "user", + creationSource: "web", + }); + yield* projections.apply({ + id: EventId.make("settle-binding:provider-thread"), + type: "provider-thread.updated", + threadId, + occurredAt: now, + payload: { + id: providerThreadId, + driver: adapter.driver, + providerInstanceId: instanceId, + providerSessionId: null, + appThreadId: threadId, + ownerNodeId: null, + nativeThreadRef: null, + nativeConversationHeadRef: null, + status: "idle", + firstRunOrdinal: 1, + lastRunOrdinal: 2, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }); + const commandItem = (ordinal: number) => TurnItemId.make(`turn-item:settle-binding:${ordinal}`); + for (const ordinal of [1, 2]) { + const runId = RunId.make(`run:settle-binding:${ordinal}`); + const attemptId = RunAttemptId.make(`attempt:settle-binding:${ordinal}`); + const nodeId = NodeId.make(`node:settle-binding:${ordinal}`); + const providerTurnId = ProviderTurnId.make(`provider-turn:settle-binding:${ordinal}`); + yield* projections.apply({ + id: EventId.make(`settle-binding:run:${ordinal}`), + type: "run.created", + threadId, + occurredAt: now, + payload: { + id: runId, + threadId, + ordinal, + providerInstanceId: instanceId, + modelSelection, + providerThreadId, + userMessageId: MessageId.make(`message:settle-binding:${ordinal}`), + rootNodeId: nodeId, + activeAttemptId: attemptId, + status: "completed", + requestedAt: now, + startedAt: now, + completedAt: now, + checkpointId: null, + contextHandoffId: null, + }, + }); + yield* projections.apply({ + id: EventId.make(`settle-binding:attempt:${ordinal}`), + type: "run-attempt.created", + threadId, + occurredAt: now, + payload: { + id: attemptId, + runId, + attemptOrdinal: 1, + rootNodeId: nodeId, + providerInstanceId: instanceId, + providerThreadId, + providerTurnId, + reason: "initial", + status: "completed", + startedAt: now, + completedAt: now, + }, + }); + yield* projections.apply({ + id: EventId.make(`settle-binding:turn:${ordinal}`), + type: "provider-turn.updated", + threadId, + occurredAt: now, + payload: { + id: providerTurnId, + providerThreadId, + nodeId, + runAttemptId: attemptId, + nativeTurnRef: null, + ordinal, + status: "completed", + startedAt: now, + completedAt: now, + }, + }); + yield* projections.apply({ + id: EventId.make(`settle-binding:item:${ordinal}`), + type: "turn-item.updated", + threadId, + runId, + occurredAt: now, + payload: { + id: commandItem(ordinal), + threadId, + runId, + nodeId, + providerThreadId, + providerTurnId, + nativeItemRef: null, + parentItemId: null, + ordinal: ordinal * 10, + status: "running", + title: `Background command ${ordinal}`, + startedAt: now, + completedAt: null, + updatedAt: now, + type: "command_execution", + input: `sleep ${ordinal}`, + }, + }); + } + const itemStatuses = Effect.map(projections.getThreadProjection(threadId), (projection) => + projection.turnItems + .flatMap((item) => (item.type === "command_execution" ? [`${item.id}:${item.status}`] : [])) + .toSorted(), + ); + // The settle that followed a Stop of run 1's turn, dispatched only after + // run 2 had settled with work of its own. + const settle = { + type: "thread.background-work.settle", + commandId: CommandId.make("stop-run-1:background-work-settled"), + threadId, + providerThreadId, + providerTurnId: ProviderTurnId.make("provider-turn:settle-binding:1"), + } as const; + yield* orchestrator.dispatch(settle); + assert.deepEqual(yield* itemStatuses, [ + `${commandItem(1)}:interrupted`, + `${commandItem(2)}:running`, + ]); + + // A settle that found nothing to end replays as a no-op, even after work + // it would match appears: its receipt is recorded with no events. + const emptySettle = { + ...settle, + commandId: CommandId.make("stop-run-1-again:background-work-settled"), + }; + const first = yield* orchestrator.dispatch(emptySettle); + assert.lengthOf(first.storedEvents, 0); + yield* projections.apply({ + id: EventId.make("settle-binding:item:late"), + type: "turn-item.updated", + threadId, + runId: RunId.make("run:settle-binding:1"), + occurredAt: now, + payload: { + id: commandItem(3), + threadId, + runId: RunId.make("run:settle-binding:1"), + nodeId: NodeId.make("node:settle-binding:1"), + providerThreadId, + providerTurnId: ProviderTurnId.make("provider-turn:settle-binding:1"), + nativeItemRef: null, + parentItemId: null, + ordinal: 30, + status: "running", + title: "Late background command", + startedAt: now, + completedAt: null, + updatedAt: now, + type: "command_execution", + input: "sleep 3", + }, + }); + const replayed = yield* orchestrator.dispatch(emptySettle); + assert.lengthOf(replayed.storedEvents, 0); + assert.deepEqual(yield* itemStatuses, [ + `${commandItem(1)}:interrupted`, + `${commandItem(2)}:running`, + `${commandItem(3)}:running`, + ]); + }).pipe(Effect.provide(testLayer)), +); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 88105bc6a3f5..df9805d504c9 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -46,7 +46,10 @@ import { ThreadId, } from "@t3tools/contracts"; import { modelSelectionsEqual } from "@t3tools/shared/model"; -import { derivePendingBackgroundWork } from "@t3tools/shared/orchestrationV2PendingBackgroundWork"; +import { + derivePendingBackgroundWork, + pendingBackgroundTurnItems, +} from "@t3tools/shared/orchestrationV2PendingBackgroundWork"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -357,6 +360,7 @@ function commandThreadId(command: OrchestrationV2ServerCommand): ThreadId { case "thread.user-input.dismiss": case "checkpoint.rollback": case "checkpoint.rollback.fail": + case "thread.background-work.settle": case "provider.switch": return command.threadId; case "delegated_task.request": @@ -7529,6 +7533,150 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); }); + /** + * Ends the background work a settled thread still shows that no provider + * process will report on: work on the provider thread whose interrupt just + * returned (`stoppedProviderThreadId`), and work on provider threads with no + * live session at all. Only the stopped run's work and older runs' is ended + * (`throughRunOrdinal`); a later run's work is its own. A dead process's + * roster goes too, as on restart. + */ + const settleBackgroundWork = (input: { + readonly command: Extract< + OrchestrationV2ServerCommand, + { readonly type: "run.interrupt" | "thread.background-work.settle" } + >; + readonly events: Ref.Ref>; + readonly projection: Pick< + OrchestrationV2ThreadProjection, + "runs" | "turnItems" | "providerThreads" + >; + readonly stoppedProviderThreadId: OrchestrationV2ProviderThread["id"] | null; + readonly throughRunOrdinal: number; + readonly now: DateTime.Utc; + }) => + Effect.gen(function* () { + const emitEvent = emit(input.events, input.command); + const runOrdinals = new Map(input.projection.runs.map((run) => [run.id, run.ordinal])); + const liveness = new Map(); + const hasLiveSession = (providerThreadId: OrchestrationV2ProviderThread["id"]) => + Effect.gen(function* () { + const known = liveness.get(providerThreadId); + if (known !== undefined) return known; + const sessionId = input.projection.providerThreads.find( + (candidate) => candidate.id === providerThreadId, + )?.providerSessionId; + const live = + sessionId !== null && + sessionId !== undefined && + Option.isSome( + yield* providerSessions + .get(sessionId) + .pipe(Effect.orElseSucceed(() => Option.none())), + ); + liveness.set(providerThreadId, live); + return live; + }); + for (const item of pendingBackgroundTurnItems({ + turnItems: input.projection.turnItems, + runs: input.projection.runs, + })) { + // An item without a run counts as the stopped run's. + const itemRunOrdinal = item.runId === null ? undefined : runOrdinals.get(item.runId); + if (itemRunOrdinal !== undefined && itemRunOrdinal > input.throughRunOrdinal) continue; + const providerThreadId = item.providerThreadId ?? null; + if ( + providerThreadId !== null && + providerThreadId !== input.stoppedProviderThreadId && + (yield* hasLiveSession(providerThreadId)) + ) { + continue; + } + const run = input.projection.runs.find((candidate) => candidate.id === item.runId); + yield* emitEvent({ + type: "turn-item.updated", + threadId: input.command.threadId, + ...(item.runId === null ? {} : { runId: item.runId }), + ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), + ...(run === undefined ? {} : { providerInstanceId: run.providerInstanceId }), + occurredAt: input.now, + payload: { + ...item, + status: "interrupted", + completedAt: input.now, + updatedAt: input.now, + }, + }); + } + for (const providerThread of input.projection.providerThreads) { + // A live process owns its roster and reports clearing it. + if ( + (providerThread.pendingBackgroundTasks?.length ?? 0) === 0 || + (yield* hasLiveSession(providerThread.id)) + ) { + continue; + } + yield* emitEvent({ + type: "provider-thread.updated", + threadId: input.command.threadId, + driver: providerThread.driver, + providerInstanceId: providerThread.providerInstanceId, + occurredAt: input.now, + payload: { ...providerThread, pendingBackgroundTasks: [], updatedAt: input.now }, + }); + } + }); + + const dispatchBackgroundWorkSettle = ( + command: Extract< + OrchestrationV2InternalCommand, + { readonly type: "thread.background-work.settle" } + >, + events: Ref.Ref>, + ) => + Effect.gen(function* () { + const [projection, stopped] = yield* Effect.all([ + projectionStore.getThreadRecords( + command.threadId, + ["runs", "attempts", "turnItems", "providerThreads"], + { + turnItemTypes: ["command_execution", "dynamic_tool", "subagent"], + turnItemStatuses: ["pending", "running", "waiting"], + }, + ), + projectionStore.getProviderControlContext(command.threadId, { + providerThreadId: command.providerThreadId, + providerTurnId: command.providerTurnId, + }), + ]).pipe( + Effect.mapError( + (cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }), + ), + ); + const stoppedRunId = projection.attempts.find( + (attempt) => attempt.id === stopped.providerTurn?.runAttemptId, + )?.runId; + const stoppedRun = projection.runs.find((run) => run.id === stoppedRunId); + // A new turn may have started since Stop; its work is not this Stop's. + if ( + stoppedRun === undefined || + projection.runs.some( + (run) => + run.status === "preparing" || run.status === "starting" || run.status === "running", + ) + ) { + return; + } + yield* settleBackgroundWork({ + command, + events, + projection, + stoppedProviderThreadId: command.providerThreadId, + throughRunOrdinal: stoppedRun.ordinal, + now: yield* DateTime.now, + }); + }); + const dispatchRunInterrupt = ( command: Extract, events: Ref.Ref>, @@ -7547,9 +7695,11 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio "turnItems", "subagents", ], + // Thread-wide, like the Waiting strip: a settled thread's background + // work can belong to an earlier run than the one Stop targets. { - turnItemTypes: ["command_execution", "run_interrupt_request", "run_interrupt_result"], - turnItemRunId: command.runId, + turnItemTypes: ["command_execution", "dynamic_tool", "subagent"], + turnItemStatuses: ["pending", "running", "waiting"], }, ); const run = projection.runs.find((candidate) => candidate.id === command.runId); @@ -7765,6 +7915,39 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio cause: `Run ${command.runId} is not interruptible.`, }); } + // Stop on a settled thread's background work. Its process may be gone + // (released, restarted) and only the projection still shows the work; + // the settle follow-up ends whatever no provider reports ending. + const settleOnly = + providerTurn.status !== "running" && + (providerThread.providerSessionId === null || + Option.isNone( + yield* providerSessions + .get(providerThread.providerSessionId) + .pipe(Effect.orElseSucceed(() => Option.none())), + )); + if (settleOnly) { + yield* emitEvent({ + type: "turn-item.updated", + threadId: command.threadId, + runId: run.id, + nodeId: rootNode.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: interruptRequestItem, + }); + if (command.holdQueue === true) yield* holdQueuedRuns; + yield* stopCompletionCohort(); + yield* settleBackgroundWork({ + command, + events, + projection, + stoppedProviderThreadId: providerThread.id, + throughRunOrdinal: run.ordinal, + now, + }); + return undefined; + } if (providerThread.providerSessionId === null) { return yield* new OrchestratorDispatchError({ commandId: command.commandId, @@ -9052,6 +9235,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio case "checkpoint.rollback.fail": yield* dispatchCheckpointRollbackFail(command, events); break; + case "thread.background-work.settle": + yield* dispatchBackgroundWorkSettle(command, events); + break; case "thread.fork": yield* dispatchThreadFork(command, events); break; @@ -9145,9 +9331,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const plan = yield* dispatchOnce(command).pipe( Effect.flatMap((planned) => - // A Stop can race with terminal provider events. Its empty plan is an - // accepted idempotent outcome; every other command must still mutate. - planned.events.length > 0 + // A settle that finds the provider already ended everything has + // nothing to record, which is its expected outcome, not a failure. + planned.events.length > 0 || command.type === "thread.background-work.settle" ? Effect.succeed(planned) : Effect.fail( new OrchestratorDispatchError({ @@ -9194,6 +9380,33 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ); + if (plan.events.length === 0) { + // A settle that ended nothing still records its receipt: a replayed Stop + // effect then finds it instead of settling work that appeared since. + const resultSequence = yield* Effect.gen(function* () { + const sequence = yield* eventSink.latestSequence({ threadId: commandThreadId(command) }); + yield* commandReceipts.insertIfAbsent({ + commandId: command.commandId, + threadId: commandThreadId(command), + commandType: command.type, + acceptedAt: yield* DateTime.now, + resultSequence: sequence, + status: "accepted", + error: null, + }); + return sequence; + }).pipe( + Effect.mapError( + (cause) => + new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause, + }), + ), + ); + return { sequence: resultSequence, storedEvents: [] } satisfies OrchestratorV2DispatchResult; + } const acceptedAt = plan.events.at(-1)?.occurredAt ?? (yield* DateTime.now); const committed = yield* eventSink .commitCommand({ diff --git a/apps/server/src/orchestration-v2/ProjectionControlReads.test.ts b/apps/server/src/orchestration-v2/ProjectionControlReads.test.ts index f57c4ab7bc81..221c773149ac 100644 --- a/apps/server/src/orchestration-v2/ProjectionControlReads.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionControlReads.test.ts @@ -267,16 +267,10 @@ for (const storage of ["sqlite", "memory"] as const) { assert.instanceOf(missingThread, ProjectionStore.ProjectionStoreThreadNotFoundError); const calls: string[] = []; - let hasBackgroundWork = false; const sessions = Layer.mock(ProviderSessionManager.ProviderSessionManagerV2)({ get: () => Effect.succeed( Option.some({ - hasPendingBackgroundWorkForThread: (target: { id: ProviderThreadId }) => - Effect.sync(() => { - assert.equal(target.id, providerThreadId); - return hasBackgroundWork; - }), interruptTurn: () => Effect.sync(() => { calls.push("interrupt"); @@ -338,9 +332,8 @@ for (const storage of ["sqlite", "memory"] as const) { occurredAt: now, payload: { ...turn, status: "completed", completedAt: now }, }); - yield* control.interrupt({ threadId, providerThreadId, providerTurnId, providerSessionId }); - assert.lengthOf(calls, 3); - hasBackgroundWork = true; + // A settled turn's Stop still reaches the adapter, which alone knows + // whether it runs background work for the thread. yield* control.interrupt({ threadId, providerThreadId, providerTurnId, providerSessionId }); assert.deepEqual(calls, ["interrupt", "Use the smaller fix.", "reply", "interrupt"]); }).pipe( diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 334ae5d9f762..6718240d9e7a 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -279,6 +279,7 @@ export interface ProjectionRecordFilter { readonly turnItemRunIds?: ReadonlyArray; readonly runIds?: ReadonlyArray; readonly turnItemTypes?: ReadonlyArray; + readonly turnItemStatuses?: ReadonlyArray; } export type ProjectionRecordField = Exclude< keyof OrchestrationV2ThreadProjection, @@ -2575,6 +2576,7 @@ export const layer: Layer.Layer = FROM orchestration_v2_projection_turn_items WHERE thread_id = ${threadId} ${filter?.turnItemTypes === undefined ? sql`` : sql`AND type IN (SELECT value FROM json_each(${encodeIdList(filter.turnItemTypes)}))`} + ${filter?.turnItemStatuses === undefined ? sql`` : sql`AND status IN (SELECT value FROM json_each(${encodeIdList(filter.turnItemStatuses)}))`} ${filter?.turnItemRunId === undefined ? sql`` : sql`AND run_id = ${filter.turnItemRunId}`} ${filter?.turnItemRunIds === undefined ? sql`` : sql`AND (run_id IN (SELECT value FROM json_each(${encodeIdList(filter.turnItemRunIds.filter((id): id is RunId => id !== null))})) OR (${filter.turnItemRunIds.includes(null) ? 1 : 0} = 1 AND run_id IS NULL))`} ORDER BY ordinal ASC, turn_item_id ASC @@ -5740,6 +5742,8 @@ export const layerMemory: Layer.Layer = Layer.effect( (filter?.turnItemRunIds === undefined || filter.turnItemRunIds.includes(row.runId)) && (filter?.turnItemTypes === undefined || filter.turnItemTypes.includes(row.type)) && + (filter?.turnItemStatuses === undefined || + filter.turnItemStatuses.includes(row.status)) && (filter?.turnItemRunId === undefined || filter.turnItemRunId === row.runId), ), }; diff --git a/apps/server/src/orchestration-v2/ProviderTurnControlService.ts b/apps/server/src/orchestration-v2/ProviderTurnControlService.ts index 3a80e2af5d85..47f7aa80bdc4 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnControlService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnControlService.ts @@ -176,12 +176,10 @@ export const layer: Layer.Layer< ? loaded.session : yield* sessions.get(input.providerSessionId); if (Option.isNone(session)) return; - if ( - loaded.providerTurn.status !== "running" && - (session.value.hasPendingBackgroundWorkForThread === undefined || - !(yield* session.value.hasPendingBackgroundWorkForThread(loaded.providerThread))) - ) - return; + // A settled turn reaches its adapter too: only the adapter knows + // whether it still runs work for the thread, and each one either + // stops it or reports there is nothing left to stop. Background work + // the projection still shows is settled by the orchestrator after. yield* session.value.interruptTurn({ providerThread: loaded.providerThread, providerTurnId: loaded.providerTurn.id, diff --git a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts index e3381514fb9c..4cc834dddd62 100644 --- a/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts +++ b/apps/server/src/orchestration-v2/testkit/OrchestratorScenario.ts @@ -47,6 +47,11 @@ export type OrchestratorV2ScenarioStep = readonly type: "await_thread_idle"; readonly threadId: ThreadId; } + | { + /** Waits until the thread's Waiting strip lists no background work. */ + readonly type: "await_no_background_work"; + readonly threadId: ThreadId; + } | { readonly type: "await_run_steerable"; readonly threadId: ThreadId; @@ -348,6 +353,30 @@ export function runOrchestratorV2Scenario( return yield* waitForThreadIdle(threadId, attemptsRemaining - 1, deadlineAt); }); + const waitForNoBackgroundWork = ( + threadId: ThreadId, + attemptsRemaining = SCENARIO_WAIT_ATTEMPTS, + deadlineAt = scenarioWaitDeadline(), + ): Effect.Effect< + void, + Orchestrator.OrchestratorV2Error | OrchestratorV2ScenarioStepError, + never + > => + Effect.gen(function* () { + const pending = (yield* orchestrator.getThreadShell(threadId))?.pendingBackgroundTasks; + if ((pending?.length ?? 0) === 0) { + return; + } + if (scenarioWaitExhausted(attemptsRemaining, deadlineAt)) { + return yield* new OrchestratorV2ScenarioStepError({ + scenario: scenario.name, + step: `await_no_background_work:${threadId}:pending=${pending?.map((task) => task.taskId).join(",")}`, + }); + } + yield* yieldToRuntime; + return yield* waitForNoBackgroundWork(threadId, attemptsRemaining - 1, deadlineAt); + }); + const waitForRunSteerable = ( threadId: ThreadId, runId: OrchestrationV2Run["id"], @@ -623,6 +652,9 @@ export function runOrchestratorV2Scenario( case "await_thread_idle": yield* waitForThreadIdle(step.threadId); break; + case "await_no_background_work": + yield* waitForNoBackgroundWork(step.threadId); + break; case "await_run_steerable": yield* waitForRunSteerable(step.threadId, step.runId); break; diff --git a/apps/server/src/orchestration-v2/testkit/fixtures/index.ts b/apps/server/src/orchestration-v2/testkit/fixtures/index.ts index 0334f9358d44..8a34d2d97e4b 100644 --- a/apps/server/src/orchestration-v2/testkit/fixtures/index.ts +++ b/apps/server/src/orchestration-v2/testkit/fixtures/index.ts @@ -155,6 +155,11 @@ import { assertToolCallReadOnlyOnRequestOutput, } from "./tool_call_read_only_on_request/output.ts"; import { toolCallReadOnlyOnRequestInput } from "./tool_call_read_only_on_request/input.ts"; +import { + stopBackgroundWorkAfterFailedTurnInput, + stopBackgroundWorkAfterReleaseInput, +} from "./stop_background_work_after_failed_turn/input.ts"; +import { assertStopBackgroundWorkAfterFailedTurnOutput } from "./stop_background_work_after_failed_turn/output.ts"; import { assertToolCallRestrictedGranularClaudeOutput } from "./tool_call_restricted_granular/claude_output.ts"; import { assertToolCallRestrictedGranularOutput } from "./tool_call_restricted_granular/codex_output.ts"; import { toolCallRestrictedGranularInput } from "./tool_call_restricted_granular/input.ts"; @@ -1529,6 +1534,37 @@ export const ORCHESTRATOR_REPLAY_FIXTURES: ReadonlyArray thread.id === projection.thread.id); + assert.deepEqual( + before?.pendingBackgroundTasks?.map((task) => task.kind), + ["command"], + "the failed turn leaves its command on the Waiting strip", + ); + + const command = projection.turnItems.find((item) => item.type === "command_execution"); + assert.equal(command?.status, "interrupted"); + assert.isNotNull(command?.completedAt); + const shell = result.shellSnapshot.threads.find((thread) => thread.id === projection.thread.id); + assert.deepEqual(shell?.pendingBackgroundTasks ?? [], []); + // Stop records the request, and never rewrites how the turn itself ended. + assert.include( + projection.turnItems.map((item) => item.type), + "run_interrupt_request", + ); + assert.equal(projection.runs[0]?.status, "failed"); +} diff --git a/apps/server/src/orchestration-v2/testkit/fixtures/stop_background_work_after_failed_turn/registry_transcript.ndjson b/apps/server/src/orchestration-v2/testkit/fixtures/stop_background_work_after_failed_turn/registry_transcript.ndjson new file mode 100644 index 000000000000..0ac4f144cb7c --- /dev/null +++ b/apps/server/src/orchestration-v2/testkit/fixtures/stop_background_work_after_failed_turn/registry_transcript.ndjson @@ -0,0 +1,10 @@ +{"type":"transcript_start","provider":"acpRegistry","protocol":"acp.ndjson-jsonrpc","version":"1","scenario":"stop_background_work_after_failed_turn","metadata":{"generatedBy":"hand-ported: tool_call_read_only_on_request registry frames + grok_prompt_error session/prompt error frame","nativeSessionId":"grok-replay-session-1"}} +{"type":"expect_outbound","label":"initialize","frame":{"kind":"request","method":"initialize","params":{"protocolVersion":2,"info":"","capabilities":"","clientCapabilities":"","clientInfo":""}}} +{"type":"emit_inbound","label":"initialized","frame":{"kind":"response","method":"initialize","result":{"protocolVersion":1,"agentCapabilities":{"loadSession":true,"mcpCapabilities":{"http":true,"sse":true},"promptCapabilities":{"audio":false,"embeddedContext":true,"image":false}},"authMethods":[{"id":"replay","name":"Replay"}],"agentInfo":{"name":"grok-replay","version":"1"}}}} +{"type":"expect_outbound","label":"session.new","frame":{"kind":"request","method":"session/new","params":{"cwd":"","mcpServers":""}}} +{"type":"emit_inbound","label":"session.created","frame":{"kind":"response","method":"session/new","result":{"sessionId":"grok-replay-session-1","models":{"currentModelId":"grok-build","availableModels":[{"modelId":"grok-build","name":"Grok Build"}]}}}} +{"type":"expect_outbound","label":"session.prompt","frame":{"kind":"request","method":"session/prompt","params":{"sessionId":"grok-replay-session-1","prompt":[{"type":"text","text":"\n\n\n\n\nStart the dev server in the background with `bun run dev` and keep it running.\n"}]}}} +{"type":"emit_inbound","label":"tool.started","frame":{"kind":"notification","method":"session/update","params":{"sessionId":"grok-replay-session-1","update":{"sessionUpdate":"tool_call","toolCallId":"dev-server","title":"Start dev server","kind":"execute","status":"pending","rawInput":{"command":"bun run dev"}}}}} +{"type":"emit_inbound","label":"tool.running","frame":{"kind":"notification","method":"session/update","params":{"sessionId":"grok-replay-session-1","update":{"sessionUpdate":"tool_call_update","toolCallId":"dev-server","status":"in_progress"}}}} +{"type":"emit_inbound","label":"prompt.failed","frame":{"kind":"response","method":"session/prompt","error":{"code":-32603,"message":"Internal error","data":{"message":"API error (status 400 Bad Request): invalid_request_error: grok prompt error fixture","http_status":400}}}} +{"type":"runtime_exit","status":"success"} diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index cafc2bb19945..cb220a5b4ccd 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -2845,6 +2845,19 @@ const OrchestrationV2InternalCommand = Schema.Union([ requestId: CommandId, message: TrimmedNonEmptyString, }), + /** + * Follows a Stop once its provider returned: background work the settled + * thread still shows on that provider thread is no longer reported by any + * provider process, so it is marked interrupted. Only the stopped turn's run + * and older runs are settled; a later run's work is its own. + */ + Schema.Struct({ + type: Schema.Literal("thread.background-work.settle"), + commandId: CommandId, + threadId: ThreadId, + providerThreadId: ProviderThreadId, + providerTurnId: ProviderTurnId, + }), ]); export type OrchestrationV2InternalCommand = typeof OrchestrationV2InternalCommand.Type; diff --git a/packages/shared/src/orchestrationV2PendingBackgroundWork.ts b/packages/shared/src/orchestrationV2PendingBackgroundWork.ts index ee78736e5c88..074b5feca534 100644 --- a/packages/shared/src/orchestrationV2PendingBackgroundWork.ts +++ b/packages/shared/src/orchestrationV2PendingBackgroundWork.ts @@ -171,6 +171,29 @@ function nativeTaskIdFromTurnItem(item: PendingBackgroundWorkTurnItem): string { return String(item.id); } +/** + * The turn items the pending-work list names, without its settled-run gate. + * Stop ends exactly these, so the list and Stop cannot disagree. + */ +export function pendingBackgroundTurnItems(input: { + readonly turnItems: ReadonlyArray; + readonly runs?: ReadonlyArray; +}): ReadonlyArray { + const rolledBackRunIds = new Set( + (input.runs ?? []).filter((run) => run.status === "rolled_back").map((run) => String(run.id)), + ); + return input.turnItems.filter( + (item) => + BACKGROUND_TURN_ITEM_TYPES.has(item.type) && + isOrchestrationV2WorkActive(item.status) && + !(item.type === "dynamic_tool" && isPersistentDynamicToolInput(item.input)) && + // Null/absent run id stays eligible; only known rolled_back runs drop. + (item.runId === undefined || + item.runId === null || + !rolledBackRunIds.has(String(item.runId))), + ); +} + /** * Derive one normalized pending-background-work list for post-settlement UI. * @@ -212,9 +235,6 @@ export function derivePendingBackgroundWork(input: { } const byTaskId = new Map(); - const rolledBackRunIds = new Set( - (input.runs ?? []).filter((run) => run.status === "rolled_back").map((run) => String(run.id)), - ); const providerThreads = input.activeProviderThreadId === undefined || input.activeProviderThreadId === null @@ -235,22 +255,7 @@ export function derivePendingBackgroundWork(input: { } } - for (const item of input.turnItems) { - if (!BACKGROUND_TURN_ITEM_TYPES.has(item.type)) { - continue; - } - if (!isOrchestrationV2WorkActive(item.status)) { - continue; - } - if (item.type === "dynamic_tool" && isPersistentDynamicToolInput(item.input)) { - continue; - } - // Null/absent run id stays eligible; only known rolled_back runs drop. - const itemRunId = item.runId; - if (itemRunId !== undefined && itemRunId !== null && rolledBackRunIds.has(String(itemRunId))) { - continue; - } - + for (const item of pendingBackgroundTurnItems(input)) { const taskId = nativeTaskIdFromTurnItem(item); if (byTaskId.has(taskId)) { continue;