From 5981000974a75f8e53ea9e4e33a8a5e0e3ffe4d8 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 15:28:14 -0400 Subject: [PATCH 1/8] fix(server): runs no longer wedge after their provider session is released Signed-off-by: Yordis Prieto --- .../src/orchestration-v2/Orchestrator.ts | 95 +++++++++++- .../ProviderSessionManager.test.ts | 137 +++++++++++++++++- .../ProviderSessionManager.ts | 54 ++++--- .../src/orchestration-v2/runtimeLayer.test.ts | 104 +++++++++++++ 4 files changed, 362 insertions(+), 28 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index c58dd46609ec..f12352971962 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -7926,14 +7926,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio // 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())), - )); + const sessionIsDead = + providerThread.providerSessionId === null || + Option.isNone( + yield* providerSessions + .get(providerThread.providerSessionId) + .pipe(Effect.orElseSucceed(() => Option.none())), + ); + const settleOnly = providerTurn.status !== "running" && sessionIsDead; if (settleOnly) { yield* emitEvent({ type: "turn-item.updated", @@ -7956,6 +7956,85 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); return undefined; } + // The projection still shows this turn running, but its provider + // session is already gone (e.g. released on idle timeout): no live + // process will ever report a terminal for it. Settle the run the same + // way as an interrupt before provider start, plus the stuck turn + // itself, instead of erroring and leaving the run wedged forever. + if (providerTurn.status === "running" && sessionIsDead) { + const attempt = projection.attempts.find( + (candidate) => candidate.id === run.activeAttemptId, + ); + if (attempt === undefined) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: `Run ${command.runId} has no active attempt to interrupt.`, + }); + } + yield* emitEvent({ + type: "turn-item.updated", + threadId: command.threadId, + runId: run.id, + nodeId: rootNode.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: interruptRequestItem, + }); + yield* emitEvent({ + type: "provider-turn.updated", + threadId: command.threadId, + runId: run.id, + nodeId: providerTurn.nodeId, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...providerTurn, status: "interrupted", completedAt: now }, + }); + yield* emitEvent({ + type: "run-attempt.updated", + threadId: command.threadId, + runId: run.id, + nodeId: rootNode.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...attempt, status: "interrupted", completedAt: now }, + }); + yield* emitEvent({ + type: "node.updated", + threadId: command.threadId, + runId: run.id, + nodeId: rootNode.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...rootNode, status: "interrupted", completedAt: now }, + }); + yield* emitEvent({ + type: "run.updated", + threadId: command.threadId, + runId: run.id, + nodeId: rootNode.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...run, status: "interrupted", completedAt: now }, + }); + if (command.holdQueue === true) yield* holdQueuedRuns; + yield* stopCompletionCohort(); + yield* settleBackgroundWork({ + command, + events, + projection, + stoppedProviderThreadId: providerThread.id, + throughRunOrdinal: run.ordinal, + now, + }); + return { + effectTypes: ["provider-turn.start", "provider-turn.restart"], + reason: `Run ${run.id} was interrupted after its provider session ${providerThread.providerSessionId} was no longer active.`, + } satisfies { + readonly effectTypes: ReadonlyArray; + readonly reason: string; + }; + } if (providerThread.providerSessionId === null) { return yield* new OrchestratorDispatchError({ commandId: command.commandId, diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 63d4c57fb663..61afeff887bf 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -195,11 +195,13 @@ function makeProviderThread(input: { readonly threadId: ThreadId; readonly providerSessionId: ProviderSessionId; readonly now: DateTime.Utc; + readonly nativeThreadId?: string; }): OrchestrationV2ProviderThread { + const nativeThreadId = input.nativeThreadId ?? "native-thread"; return { id: input.idAllocator.derive.providerThread({ driver: CODEX_DRIVER, - nativeThreadId: "native-thread", + nativeThreadId, }), driver: CODEX_DRIVER, providerInstanceId: modelSelection.instanceId, @@ -208,7 +210,7 @@ function makeProviderThread(input: { ownerNodeId: null, nativeThreadRef: { driver: CODEX_DRIVER, - nativeId: "native-thread", + nativeId: nativeThreadId, strength: "strong", }, nativeConversationHeadRef: null, @@ -2561,12 +2563,14 @@ it.effect( threadId: firstThreadId, providerSessionId, now, + nativeThreadId: "native-thread-a", }); const secondProviderThread = makeProviderThread({ idAllocator, threadId: secondThreadId, providerSessionId, now, + nativeThreadId: "native-thread-b", }); const firstRunId = idAllocator.derive.run({ threadId: firstThreadId, ordinal: 1 }); const secondRunId = idAllocator.derive.run({ threadId: secondThreadId, ordinal: 1 }); @@ -2677,6 +2681,135 @@ it.effect( }), ); +it.effect( + "ProviderSessionManagerV2 ignores a stray turn.terminal for a provider thread that never started a turn", + () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const projectId = yield* idAllocator.allocate.project({ + fixtureName: "provider-session-manager-stray-terminal", + }); + const firstThreadId = yield* idAllocator.allocate.thread({ + fixtureName: "provider-session-manager-stray-terminal-a", + projectId, + }); + const secondThreadId = yield* idAllocator.allocate.thread({ + fixtureName: "provider-session-manager-stray-terminal-b", + projectId, + }); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId: firstThreadId, + }); + const firstProviderThread = makeProviderThread({ + idAllocator, + threadId: firstThreadId, + providerSessionId, + now, + nativeThreadId: "native-thread-a", + }); + const secondProviderThread = makeProviderThread({ + idAllocator, + threadId: secondThreadId, + providerSessionId, + now, + nativeThreadId: "native-thread-b", + }); + const firstRunId = idAllocator.derive.run({ threadId: firstThreadId, ordinal: 1 }); + const strayProviderTurnId = idAllocator.derive.providerTurn({ + driver: CODEX_DRIVER, + nativeTurnId: "native-turn-stray", + }); + + yield* eventSink.write({ + events: [ + yield* makeThreadCreatedEvent({ idAllocator, threadId: firstThreadId, now }), + yield* makeThreadCreatedEvent({ idAllocator, threadId: secondThreadId, now }), + ], + }); + const runtime = yield* manager.open({ + threadId: firstThreadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* manager.open({ + threadId: secondThreadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* runtime.events.pipe(Stream.runDrain, Effect.forkScoped); + const firstAppThread = (yield* projectionStore.getThreadProjection(firstThreadId)).thread; + // Only the first thread ever starts a turn. The second thread's + // provider thread id never calls markBusy on this session. + yield* runtime.startTurn({ + appThread: firstAppThread, + threadId: firstThreadId, + runId: firstRunId, + runOrdinal: 1, + providerTurnOrdinal: 1, + attemptId: idAllocator.derive.runAttempt({ runId: firstRunId, attemptOrdinal: 1 }), + rootNodeId: idAllocator.derive.rootNode({ runId: firstRunId }), + providerThread: firstProviderThread, + message: { + createdBy: "user", + creationSource: "web", + messageId: yield* idAllocator.allocate.message({ threadId: firstThreadId, ordinal: 1 }), + text: "first", + attachments: [], + }, + modelSelection, + runtimePolicy, + }); + + const queue = (yield* Ref.get(state)).eventQueues.get(String(providerSessionId)); + assert.isDefined(queue); + // A stray terminal arrives for the second (never-busy) provider + // thread, e.g. a turn cancelled moments after it started on another + // app thread sharing this session. With a session-wide busyCount + // this would zero it out and release the session even though the + // first thread's turn is still genuinely running. + yield* Queue.offer(queue!, { + type: "turn.terminal", + driver: CODEX_DRIVER, + providerThreadId: secondProviderThread.id, + providerTurnId: strayProviderTurnId, + runOrdinal: 1, + status: "cancelled", + failure: null, + threadDisposition: "reusable", + }); + yield* TestClock.adjust("2 seconds"); + yield* Effect.yieldNow; + assert.equal((yield* Ref.get(state)).closeCount, 0); + + // The same stray terminal arriving again is also a no-op. + yield* Queue.offer(queue!, { + type: "turn.terminal", + driver: CODEX_DRIVER, + providerThreadId: secondProviderThread.id, + providerTurnId: strayProviderTurnId, + runOrdinal: 1, + status: "cancelled", + failure: null, + threadDisposition: "reusable", + }); + yield* TestClock.adjust("2 seconds"); + yield* Effect.yieldNow; + assert.equal((yield* Ref.get(state)).closeCount, 0); + }); + + yield* effect.pipe(Effect.provide(makeTestLayer({ state, idleTimeoutMs: 1000 }))); + }), +); + it.effect( "ProviderSessionManagerV2 opens one shared runtime, broadcasts events, and detaches threads independently", () => diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 786eb5a6bf0a..36a9d97aeae8 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -7,6 +7,7 @@ import { OrchestrationV2RuntimeRequest, ProviderInstanceId, ProviderSessionId, + type ProviderThreadId, ThreadId, } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; @@ -204,7 +205,13 @@ interface LiveSessionEntry { readonly requestEventPermit: Semaphore.Semaphore; readonly scope: Scope.Closeable; readonly idleGeneration: number; - readonly busyCount: number; + /** + * Provider threads with a turn in flight, keyed by provider thread id + * rather than counted: a stray or duplicate `turn.terminal` for a thread + * that never started a turn here (or already settled one) is then a no-op + * instead of zeroing out another thread's genuinely running turn. + */ + readonly busyProviderThreadIds: ReadonlySet; readonly lastActivityAtMs: number; readonly idleFiber: Fiber.Fiber | null; /** Set when idle release is deferred for pending background work; bounds total deferral. */ @@ -725,7 +732,8 @@ export const layerWithOptions = ( } if ( input.onlyIfIdleGeneration !== undefined && - (existing.busyCount > 0 || existing.idleGeneration !== input.onlyIfIdleGeneration) + (existing.busyProviderThreadIds.size > 0 || + existing.idleGeneration !== input.onlyIfIdleGeneration) ) { return [Option.none(), current] as const; } @@ -871,7 +879,7 @@ export const layerWithOptions = ( const entry = current.get(key); if ( entry === undefined || - entry.busyCount > 0 || + entry.busyProviderThreadIds.size > 0 || entry.idleGeneration !== input.generation ) { return; @@ -893,7 +901,7 @@ export const layerWithOptions = ( const latestEntry = latest.get(key); if ( latestEntry === undefined || - latestEntry.busyCount > 0 || + latestEntry.busyProviderThreadIds.size > 0 || latestEntry.idleGeneration !== input.generation || latestEntry.runtime !== probedRuntime ) { @@ -925,8 +933,8 @@ export const layerWithOptions = ( } // hasPendingBackgroundWork yields to the adapter, so the idle // decision above can go stale; the generation guard revalidates - // busyCount and idleGeneration inside releaseEntry's atomic - // entry removal. + // busyProviderThreadIds and idleGeneration inside releaseEntry's + // atomic entry removal. yield* releaseEntry({ providerSessionId: input.providerSessionId, reason: "idle_timeout", @@ -962,7 +970,7 @@ export const layerWithOptions = ( const key = sessionKey(providerSessionId); const current = yield* Ref.get(sessions); const entry = current.get(key); - if (entry === undefined || entry.busyCount > 0) { + if (entry === undefined || entry.busyProviderThreadIds.size > 0) { return; } @@ -975,7 +983,7 @@ export const layerWithOptions = ( const lastActivityAtMs = yield* Clock.currentTimeMillis; yield* Ref.update(sessions, (latest) => { const latestEntry = latest.get(key); - if (latestEntry === undefined || latestEntry.busyCount > 0) { + if (latestEntry === undefined || latestEntry.busyProviderThreadIds.size > 0) { return latest; } const updated = new Map(latest); @@ -1160,7 +1168,7 @@ export const layerWithOptions = ( ); }); - const markBusy = (providerSessionId: ProviderSessionId) => + const markBusy = (providerSessionId: ProviderSessionId, providerThreadId: ProviderThreadId) => withActivityError( providerSessionId, Effect.gen(function* () { @@ -1172,9 +1180,11 @@ export const layerWithOptions = ( return [null, current] as const; } const updated = new Map(current); + const busyProviderThreadIds = new Set(entry.busyProviderThreadIds); + busyProviderThreadIds.add(providerThreadId); updated.set(key, { ...entry, - busyCount: entry.busyCount + 1, + busyProviderThreadIds, idleFiber: null, lastActivityAtMs: now, pinnedSinceMs: null, @@ -1185,7 +1195,7 @@ export const layerWithOptions = ( }), ); - const markIdle = (providerSessionId: ProviderSessionId) => + const markIdle = (providerSessionId: ProviderSessionId, providerThreadId: ProviderThreadId) => withActivityError( providerSessionId, Effect.gen(function* () { @@ -1197,9 +1207,11 @@ export const layerWithOptions = ( return current; } const updated = new Map(current); + const busyProviderThreadIds = new Set(entry.busyProviderThreadIds); + busyProviderThreadIds.delete(providerThreadId); updated.set(key, { ...entry, - busyCount: Math.max(0, entry.busyCount - 1), + busyProviderThreadIds, lastActivityAtMs: now, }); return updated; @@ -1373,12 +1385,18 @@ export const layerWithOptions = ( providerInstanceId: runtime.instanceId, }), ).pipe( - Effect.andThen(observeActivity(providerSessionId, markBusy(providerSessionId))), + Effect.andThen( + observeActivity( + providerSessionId, + markBusy(providerSessionId, input.providerThread.id), + ), + ), Effect.andThen(runtime.startTurn(input)), Effect.catch((error) => - observeActivity(providerSessionId, markIdle(providerSessionId)).pipe( - Effect.andThen(Effect.fail(error)), - ), + observeActivity( + providerSessionId, + markIdle(providerSessionId, input.providerThread.id), + ).pipe(Effect.andThen(Effect.fail(error))), ), ), steerTurn: (input) => @@ -1435,7 +1453,7 @@ export const layerWithOptions = ( return observeActivity( entry.runtime.providerSessionId, event.type === "turn.terminal" - ? markIdle(entry.runtime.providerSessionId) + ? markIdle(entry.runtime.providerSessionId, event.providerThreadId) : touchActivity(entry.runtime.providerSessionId), ).pipe( Effect.andThen( @@ -1678,7 +1696,7 @@ export const layerWithOptions = ( requestEventPermit: yield* Semaphore.make(1), scope: sessionScope, idleGeneration: 0, - busyCount: 0, + busyProviderThreadIds: new Set(), lastActivityAtMs: now, idleFiber: null, pinnedSinceMs: null, diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 2fdd2a1f7f35..4768f949dfd1 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -2995,6 +2995,110 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }), ); + it.effect( + "settles a run left running after its provider session was released out from under it", + () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const eventSink = yield* EventSink.EventSinkV2; + const threadId = ThreadId.make("runtime-layer-interrupt-dead-session"); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${threadId}:create`), + threadId, + projectId: ProjectId.make(`${threadId}:project`), + title: "Interrupt dead session", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: process.cwd(), + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${threadId}:message:0`), + threadId, + messageId: MessageId.make(`${threadId}:message:0`), + text: "Active", + attachments: [], + dispatchMode: { type: "start_immediately" }, + }); + const before = yield* orchestrator.getThreadProjection(threadId); + const activeRun = before.runs[0]!; + const providerThread = before.providerThreads[0]!; + const now = yield* DateTime.now; + const providerTurn = { + id: ProviderTurnId.make(`${threadId}:turn`), + providerThreadId: providerThread.id, + nodeId: activeRun.rootNodeId!, + runAttemptId: activeRun.activeAttemptId, + nativeTurnRef: null, + ordinal: 1, + status: "running" as const, + startedAt: now, + completedAt: null, + }; + // Drive the run and its provider turn into "running" the same way a + // real provider session reaching that state would, without actually + // opening one in ProviderSessionManagerV2. No server restart happens + // between a run starting and its shared session later getting + // released on idle timeout, so by the time interrupt is dispatched + // the projection still says "running" while the live session is + // already gone: `sessions.get` for it returns none, exactly as it + // would after a real release. + yield* eventSink.write({ + commandId: CommandId.make(`${threadId}:force-running`), + events: [ + { + id: EventId.make(`${threadId}:run-running`), + type: "run.updated", + threadId, + runId: activeRun.id, + occurredAt: now, + payload: { ...activeRun, status: "running", startedAt: now }, + }, + { + id: EventId.make(`${threadId}:turn-running`), + type: "provider-turn.updated", + threadId, + runId: activeRun.id, + occurredAt: now, + payload: providerTurn, + }, + ], + }); + + yield* orchestrator.dispatch({ + type: "run.interrupt", + commandId: CommandId.make(`${threadId}:interrupt`), + threadId, + runId: activeRun.id, + }); + + const after = yield* orchestrator.getThreadProjection(threadId); + const interruptedRun = after.runs.find((run) => run.id === activeRun.id); + assert.equal(interruptedRun?.status, "interrupted"); + assert.equal( + after.providerTurns.find( + (providerTurn) => providerTurn.runAttemptId === activeRun.activeAttemptId, + )?.status, + "interrupted", + ); + assert.equal( + after.attempts.find((attempt) => attempt.id === activeRun.activeAttemptId)?.status, + "interrupted", + ); + assert.equal( + after.nodes.find((node) => node.id === activeRun.rootNodeId)?.status, + "interrupted", + ); + }), + ); + for (const trigger of ["startup", "shutdown"] as const) { it.effect( `preserves and holds queued messages across ${trigger} until explicitly resumed`, From 354d5d12a2274d7c92fb1a6ba9ca84fac8e08e86 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 16:03:35 -0400 Subject: [PATCH 2/8] fix(server): dead-session interrupts finalize in-flight work and mark the run interrupted Signed-off-by: Yordis Prieto --- .../src/orchestration-v2/Orchestrator.ts | 105 ++++++++++++++++- .../orchestration-v2/RunExecutionService.ts | 2 +- .../src/orchestration-v2/runtimeLayer.test.ts | 110 ++++++++++++++++++ 3 files changed, 215 insertions(+), 2 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index f12352971962..4330ec5736fe 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -39,6 +39,7 @@ import { type OrchestrationV2Subagent, type OrchestrationV2ThreadProjection, type OrchestrationV2TurnItem, + isOrchestrationV2WorkActive, ProviderInstanceId, type ProviderSessionId, RunId, @@ -99,6 +100,7 @@ import { makeProviderFailure } from "./ProviderFailure.ts"; import { ProviderSessionManagerV2 } from "./ProviderSessionManager.ts"; import { ProviderSwitchServiceV2 } from "./ProviderSwitchService.ts"; import { isAutomaticCompletionRun, queuedRunsInDeliveryOrder } from "./QueuedRunOrder.ts"; +import { makeInterruptResultTurnItem } from "./RunExecutionService.ts"; import { RuntimePolicyV2 } from "./RuntimePolicy.ts"; import { makeSubagentChildThread, @@ -392,6 +394,15 @@ function pendingThreadTitleGenerationEffect( const WORKSPACE_PREPARATION_INPUT = "Preparing workspace"; +// Turn item types a live provider process owns and reports the end of itself +// (RunExecutionService/ProviderRuntimeRecoveryService background-work sweeps). +// Interrupt finalization must not re-settle these outside that sweep. +const BACKGROUND_CAPABLE_TURN_ITEM_TYPES: ReadonlySet = new Set([ + "command_execution", + "dynamic_tool", + "subagent", +]); + function isBlockingRun(run: OrchestrationV2Run): boolean { return ( run.status === "preparing" || @@ -7706,7 +7717,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio // 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", "dynamic_tool", "subagent"], + turnItemTypes: [...BACKGROUND_CAPABLE_TURN_ITEM_TYPES], turnItemStatuses: ["pending", "running", "waiting"], }, ); @@ -8017,6 +8028,98 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio occurredAt: now, payload: { ...run, status: "interrupted", completedAt: now }, }); + // No live process remains to settle the rest of the turn's in-flight + // work (settleBackgroundWork below only closes the background-capable + // roster). Finalize everything else the same way a dead-process + // reconciliation sweep would (ProviderRuntimeRecoveryService), so the + // thread does not keep showing open items, streaming text, or running + // subagents once the run itself says interrupted. request/approval + // items keep their runtime-request lifecycle untouched, same as + // elsewhere in this dispatch. + const residualTurnItems = yield* loadProjectionForCommand(command, ["turnItems"], { + turnItemRunId: run.id, + turnItemStatuses: ["pending", "running", "waiting"], + }); + for (const item of residualTurnItems.turnItems) { + if (BACKGROUND_CAPABLE_TURN_ITEM_TYPES.has(item.type) || "requestId" in item) continue; + yield* emitEvent({ + type: "turn-item.updated", + threadId: command.threadId, + runId: run.id, + ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { + ...item, + status: "interrupted", + completedAt: now, + updatedAt: now, + ...("streaming" in item ? { streaming: false } : {}), + }, + }); + } + for (const node of projection.nodes.filter( + (candidate) => + candidate.id !== rootNode.id && + candidate.runId === run.id && + isOrchestrationV2WorkActive(candidate.status), + )) { + yield* emitEvent({ + type: "node.updated", + threadId: command.threadId, + runId: run.id, + nodeId: node.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...node, status: "interrupted", completedAt: now }, + }); + } + for (const message of projection.messages.filter( + (candidate) => candidate.runId === run.id && candidate.streaming, + )) { + yield* emitEvent({ + type: "message.updated", + threadId: command.threadId, + runId: run.id, + ...(message.nodeId === null ? {} : { nodeId: message.nodeId }), + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...message, streaming: false, updatedAt: now }, + }); + } + for (const subagent of projection.subagents.filter( + (candidate) => + candidate.runId === run.id && isOrchestrationV2WorkActive(candidate.status), + )) { + yield* emitEvent({ + type: "subagent.updated", + threadId: command.threadId, + runId: run.id, + nodeId: subagent.id, + driver: subagent.driver, + providerInstanceId: subagent.providerInstanceId, + occurredAt: now, + payload: { ...subagent, status: "interrupted", completedAt: now, updatedAt: now }, + }); + } + // Same marker a live provider's turn.terminal interrupted would leave + // behind (RunExecutionService.writeFinalRunEvents); the web timeline + // keys the "Run interrupted" label off this item. + yield* emitEvent({ + type: "turn-item.updated", + threadId: command.threadId, + runId: run.id, + nodeId: rootNode.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: makeInterruptResultTurnItem({ + idAllocator, + run, + rootNode, + providerThread, + completedAt: now, + }), + }); if (command.holdQueue === true) yield* holdQueuedRuns; yield* stopCompletionCohort(); yield* settleBackgroundWork({ diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index 71211c126545..aab019166416 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -1433,7 +1433,7 @@ export const layer: Layer.Layer< }), ); -function makeInterruptResultTurnItem(input: { +export function makeInterruptResultTurnItem(input: { readonly idAllocator: IdAllocator.IdAllocatorV2Shape; readonly run: OrchestrationV2Run; readonly rootNode: OrchestrationV2ExecutionNode; diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 4768f949dfd1..7659788b6399 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -3042,6 +3042,9 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { startedAt: now, completedAt: null, }; + const streamingMessageId = MessageId.make(`${threadId}:assistant-message`); + const streamingItemId = TurnItemId.make(`${threadId}:assistant-item`); + const subagentId = NodeId.make(`${threadId}:subagent`); // Drive the run and its provider turn into "running" the same way a // real provider session reaching that state would, without actually // opening one in ProviderSessionManagerV2. No server restart happens @@ -3050,6 +3053,12 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { // the projection still says "running" while the live session is // already gone: `sessions.get` for it returns none, exactly as it // would after a real release. + // + // Also seed an in-flight turn item, a streaming message, and a + // running subagent on the same run, mirroring the in-flight work a + // real provider turn would leave open. None of these are reported + // terminal by a live process, so interrupt finalization must settle + // them itself. yield* eventSink.write({ commandId: CommandId.make(`${threadId}:force-running`), events: [ @@ -3069,6 +3078,87 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { occurredAt: now, payload: providerTurn, }, + { + id: EventId.make(`${threadId}:streaming-message`), + type: "message.updated", + threadId, + runId: activeRun.id, + nodeId: activeRun.rootNodeId!, + occurredAt: now, + payload: { + createdBy: "agent", + creationSource: "provider", + id: streamingMessageId, + threadId, + runId: activeRun.id, + nodeId: activeRun.rootNodeId!, + role: "assistant", + text: "Working on it", + attachments: [], + streaming: true, + createdAt: now, + updatedAt: now, + }, + }, + { + id: EventId.make(`${threadId}:streaming-item`), + type: "turn-item.updated", + threadId, + runId: activeRun.id, + nodeId: activeRun.rootNodeId!, + occurredAt: now, + payload: { + id: streamingItemId, + type: "assistant_message", + threadId, + runId: activeRun.id, + nodeId: activeRun.rootNodeId!, + providerThreadId: providerThread.id, + providerTurnId: providerTurn.id, + nativeItemRef: null, + parentItemId: null, + ordinal: 0, + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + messageId: streamingMessageId, + text: "Working on it", + streaming: true, + }, + }, + { + id: EventId.make(`${threadId}:subagent-running`), + type: "subagent.updated", + threadId, + runId: activeRun.id, + nodeId: subagentId, + driver: providerThread.driver, + providerInstanceId: providerThread.providerInstanceId, + occurredAt: now, + payload: { + id: subagentId, + threadId, + runId: activeRun.id, + parentNodeId: activeRun.rootNodeId!, + origin: "app_owned", + createdBy: "agent", + driver: providerThread.driver, + providerInstanceId: providerThread.providerInstanceId, + providerThreadId: null, + childThreadId: null, + nativeTaskRef: null, + prompt: "Do the subtask", + title: null, + model: null, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, ], }); @@ -3096,6 +3186,26 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { after.nodes.find((node) => node.id === activeRun.rootNodeId)?.status, "interrupted", ); + const streamingItem = after.turnItems.find((item) => item.id === streamingItemId); + assert.equal(streamingItem?.status, "interrupted"); + assert.equal( + streamingItem !== undefined && "streaming" in streamingItem + ? streamingItem.streaming + : undefined, + false, + ); + assert.equal( + after.messages.find((message) => message.id === streamingMessageId)?.streaming, + false, + ); + assert.equal( + after.subagents.find((subagent) => subagent.id === subagentId)?.status, + "interrupted", + ); + const interruptResult = after.turnItems.find( + (item) => item.type === "run_interrupt_result" && item.runId === activeRun.id, + ); + assert.isDefined(interruptResult); }), ); From 4c6ed668a8a09ba63042b81a9a287daedeb5b74e Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 16:09:50 -0400 Subject: [PATCH 3/8] fix(server): a failed overlapping turn start no longer idles a running turn Signed-off-by: Yordis Prieto --- .../ProviderSessionManager.test.ts | 96 ++++++++++++++++++- .../ProviderSessionManager.ts | 42 +++++--- 2 files changed, 125 insertions(+), 13 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 61afeff887bf..55360e2f0d77 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -248,6 +248,7 @@ function makeProviderAdapter( readonly hasPendingBackgroundWork?: Effect.Effect; readonly hangSessionScopeClose?: boolean; readonly beforeUnload?: Effect.Effect; + readonly startTurn?: ProviderAdapterV2SessionRuntime["startTurn"]; } = {}, ): ProviderAdapterV2Shape { return { @@ -318,7 +319,7 @@ function makeProviderAdapter( ...current, resumeCount: current.resumeCount + 1, })).pipe(Effect.as(threadInput.providerThread)), - startTurn: () => Effect.void, + startTurn: options.startTurn ?? (() => Effect.void), steerTurn: () => Effect.void, interruptTurn: () => Ref.update(state, (current) => ({ @@ -363,6 +364,7 @@ function makeTestLayer(input: { readonly hasPendingBackgroundWork?: Effect.Effect; readonly hangSessionScopeClose?: boolean; readonly beforeUnload?: Effect.Effect; + readonly startTurn?: ProviderAdapterV2SessionRuntime["startTurn"]; readonly serverSettingsLayer?: ReturnType; readonly projectServiceLayer?: Layer.Layer; }) { @@ -382,6 +384,7 @@ function makeTestLayer(input: { ? {} : { hangSessionScopeClose: input.hangSessionScopeClose }), ...(input.beforeUnload === undefined ? {} : { beforeUnload: input.beforeUnload }), + ...(input.startTurn === undefined ? {} : { startTurn: input.startTurn }), }), ); const providerEventIngestorTestLayer = ProviderEventIngestor.layer.pipe( @@ -2810,6 +2813,97 @@ it.effect( }), ); +it.effect( + "ProviderSessionManagerV2 keeps a running turn busy when an overlapping startTurn on the same provider thread fails", + () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const startTurnCalls = yield* Ref.make(0); + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const projectId = yield* idAllocator.allocate.project({ + fixtureName: "provider-session-manager-overlapping-start", + }); + const threadId = yield* idAllocator.allocate.thread({ + fixtureName: "provider-session-manager-overlapping-start", + projectId, + }); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + const providerThread = makeProviderThread({ + idAllocator, + threadId, + providerSessionId, + now, + nativeThreadId: "native-thread-overlap", + }); + yield* eventSink.write({ + events: [yield* makeThreadCreatedEvent({ idAllocator, threadId, now })], + }); + const runtime = yield* manager.open({ + threadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* runtime.events.pipe(Stream.runDrain, Effect.forkScoped); + const appThread = (yield* projectionStore.getThreadProjection(threadId)).thread; + const startInput = (ordinal: number) => + Effect.gen(function* () { + const runId = idAllocator.derive.run({ threadId, ordinal }); + return { + appThread, + threadId, + runId, + runOrdinal: ordinal, + providerTurnOrdinal: ordinal, + attemptId: idAllocator.derive.runAttempt({ runId, attemptOrdinal: 1 }), + rootNodeId: idAllocator.derive.rootNode({ runId }), + providerThread, + message: { + createdBy: "user" as const, + creationSource: "web" as const, + messageId: yield* idAllocator.allocate.message({ threadId, ordinal }), + text: `turn ${ordinal}`, + attachments: [], + }, + modelSelection, + runtimePolicy, + }; + }); + + yield* runtime.startTurn(yield* startInput(1)); + const overlapping = yield* runtime.startTurn(yield* startInput(2)).pipe(Effect.exit); + assert.isTrue(overlapping._tag === "Failure"); + + yield* TestClock.adjust("2 seconds"); + yield* Effect.yieldNow; + assert.equal((yield* Ref.get(state)).closeCount, 0); + }); + + yield* effect.pipe( + Effect.provide( + makeTestLayer({ + state, + idleTimeoutMs: 1000, + startTurn: () => + Ref.getAndUpdate(startTurnCalls, (calls) => calls + 1).pipe( + Effect.flatMap((calls) => + calls === 0 ? Effect.void : unimplemented("overlapping turn rejected"), + ), + ), + }), + ), + ); + }), +); + it.effect( "ProviderSessionManagerV2 opens one shared runtime, broadcasts events, and detaches threads independently", () => diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 36a9d97aeae8..24127fb9bfea 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -1174,10 +1174,10 @@ export const layerWithOptions = ( Effect.gen(function* () { const key = sessionKey(providerSessionId); const now = yield* Clock.currentTimeMillis; - const idleFiber = yield* Ref.modify(sessions, (current) => { + const [idleFiber, added] = yield* Ref.modify(sessions, (current) => { const entry = current.get(key); if (entry === undefined) { - return [null, current] as const; + return [[null, false] as const, current] as const; } const updated = new Map(current); const busyProviderThreadIds = new Set(entry.busyProviderThreadIds); @@ -1189,9 +1189,13 @@ export const layerWithOptions = ( lastActivityAtMs: now, pinnedSinceMs: null, }); - return [entry.idleFiber, updated] as const; + return [ + [entry.idleFiber, !entry.busyProviderThreadIds.has(providerThreadId)] as const, + updated, + ] as const; }); yield* cancelIdleFiber(idleFiber); + return added; }), ); @@ -1386,17 +1390,31 @@ export const layerWithOptions = ( }), ).pipe( Effect.andThen( - observeActivity( - providerSessionId, - markBusy(providerSessionId, input.providerThread.id), + markBusy(providerSessionId, input.providerThread.id).pipe( + Effect.catchCause((cause) => + Effect.logWarning("orchestration-v2.driver-session.activity-failed", { + providerSessionId, + cause, + }).pipe(Effect.as(false)), + ), ), ), - Effect.andThen(runtime.startTurn(input)), - Effect.catch((error) => - observeActivity( - providerSessionId, - markIdle(providerSessionId, input.providerThread.id), - ).pipe(Effect.andThen(Effect.fail(error))), + // Only the startTurn that marked the thread busy may clear it on + // failure; an overlapping attempt must not idle a running turn. + Effect.flatMap((markedBusy) => + runtime + .startTurn(input) + .pipe( + Effect.catch((error) => + (markedBusy + ? observeActivity( + providerSessionId, + markIdle(providerSessionId, input.providerThread.id), + ) + : Effect.void + ).pipe(Effect.andThen(Effect.fail(error))), + ), + ), ), ), steerTurn: (input) => From 7d6d7525ce63b0f1074ad3f3732701ee5fe754c4 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 17:01:16 -0400 Subject: [PATCH 4/8] fix(server): dead-session interrupts reconcile like a lost provider process Signed-off-by: Yordis Prieto --- .../src/orchestration-v2/Orchestrator.ts | 249 ++-- .../ProviderRuntimeRecoveryService.ts | 1090 ++++++++++------- .../orchestration-v2/RunExecutionService.ts | 2 +- .../src/orchestration-v2/runtimeLayer.test.ts | 156 ++- 4 files changed, 886 insertions(+), 611 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 4330ec5736fe..088c70f442b8 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -39,7 +39,6 @@ import { type OrchestrationV2Subagent, type OrchestrationV2ThreadProjection, type OrchestrationV2TurnItem, - isOrchestrationV2WorkActive, ProviderInstanceId, type ProviderSessionId, RunId, @@ -75,7 +74,12 @@ import { notificationTurnItem } from "./Notification.ts"; import { isRestartNoteSource } from "./RestartBackgroundNote.ts"; import { isUndeliveredMailboxSteer } from "./NotificationMailbox.ts"; import { EventSinkV2 } from "./EventSink.ts"; -import type { OrchestrationEffectRequestV2, PendingOrchestrationEffectV2 } from "./EffectOutbox.ts"; +import { + EffectOutboxV2, + PROCESS_BOUND_EFFECT_TYPES, + type OrchestrationEffectRequestV2, + type PendingOrchestrationEffectV2, +} from "./EffectOutbox.ts"; import { IdAllocatorV2 } from "./IdAllocator.ts"; import { ThreadCommandExecutor, @@ -97,10 +101,13 @@ import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import { ProviderAdapterRegistryV2 } from "./ProviderAdapterRegistry.ts"; import { ProviderContinuationRequests } from "./ProviderContinuationRequests.ts"; import { makeProviderFailure } from "./ProviderFailure.ts"; +import { + planProcessLossReconciliation, + planThreadReconciliation, +} from "./ProviderRuntimeRecoveryService.ts"; import { ProviderSessionManagerV2 } from "./ProviderSessionManager.ts"; import { ProviderSwitchServiceV2 } from "./ProviderSwitchService.ts"; import { isAutomaticCompletionRun, queuedRunsInDeliveryOrder } from "./QueuedRunOrder.ts"; -import { makeInterruptResultTurnItem } from "./RunExecutionService.ts"; import { RuntimePolicyV2 } from "./RuntimePolicy.ts"; import { makeSubagentChildThread, @@ -662,6 +669,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const eventSink = yield* EventSinkV2; const commandReceipts = yield* CommandReceiptStoreV2; const idAllocator = yield* IdAllocatorV2; + const outbox = yield* EffectOutboxV2; const projects = yield* ProjectStore.ProjectStoreV2; const projectionStore = yield* ProjectionStoreV2; const nextTurnItemOrdinal = ( @@ -7973,10 +7981,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio // way as an interrupt before provider start, plus the stuck turn // itself, instead of erroring and leaving the run wedged forever. if (providerTurn.status === "running" && sessionIsDead) { - const attempt = projection.attempts.find( - (candidate) => candidate.id === run.activeAttemptId, - ); - if (attempt === undefined) { + if (!projection.attempts.some((candidate) => candidate.id === run.activeAttemptId)) { return yield* new OrchestratorDispatchError({ commandId: command.commandId, commandType: command.type, @@ -7992,134 +7997,111 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio occurredAt: now, payload: interruptRequestItem, }); - yield* emitEvent({ - type: "provider-turn.updated", - threadId: command.threadId, - runId: run.id, - nodeId: providerTurn.nodeId, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...providerTurn, status: "interrupted", completedAt: now }, - }); - yield* emitEvent({ - type: "run-attempt.updated", - threadId: command.threadId, - runId: run.id, - nodeId: rootNode.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...attempt, status: "interrupted", completedAt: now }, - }); - yield* emitEvent({ - type: "node.updated", - threadId: command.threadId, - runId: run.id, - nodeId: rootNode.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...rootNode, status: "interrupted", completedAt: now }, - }); - yield* emitEvent({ - type: "run.updated", - threadId: command.threadId, + // No live process remains to settle this run, so finalize it exactly + // the way ProviderRuntimeRecoveryService reconciles a dead provider + // process: scoped to this run and its own provider thread, so sibling + // live provider threads on the same orchestration thread (e.g. a + // concurrently-running native subagent) are left untouched. + const recoveryProjection = yield* projectionStore + .getRuntimeRecoveryProjection(command.threadId) + .pipe( + Effect.mapError( + (cause) => new OrchestratorProjectionError({ threadId: command.threadId, cause }), + ), + ); + const plan = yield* planProcessLossReconciliation({ + projection: recoveryProjection, runId: run.id, - nodeId: rootNode.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...run, status: "interrupted", completedAt: now }, - }); - // No live process remains to settle the rest of the turn's in-flight - // work (settleBackgroundWork below only closes the background-capable - // roster). Finalize everything else the same way a dead-process - // reconciliation sweep would (ProviderRuntimeRecoveryService), so the - // thread does not keep showing open items, streaming text, or running - // subagents once the run itself says interrupted. request/approval - // items keep their runtime-request lifecycle untouched, same as - // elsewhere in this dispatch. - const residualTurnItems = yield* loadProjectionForCommand(command, ["turnItems"], { - turnItemRunId: run.id, - turnItemStatuses: ["pending", "running", "waiting"], - }); - for (const item of residualTurnItems.turnItems) { - if (BACKGROUND_CAPABLE_TURN_ITEM_TYPES.has(item.type) || "requestId" in item) continue; - yield* emitEvent({ - type: "turn-item.updated", - threadId: command.threadId, - runId: run.id, - ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { - ...item, - status: "interrupted", - completedAt: now, - updatedAt: now, - ...("streaming" in item ? { streaming: false } : {}), + providerThreadId: providerThread.id, + commandId: command.commandId, + now, + ids: idAllocator, + outbox, + }).pipe( + Effect.mapError( + (cause) => + new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause, + }), + ), + ); + yield* Ref.update(events, (existing) => [...existing, ...plan.events]); + yield* Ref.update(effects, (existing) => [...existing, ...plan.effects]); + // The workspace may have changed before the session died, so the run + // still gets its end-of-turn checkpoint like any interrupted turn. + const checkpointScopeId = rootNode.checkpointScopeId; + if (checkpointScopeId !== null) { + yield* Ref.update(effects, (existing) => [ + ...existing, + { + id: `effect:checkpoint.capture:${run.id}`, + commandId: CommandId.make(`command:effect:checkpoint.capture:${run.id}`), + threadId: run.threadId, + request: { + type: "checkpoint.capture" as const, + runId: run.id, + scopeId: checkpointScopeId, + }, }, - }); + ]); } - for (const node of projection.nodes.filter( - (candidate) => - candidate.id !== rootNode.id && - candidate.runId === run.id && - isOrchestrationV2WorkActive(candidate.status), - )) { - yield* emitEvent({ - type: "node.updated", - threadId: command.threadId, - runId: run.id, - nodeId: node.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...node, status: "interrupted", completedAt: now }, - }); - } - for (const message of projection.messages.filter( - (candidate) => candidate.runId === run.id && candidate.streaming, - )) { - yield* emitEvent({ - type: "message.updated", - threadId: command.threadId, - runId: run.id, - ...(message.nodeId === null ? {} : { nodeId: message.nodeId }), - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...message, streaming: false, updatedAt: now }, - }); - } - for (const subagent of projection.subagents.filter( - (candidate) => - candidate.runId === run.id && isOrchestrationV2WorkActive(candidate.status), - )) { - yield* emitEvent({ - type: "subagent.updated", - threadId: command.threadId, - runId: run.id, - nodeId: subagent.id, - driver: subagent.driver, - providerInstanceId: subagent.providerInstanceId, - occurredAt: now, - payload: { ...subagent, status: "interrupted", completedAt: now, updatedAt: now }, - }); + // A linked subagent child thread is a genuinely separate thread: its + // own open work needs its own reconciliation plan and its own effect + // cancellation, since cancelUnsettledEffects below only scopes the + // primary thread. + const linkedChildThreadIds = recoveryProjection.subagents + .filter((subagent) => subagent.runId === run.id) + .map((subagent) => subagent.childThreadId) + .filter((childThreadId): childThreadId is ThreadId => childThreadId !== null); + for (const childThreadId of linkedChildThreadIds) { + const childProjection = yield* projectionStore + .getRuntimeRecoveryProjection(childThreadId) + .pipe( + Effect.mapError( + (cause) => new OrchestratorProjectionError({ threadId: childThreadId, cause }), + ), + ); + const childPlan = yield* planThreadReconciliation({ + projection: childProjection, + trigger: "process-loss", + continueAfterRestart: false, + commandId: command.commandId, + now, + ids: idAllocator, + outbox, + }).pipe( + Effect.mapError( + (cause) => + new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause, + }), + ), + ); + yield* Ref.update(events, (existing) => [...existing, ...childPlan.events]); + yield* Ref.update(effects, (existing) => [...existing, ...childPlan.effects]); + if (childPlan.events.length === 0) continue; + const retiredEffectIds = yield* outbox + .cancelUnsettled({ + threadId: childThreadId, + effectTypes: PROCESS_BOUND_EFFECT_TYPES, + reason: childPlan.detail, + }) + .pipe( + Effect.mapError( + (cause) => + new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause, + }), + ), + ); + yield* outbox.signalCancellations(retiredEffectIds); } - // Same marker a live provider's turn.terminal interrupted would leave - // behind (RunExecutionService.writeFinalRunEvents); the web timeline - // keys the "Run interrupted" label off this item. - yield* emitEvent({ - type: "turn-item.updated", - threadId: command.threadId, - runId: run.id, - nodeId: rootNode.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: makeInterruptResultTurnItem({ - idAllocator, - run, - rootNode, - providerThread, - completedAt: now, - }), - }); if (command.holdQueue === true) yield* holdQueuedRuns; yield* stopCompletionCohort(); yield* settleBackgroundWork({ @@ -8131,8 +8113,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio now, }); return { - effectTypes: ["provider-turn.start", "provider-turn.restart"], - reason: `Run ${run.id} was interrupted after its provider session ${providerThread.providerSessionId} was no longer active.`, + effectTypes: PROCESS_BOUND_EFFECT_TYPES, + reason: plan.detail, } satisfies { readonly effectTypes: ReadonlyArray; readonly reason: string; @@ -9892,6 +9874,7 @@ export const layer: Layer.Layer< | CommandPolicyV2 | CommandReceiptStoreV2 | ContextHandoffServiceV2 + | EffectOutboxV2 | EventSinkV2 | IdAllocatorV2 | ProjectStore.ProjectStoreV2 diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index 67ae22f1883b..c2fa867f23ba 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -5,6 +5,7 @@ import { type ProviderThreadId, type OrchestrationV2RestartCancelledBackgroundWork, type OrchestrationV2ThreadProjection, + type RunId, ThreadId, } from "@t3tools/contracts"; import * as Context from "effect/Context"; @@ -54,6 +55,24 @@ export interface ProviderRuntimeReconciliationSummary { readonly requeuedEffects: number; } +/** + * Why a projection is being reconciled. "startup"/"shutdown" assume every + * provider process in the projection is dead. "process-loss" is narrower: one + * provider session died while the rest of the thread may still be live, so + * callers using it must scope the projection to what actually died (see + * `scopeProjectionToRun`). + */ +export type ReconciliationTrigger = "startup" | "shutdown" | "process-loss"; + +export interface ThreadReconciliationPlan { + readonly events: ReadonlyArray; + readonly effects: ReadonlyArray; + readonly detail: string; + readonly terminalizedRuns: number; + readonly stoppedSessions: number; + readonly closedRequests: number; +} + export class ProviderRuntimeRecoveryService extends Context.Service< ProviderRuntimeRecoveryService, { @@ -173,61 +192,101 @@ function latestStartedRun( ); } -export const make = Effect.gen(function* () { - const settings = yield* ServerSettings.ServerSettingsService; - const projections = yield* ProjectionStore.ProjectionStoreV2; - const eventSink = yield* EventSink.EventSinkV2; - const ids = yield* IdAllocator.IdAllocatorV2; - const outbox = yield* EffectOutbox.EffectOutboxV2; - const reconcileProjection = Effect.fn("ProviderRuntimeRecoveryService.reconcileProjection")( - function* ( - projection: ProjectionStore.ProjectionRuntimeRecoveryState, - trigger: "startup" | "shutdown", - continueAfterRestart: boolean, - ) { - const now = yield* DateTime.now; - const runs = [] as Array; - for (const run of nonterminalRuns(projection)) { - if (run.status === "waiting") { - const checkpointEffects = yield* outbox - .listByCommandId(CommandId.make(`command:effect:checkpoint.capture:${run.id}`)) - .pipe( - Effect.mapError( - (cause) => - new ProviderRuntimeRecoveryError({ - operation: "reconcile", - threadId: projection.thread.id, - cause, - }), - ), - ); - const hasReplayableCheckpoint = checkpointEffects.some( - (effect) => - effect.request.type === "checkpoint.capture" && - effect.request.runId === run.id && - (effect.status === "pending" || effect.status === "running"), - ); - if (hasReplayableCheckpoint) continue; - } - runs.push(run); - } - const messageRequestNodeIds = new Set( - projection.runtimeRequests - .filter( - (request) => - request.status === "pending" && request.responseCapability.type === "message", - ) - .map((request) => request.nodeId), - ); - const requests = projection.runtimeRequests.filter( - (request) => request.status === "pending" && request.responseCapability.type !== "message", - ); - const detail = `Cancelled because the server ${trigger === "startup" ? "restarted" : "shut down"} before the provider work completed.`; - const commandId = CommandId.make( - `command:runtime-reconcile:${trigger}:${projection.thread.id}:${DateTime.formatIso(now)}`, - ); - const allocateEventId = () => - ids.allocate.event({ threadId: projection.thread.id, commandId }).pipe( +function reconciliationDetail(trigger: ReconciliationTrigger): string { + switch (trigger) { + case "startup": + return "Cancelled because the server restarted before the provider work completed."; + case "shutdown": + return "Cancelled because the server shut down before the provider work completed."; + case "process-loss": + return "Cancelled because its provider session ended before the work completed."; + } +} + +function reconciliationRequestReason(trigger: ReconciliationTrigger): string { + switch (trigger) { + case "startup": + return "The server restarted before this runtime request was resolved."; + case "shutdown": + return "The server shut down before this runtime request was resolved."; + case "process-loss": + return "The provider session ended before this runtime request was resolved."; + } +} + +/** + * Narrow a thread's recovery projection to one run and the provider thread + * whose session died for it: that run's own records, plus runless native + * subagent nodes/items owned by the same provider thread. Other runs, + * provider threads and sessions on the same orchestration thread are left + * out entirely, so a "process-loss" plan built from this view cannot cancel + * work that is still alive elsewhere on the thread. + */ +export function scopeProjectionToRun( + projection: ProjectionStore.ProjectionRuntimeRecoveryState, + input: { readonly runId: RunId; readonly providerThreadId: ProviderThreadId }, +): ProjectionStore.ProjectionRuntimeRecoveryState { + const { runId, providerThreadId } = input; + const attempts = projection.attempts.filter((attempt) => attempt.runId === runId); + const attemptIds = new Set(attempts.map((attempt) => attempt.id)); + const nodes = projection.nodes.filter( + (node) => + node.runId === runId || (node.runId === null && node.providerThreadId === providerThreadId), + ); + const nodeIds = new Set(nodes.map((node) => node.id)); + return { + thread: projection.thread, + runs: projection.runs.filter((run) => run.id === runId), + attempts, + nodes, + subagents: projection.subagents.filter((subagent) => subagent.runId === runId), + providerSessions: projection.providerSessions.filter( + (session) => + session.id === + projection.providerThreads.find((providerThread) => providerThread.id === providerThreadId) + ?.providerSessionId, + ), + providerThreads: projection.providerThreads.filter( + (providerThread) => providerThread.id === providerThreadId, + ), + providerTurns: projection.providerTurns.filter( + (providerTurn) => + providerTurn.runAttemptId !== null && attemptIds.has(providerTurn.runAttemptId), + ), + runtimeRequests: projection.runtimeRequests.filter((request) => nodeIds.has(request.nodeId)), + messages: projection.messages.filter((message) => message.runId === runId), + turnItems: projection.turnItems.filter( + (item) => + item.runId === runId || (item.runId === null && item.providerThreadId === providerThreadId), + ), + }; +} + +/** + * Build the events and effects that cancel a projection's open work. Pure + * planning only: the caller decides how to commit (a standalone reconcile + * command, or folded into a larger command's own commit). "process-loss" + * callers must scope `projection` first (see `scopeProjectionToRun`), since + * this plans against everything the projection contains. + */ +export const planThreadReconciliation = Effect.fn( + "ProviderRuntimeRecoveryService.planThreadReconciliation", +)(function* (input: { + readonly projection: ProjectionStore.ProjectionRuntimeRecoveryState; + readonly trigger: ReconciliationTrigger; + readonly continueAfterRestart: boolean; + readonly commandId: CommandId; + readonly now: DateTime.Utc; + readonly ids: IdAllocator.IdAllocatorV2Shape; + readonly outbox: EffectOutbox.EffectOutboxV2Shape; +}) { + const { projection, trigger, continueAfterRestart, commandId, now, ids, outbox } = input; + const runs = [] as Array; + for (const run of nonterminalRuns(projection)) { + if (run.status === "waiting") { + const checkpointEffects = yield* outbox + .listByCommandId(CommandId.make(`command:effect:checkpoint.capture:${run.id}`)) + .pipe( Effect.mapError( (cause) => new ProviderRuntimeRecoveryError({ @@ -237,425 +296,528 @@ export const make = Effect.gen(function* () { }), ), ); - const events: Array = []; - // Background work that outlived its settled turn. The provider transcript - // cannot record its death, so the next provider turn is told instead. - // Shutdown records it too: a graceful restart cancels it there first. - // Keyed by the provider thread that lost the work: only its turns are told. - const cancelledBackgroundWork = new Map< - ProviderThreadId, - Array - >(); - const cancelledBackgroundNativeIds = new Set(); - const recordCancelledBackgroundWork = ( - providerThreadId: ProviderThreadId | null | undefined, - work: OrchestrationV2RestartCancelledBackgroundWork, - ) => { - if (providerThreadId == null) return; - const existing = cancelledBackgroundWork.get(providerThreadId); - if (existing === undefined) cancelledBackgroundWork.set(providerThreadId, [work]); - else existing.push(work); - }; - const recordCancelledBackgroundItem = ( - item: OrchestrationV2ThreadProjection["turnItems"][number], - ) => { - if (!isBackgroundCapableTurnItemType(item.type)) return; - const work = cancelledTurnItemWork(item); - if (work === undefined) return; - recordCancelledBackgroundWork( - item.providerThreadId ?? - projection.runs.find((run) => run.id === item.runId)?.providerThreadId, - work, - ); - if (item.nativeItemRef?.nativeId != null) { - cancelledBackgroundNativeIds.add(item.nativeItemRef.nativeId); - } - }; - // Queued runs have not started provider work. Preserve their execution - // identities and order, but require explicit consent before draining them. - for (const run of projection.runs) { - if (run.status !== "queued" || run.queueHeld === true) continue; - events.push({ - id: yield* allocateEventId(), - type: "run.updated", - threadId: projection.thread.id, - runId: run.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...run, queueHeld: true }, - }); - } - for (const request of requests) { - events.push({ - id: yield* allocateEventId(), - type: "runtime-request.updated", - threadId: projection.thread.id, - nodeId: request.nodeId, - occurredAt: now, - payload: { - ...request, - status: trigger === "startup" ? "expired" : "cancelled", - responseCapability: { - type: "not_resumable", - reason: `The server ${trigger === "startup" ? "restarted" : "shut down"} before this runtime request was resolved.`, - }, - resolvedAt: now, - }, - }); - } - for (const run of runs) { - events.push({ - id: yield* allocateEventId(), - type: "run.updated", - threadId: projection.thread.id, - runId: run.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...run, status: "cancelled", queuePosition: null, completedAt: now }, - }); - for (const attempt of projection.attempts.filter( - (candidate) => - candidate.runId === run.id && - (candidate.status === "pending" || candidate.status === "running"), - )) { - events.push({ - id: yield* allocateEventId(), - type: "run-attempt.updated", - threadId: projection.thread.id, - runId: run.id, - nodeId: attempt.rootNodeId, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...attempt, status: "cancelled", completedAt: now }, - }); - } - for (const node of projection.nodes.filter( - (candidate) => - candidate.runId === run.id && - !messageRequestNodeIds.has(candidate.id) && - (candidate.status === "pending" || - candidate.status === "running" || - candidate.status === "waiting"), - )) { - events.push({ - id: yield* allocateEventId(), - type: "node.updated", - threadId: projection.thread.id, - runId: run.id, - nodeId: node.id, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...node, status: "cancelled", completedAt: now }, - }); - } - for (const subagent of projection.subagents.filter( - (candidate) => - candidate.runId === run.id && - (candidate.status === "pending" || - candidate.status === "running" || - candidate.status === "waiting"), - )) { - events.push({ - id: yield* allocateEventId(), - type: "subagent.updated", - threadId: projection.thread.id, - runId: run.id, - nodeId: subagent.id, - driver: subagent.driver, - providerInstanceId: subagent.providerInstanceId, - occurredAt: now, - payload: { ...subagent, status: "cancelled", completedAt: now, updatedAt: now }, - }); - } - for (const providerTurn of projection.providerTurns.filter( - (candidate) => - candidate.runAttemptId !== null && - projection.attempts.some( - (attempt) => attempt.id === candidate.runAttemptId && attempt.runId === run.id, - ) && - (candidate.status === "pending" || candidate.status === "running"), - )) { - events.push({ - id: yield* allocateEventId(), - type: "provider-turn.updated", - threadId: projection.thread.id, - runId: run.id, - nodeId: providerTurn.nodeId, - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...providerTurn, status: "cancelled", completedAt: now }, - }); - } - for (const message of projection.messages.filter( - (candidate) => candidate.runId === run.id && candidate.streaming, - )) { - events.push({ - id: yield* allocateEventId(), - type: "message.updated", - threadId: projection.thread.id, - runId: run.id, - ...(message.nodeId === null ? {} : { nodeId: message.nodeId }), - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...message, streaming: false, updatedAt: now }, - }); - } - for (const item of projection.turnItems.filter( - (candidate) => - candidate.runId === run.id && - (candidate.nodeId === null || !messageRequestNodeIds.has(candidate.nodeId)) && - (candidate.status === "pending" || - candidate.status === "running" || - candidate.status === "waiting"), - )) { - // A waiting run's provider turn already settled; its open items are - // background work. A running run's items die with its turn. - if (run.status === "waiting") recordCancelledBackgroundItem(item); - events.push({ - id: yield* allocateEventId(), - type: "turn-item.updated", - threadId: projection.thread.id, - runId: run.id, - ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), - providerInstanceId: run.providerInstanceId, - occurredAt: now, - payload: { ...item, status: "cancelled", completedAt: now, updatedAt: now }, - }); - } - } - // Process loss also orphans background-capable turn items on already- - // settled runs (e.g. post-settle Waiting work). Skip items already - // cancelled above for recovered nonterminal runs to avoid duplicate - // cancellation events. - const recoveredNonterminalRunIds = new Set(runs.map((run) => run.id)); - const cancelledStaleNodeIds = new Set(); - for (const item of projection.turnItems ?? []) { - if (item.runId !== null && recoveredNonterminalRunIds.has(item.runId)) { - continue; - } - if (!isBackgroundCapableTurnItemType(item.type)) { - continue; - } - if (!isNonterminalTurnItemStatus(item.status)) { - continue; - } - const providerInstanceId = resolveStaleBackgroundItemProviderInstanceId(item, projection); - recordCancelledBackgroundItem(item); - events.push({ - id: yield* allocateEventId(), - type: "turn-item.updated", - threadId: projection.thread.id, - ...(item.runId === null ? {} : { runId: item.runId }), - ...(item.nodeId === null || item.nodeId === undefined ? {} : { nodeId: item.nodeId }), - providerInstanceId, - occurredAt: now, - payload: { ...item, status: "cancelled", completedAt: now, updatedAt: now }, - }); - if (item.nodeId !== null && item.nodeId !== undefined) { - const staleItemNode = projection.nodes.find( - (candidate) => - candidate.id === item.nodeId && isNonterminalNodeStatus(candidate.status), - ); - if (staleItemNode !== undefined && !cancelledStaleNodeIds.has(staleItemNode.id)) { - cancelledStaleNodeIds.add(staleItemNode.id); - events.push({ - id: yield* allocateEventId(), - type: "node.updated", - threadId: projection.thread.id, - ...(item.runId === null ? {} : { runId: item.runId }), - nodeId: staleItemNode.id, - providerInstanceId, - occurredAt: now, - payload: { ...staleItemNode, status: "cancelled", completedAt: now }, - }); - } - } - if (item.type !== "subagent") { - continue; - } - // Cancelling only the turn item would leave the linked subagent entity - // non-terminal forever, since the dead provider process can no longer - // emit its terminal event. Match the exact linked id so a subagent - // that already finished is never overwritten. - const staleSubagent = projection.subagents.find( - (candidate) => - candidate.id === item.subagentId && isNonterminalSubagentStatus(candidate.status), - ); - if (staleSubagent !== undefined) { - events.push({ - id: yield* allocateEventId(), - type: "subagent.updated", - threadId: projection.thread.id, - ...(item.runId === null ? {} : { runId: item.runId }), - nodeId: staleSubagent.id, - driver: staleSubagent.driver, - providerInstanceId: staleSubagent.providerInstanceId, - occurredAt: now, - payload: { ...staleSubagent, status: "cancelled", completedAt: now, updatedAt: now }, - }); - } - const staleSubagentNode = projection.nodes.find( - (candidate) => - candidate.id === item.subagentId && isNonterminalNodeStatus(candidate.status), - ); - if (staleSubagentNode !== undefined && !cancelledStaleNodeIds.has(staleSubagentNode.id)) { - cancelledStaleNodeIds.add(staleSubagentNode.id); - events.push({ - id: yield* allocateEventId(), - type: "node.updated", + const hasReplayableCheckpoint = checkpointEffects.some( + (effect) => + effect.request.type === "checkpoint.capture" && + effect.request.runId === run.id && + (effect.status === "pending" || effect.status === "running"), + ); + if (hasReplayableCheckpoint) continue; + } + runs.push(run); + } + const messageRequestNodeIds = new Set( + projection.runtimeRequests + .filter( + (request) => request.status === "pending" && request.responseCapability.type === "message", + ) + .map((request) => request.nodeId), + ); + const requests = projection.runtimeRequests.filter( + (request) => request.status === "pending" && request.responseCapability.type !== "message", + ); + const detail = reconciliationDetail(trigger); + const allocateEventId = () => + ids.allocate.event({ threadId: projection.thread.id, commandId }).pipe( + Effect.mapError( + (cause) => + new ProviderRuntimeRecoveryError({ + operation: "reconcile", threadId: projection.thread.id, - ...(item.runId === null ? {} : { runId: item.runId }), - nodeId: staleSubagentNode.id, - providerInstanceId, - occurredAt: now, - payload: { ...staleSubagentNode, status: "cancelled", completedAt: now }, - }); - } - } - // A provider-native subagent thread has no runs: its work is a runless - // root turn, plus items under it (Claude's live progress item), that - // only the dead provider process could settle. Left running, the child - // would show as working forever. - const cancelledStaleItemIds = new Set( - events.flatMap((event) => (event.type === "turn-item.updated" ? [event.payload.id] : [])), + cause, + }), + ), + ); + const events: Array = []; + // Background work that outlived its settled turn. The provider transcript + // cannot record its death, so the next provider turn is told instead. + // Shutdown records it too: a graceful restart cancels it there first. + // Keyed by the provider thread that lost the work: only its turns are told. + const cancelledBackgroundWork = new Map< + ProviderThreadId, + Array + >(); + const cancelledBackgroundNativeIds = new Set(); + const recordCancelledBackgroundWork = ( + providerThreadId: ProviderThreadId | null | undefined, + work: OrchestrationV2RestartCancelledBackgroundWork, + ) => { + if (providerThreadId == null) return; + const existing = cancelledBackgroundWork.get(providerThreadId); + if (existing === undefined) cancelledBackgroundWork.set(providerThreadId, [work]); + else existing.push(work); + }; + const recordCancelledBackgroundItem = ( + item: OrchestrationV2ThreadProjection["turnItems"][number], + ) => { + if (!isBackgroundCapableTurnItemType(item.type)) return; + const work = cancelledTurnItemWork(item); + if (work === undefined) return; + recordCancelledBackgroundWork( + item.providerThreadId ?? + projection.runs.find((run) => run.id === item.runId)?.providerThreadId, + work, + ); + if (item.nativeItemRef?.nativeId != null) { + cancelledBackgroundNativeIds.add(item.nativeItemRef.nativeId); + } + }; + // Queued runs have not started provider work. Preserve their execution + // identities and order, but require explicit consent before draining them. + for (const run of projection.runs) { + if (run.status !== "queued" || run.queueHeld === true) continue; + events.push({ + id: yield* allocateEventId(), + type: "run.updated", + threadId: projection.thread.id, + runId: run.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...run, queueHeld: true }, + }); + } + for (const request of requests) { + events.push({ + id: yield* allocateEventId(), + type: "runtime-request.updated", + threadId: projection.thread.id, + nodeId: request.nodeId, + occurredAt: now, + payload: { + ...request, + status: trigger === "startup" ? "expired" : "cancelled", + responseCapability: { + type: "not_resumable", + reason: reconciliationRequestReason(trigger), + }, + resolvedAt: now, + }, + }); + } + for (const run of runs) { + events.push({ + id: yield* allocateEventId(), + type: "run.updated", + threadId: projection.thread.id, + runId: run.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...run, status: "cancelled", queuePosition: null, completedAt: now }, + }); + for (const attempt of projection.attempts.filter( + (candidate) => + candidate.runId === run.id && + (candidate.status === "pending" || candidate.status === "running"), + )) { + events.push({ + id: yield* allocateEventId(), + type: "run-attempt.updated", + threadId: projection.thread.id, + runId: run.id, + nodeId: attempt.rootNodeId, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...attempt, status: "cancelled", completedAt: now }, + }); + } + for (const node of projection.nodes.filter( + (candidate) => + candidate.runId === run.id && + !messageRequestNodeIds.has(candidate.id) && + (candidate.status === "pending" || + candidate.status === "running" || + candidate.status === "waiting"), + )) { + events.push({ + id: yield* allocateEventId(), + type: "node.updated", + threadId: projection.thread.id, + runId: run.id, + nodeId: node.id, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...node, status: "cancelled", completedAt: now }, + }); + } + for (const subagent of projection.subagents.filter( + (candidate) => + candidate.runId === run.id && + (candidate.status === "pending" || + candidate.status === "running" || + candidate.status === "waiting"), + )) { + events.push({ + id: yield* allocateEventId(), + type: "subagent.updated", + threadId: projection.thread.id, + runId: run.id, + nodeId: subagent.id, + driver: subagent.driver, + providerInstanceId: subagent.providerInstanceId, + occurredAt: now, + payload: { ...subagent, status: "cancelled", completedAt: now, updatedAt: now }, + }); + } + for (const providerTurn of projection.providerTurns.filter( + (candidate) => + candidate.runAttemptId !== null && + projection.attempts.some( + (attempt) => attempt.id === candidate.runAttemptId && attempt.runId === run.id, + ) && + (candidate.status === "pending" || candidate.status === "running"), + )) { + events.push({ + id: yield* allocateEventId(), + type: "provider-turn.updated", + threadId: projection.thread.id, + runId: run.id, + nodeId: providerTurn.nodeId, + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...providerTurn, status: "cancelled", completedAt: now }, + }); + } + for (const message of projection.messages.filter( + (candidate) => candidate.runId === run.id && candidate.streaming, + )) { + events.push({ + id: yield* allocateEventId(), + type: "message.updated", + threadId: projection.thread.id, + runId: run.id, + ...(message.nodeId === null ? {} : { nodeId: message.nodeId }), + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { ...message, streaming: false, updatedAt: now }, + }); + } + for (const item of projection.turnItems.filter( + (candidate) => + candidate.runId === run.id && + (candidate.nodeId === null || !messageRequestNodeIds.has(candidate.nodeId)) && + (candidate.status === "pending" || + candidate.status === "running" || + candidate.status === "waiting"), + )) { + // A waiting run's provider turn already settled; its open items are + // background work. A running run's items die with its turn. + if (run.status === "waiting") recordCancelledBackgroundItem(item); + events.push({ + id: yield* allocateEventId(), + type: "turn-item.updated", + threadId: projection.thread.id, + runId: run.id, + ...(item.nodeId === null ? {} : { nodeId: item.nodeId }), + providerInstanceId: run.providerInstanceId, + occurredAt: now, + payload: { + ...item, + status: "cancelled", + completedAt: now, + updatedAt: now, + ...(item.type === "reasoning" || item.type === "assistant_message" + ? { streaming: false } + : {}), + }, + }); + } + } + // Process loss also orphans background-capable turn items on already- + // settled runs (e.g. post-settle Waiting work). Skip items already + // cancelled above for recovered nonterminal runs to avoid duplicate + // cancellation events. + const recoveredNonterminalRunIds = new Set(runs.map((run) => run.id)); + const cancelledStaleNodeIds = new Set(); + for (const item of projection.turnItems ?? []) { + if (item.runId !== null && recoveredNonterminalRunIds.has(item.runId)) { + continue; + } + if (!isBackgroundCapableTurnItemType(item.type)) { + continue; + } + if (!isNonterminalTurnItemStatus(item.status)) { + continue; + } + const providerInstanceId = resolveStaleBackgroundItemProviderInstanceId(item, projection); + recordCancelledBackgroundItem(item); + events.push({ + id: yield* allocateEventId(), + type: "turn-item.updated", + threadId: projection.thread.id, + ...(item.runId === null ? {} : { runId: item.runId }), + ...(item.nodeId === null || item.nodeId === undefined ? {} : { nodeId: item.nodeId }), + providerInstanceId, + occurredAt: now, + payload: { ...item, status: "cancelled", completedAt: now, updatedAt: now }, + }); + if (item.nodeId !== null && item.nodeId !== undefined) { + const staleItemNode = projection.nodes.find( + (candidate) => candidate.id === item.nodeId && isNonterminalNodeStatus(candidate.status), ); - for (const node of projection.nodes) { - if ( - node.kind !== "root_turn" || - node.runId !== null || - !isNonterminalNodeStatus(node.status) || - cancelledStaleNodeIds.has(node.id) - ) { - continue; - } - cancelledStaleNodeIds.add(node.id); + if (staleItemNode !== undefined && !cancelledStaleNodeIds.has(staleItemNode.id)) { + cancelledStaleNodeIds.add(staleItemNode.id); events.push({ id: yield* allocateEventId(), type: "node.updated", threadId: projection.thread.id, - nodeId: node.id, - providerInstanceId: projection.thread.providerInstanceId, + ...(item.runId === null ? {} : { runId: item.runId }), + nodeId: staleItemNode.id, + providerInstanceId, occurredAt: now, - payload: { ...node, status: "cancelled", completedAt: now }, + payload: { ...staleItemNode, status: "cancelled", completedAt: now }, }); - for (const item of projection.turnItems) { - if ( - item.nodeId !== node.id || - item.runId !== null || - !isNonterminalTurnItemStatus(item.status) || - cancelledStaleItemIds.has(item.id) - ) { - continue; - } - cancelledStaleItemIds.add(item.id); - events.push({ - id: yield* allocateEventId(), - type: "turn-item.updated", - threadId: projection.thread.id, - nodeId: node.id, - providerInstanceId: projection.thread.providerInstanceId, - occurredAt: now, - payload: { - ...item, - status: "cancelled", - completedAt: now, - updatedAt: now, - ...(item.type === "reasoning" || item.type === "assistant_message" - ? { streaming: false } - : {}), - }, - }); - } } - // All provider processes are gone on startup/shutdown: clear any - // persisted Waiting roster (including idle threads from settled roots) - // and idle active threads without resurrecting active status. - for (const providerThread of projection.providerThreads ?? []) { - const needsIdle = providerThread.status === "active"; - const needsRosterClear = providerThreadHasPendingBackgroundTasks(providerThread); - if (!needsIdle && !needsRosterClear) { - continue; - } - if (providerThread.ownerNodeId === null) { - for (const task of providerThread.pendingBackgroundTasks ?? []) { - if (cancelledBackgroundNativeIds.has(task.taskId)) continue; - cancelledBackgroundNativeIds.add(task.taskId); - recordCancelledBackgroundWork(providerThread.id, cancelledRosterTaskWork(task)); - } - } - events.push({ - id: yield* allocateEventId(), - type: "provider-thread.updated", - threadId: projection.thread.id, - driver: providerThread.driver, - providerInstanceId: providerThread.providerInstanceId, - occurredAt: now, - payload: { - ...providerThread, - status: needsIdle ? "idle" : providerThread.status, - pendingBackgroundTasks: [], - updatedAt: now, - }, - }); + } + if (item.type !== "subagent") { + continue; + } + // Cancelling only the turn item would leave the linked subagent entity + // non-terminal forever, since the dead provider process can no longer + // emit its terminal event. Match the exact linked id so a subagent + // that already finished is never overwritten. + const staleSubagent = projection.subagents.find( + (candidate) => + candidate.id === item.subagentId && isNonterminalSubagentStatus(candidate.status), + ); + if (staleSubagent !== undefined) { + events.push({ + id: yield* allocateEventId(), + type: "subagent.updated", + threadId: projection.thread.id, + ...(item.runId === null ? {} : { runId: item.runId }), + nodeId: staleSubagent.id, + driver: staleSubagent.driver, + providerInstanceId: staleSubagent.providerInstanceId, + occurredAt: now, + payload: { ...staleSubagent, status: "cancelled", completedAt: now, updatedAt: now }, + }); + } + const staleSubagentNode = projection.nodes.find( + (candidate) => candidate.id === item.subagentId && isNonterminalNodeStatus(candidate.status), + ); + if (staleSubagentNode !== undefined && !cancelledStaleNodeIds.has(staleSubagentNode.id)) { + cancelledStaleNodeIds.add(staleSubagentNode.id); + events.push({ + id: yield* allocateEventId(), + type: "node.updated", + threadId: projection.thread.id, + ...(item.runId === null ? {} : { runId: item.runId }), + nodeId: staleSubagentNode.id, + providerInstanceId, + occurredAt: now, + payload: { ...staleSubagentNode, status: "cancelled", completedAt: now }, + }); + } + } + // A provider-native subagent thread has no runs: its work is a runless + // root turn, plus items under it (Claude's live progress item), that + // only the dead provider process could settle. Left running, the child + // would show as working forever. + const cancelledStaleItemIds = new Set( + events.flatMap((event) => (event.type === "turn-item.updated" ? [event.payload.id] : [])), + ); + for (const node of projection.nodes) { + if ( + node.kind !== "root_turn" || + node.runId !== null || + !isNonterminalNodeStatus(node.status) || + cancelledStaleNodeIds.has(node.id) + ) { + continue; + } + cancelledStaleNodeIds.add(node.id); + events.push({ + id: yield* allocateEventId(), + type: "node.updated", + threadId: projection.thread.id, + nodeId: node.id, + providerInstanceId: projection.thread.providerInstanceId, + occurredAt: now, + payload: { ...node, status: "cancelled", completedAt: now }, + }); + for (const item of projection.turnItems) { + if ( + item.nodeId !== node.id || + item.runId !== null || + !isNonterminalTurnItemStatus(item.status) || + cancelledStaleItemIds.has(item.id) + ) { + continue; } - for (const session of projection.providerSessions.filter( - (candidate) => candidate.status !== "stopped" && candidate.status !== "error", - )) { - events.push({ - id: yield* allocateEventId(), - type: "provider-session.updated", - threadId: projection.thread.id, - driver: session.driver, - providerInstanceId: session.providerInstanceId, - occurredAt: now, - payload: { ...session, status: "stopped", updatedAt: now, lastError: null }, - }); + cancelledStaleItemIds.add(item.id); + events.push({ + id: yield* allocateEventId(), + type: "turn-item.updated", + threadId: projection.thread.id, + nodeId: node.id, + providerInstanceId: projection.thread.providerInstanceId, + occurredAt: now, + payload: { + ...item, + status: "cancelled", + completedAt: now, + updatedAt: now, + ...(item.type === "reasoning" || item.type === "assistant_message" + ? { streaming: false } + : {}), + }, + }); + } + } + // All provider processes are gone on startup/shutdown: clear any + // persisted Waiting roster (including idle threads from settled roots) + // and idle active threads without resurrecting active status. + for (const providerThread of projection.providerThreads ?? []) { + const needsIdle = providerThread.status === "active"; + const needsRosterClear = providerThreadHasPendingBackgroundTasks(providerThread); + if (!needsIdle && !needsRosterClear) { + continue; + } + if (providerThread.ownerNodeId === null) { + for (const task of providerThread.pendingBackgroundTasks ?? []) { + if (cancelledBackgroundNativeIds.has(task.taskId)) continue; + cancelledBackgroundNativeIds.add(task.taskId); + recordCancelledBackgroundWork(providerThread.id, cancelledRosterTaskWork(task)); } - for (const [providerThreadId, work] of cancelledBackgroundWork) { - const noteRun = latestStartedRun(projection, providerThreadId); - if (noteRun === undefined) continue; - // Its own event: a run snapshot read before this commit could regress - // a lifecycle change (e.g. a checkpoint completing the run) made since. - events.push({ - id: yield* allocateEventId(), - type: "run.background-work-cancelled", + } + events.push({ + id: yield* allocateEventId(), + type: "provider-thread.updated", + threadId: projection.thread.id, + driver: providerThread.driver, + providerInstanceId: providerThread.providerInstanceId, + occurredAt: now, + payload: { + ...providerThread, + status: needsIdle ? "idle" : providerThread.status, + pendingBackgroundTasks: [], + updatedAt: now, + }, + }); + } + for (const session of projection.providerSessions.filter( + (candidate) => candidate.status !== "stopped" && candidate.status !== "error", + )) { + events.push({ + id: yield* allocateEventId(), + type: "provider-session.updated", + threadId: projection.thread.id, + driver: session.driver, + providerInstanceId: session.providerInstanceId, + occurredAt: now, + payload: { ...session, status: "stopped", updatedAt: now, lastError: null }, + }); + } + for (const [providerThreadId, work] of cancelledBackgroundWork) { + const noteRun = latestStartedRun(projection, providerThreadId); + if (noteRun === undefined) continue; + // Its own event: a run snapshot read before this commit could regress + // a lifecycle change (e.g. a checkpoint completing the run) made since. + events.push({ + id: yield* allocateEventId(), + type: "run.background-work-cancelled", + threadId: projection.thread.id, + runId: noteRun.id, + providerInstanceId: noteRun.providerInstanceId, + occurredAt: now, + payload: { + runId: noteRun.id, + restartCancelledBackgroundWork: mergeRestartCancelledBackgroundWork( + noteRun.restartCancelledBackgroundWork ?? [], + work, + ), + }, + }); + } + const continuationRun = + continueAfterRestart && trigger === "startup" + ? restartContinuationRun(projection, new Set(cancelledBackgroundWork.keys())) + : undefined; + const effects: Array = continuationRun + ? [ + { + id: `effect:restart-continuation:${continuationRun.id}`, + commandId, threadId: projection.thread.id, - runId: noteRun.id, - providerInstanceId: noteRun.providerInstanceId, - occurredAt: now, - payload: { - runId: noteRun.id, - restartCancelledBackgroundWork: mergeRestartCancelledBackgroundWork( - noteRun.restartCancelledBackgroundWork ?? [], - work, - ), - }, - }); - } - const continuationRun = - continueAfterRestart && trigger === "startup" - ? restartContinuationRun(projection, new Set(cancelledBackgroundWork.keys())) - : undefined; - const effects: Array = continuationRun - ? [ - { - id: `effect:restart-continuation:${continuationRun.id}`, - commandId, - threadId: projection.thread.id, - request: { type: "provider-runtime.continue", sourceRunId: continuationRun.id }, - }, - ] - : []; - const stoppedSessions = projection.providerSessions.filter( - (candidate) => candidate.status !== "stopped" && candidate.status !== "error", - ).length; + request: { type: "provider-runtime.continue", sourceRunId: continuationRun.id }, + }, + ] + : []; + const stoppedSessions = projection.providerSessions.filter( + (candidate) => candidate.status !== "stopped" && candidate.status !== "error", + ).length; + return { + events, + effects, + detail, + terminalizedRuns: runs.length, + stoppedSessions, + closedRequests: requests.length, + } satisfies ThreadReconciliationPlan; +}); + +/** + * Narrow a thread's recovery plan to a single run whose provider session + * died while the rest of the thread may still be live: scopes the projection + * first, then plans against only what died. + */ +export const planProcessLossReconciliation = Effect.fn( + "ProviderRuntimeRecoveryService.planProcessLossReconciliation", +)(function* (input: { + readonly projection: ProjectionStore.ProjectionRuntimeRecoveryState; + readonly runId: RunId; + readonly providerThreadId: ProviderThreadId; + readonly commandId: CommandId; + readonly now: DateTime.Utc; + readonly ids: IdAllocator.IdAllocatorV2Shape; + readonly outbox: EffectOutbox.EffectOutboxV2Shape; +}) { + return yield* planThreadReconciliation({ + projection: scopeProjectionToRun(input.projection, { + runId: input.runId, + providerThreadId: input.providerThreadId, + }), + trigger: "process-loss", + continueAfterRestart: false, + commandId: input.commandId, + now: input.now, + ids: input.ids, + outbox: input.outbox, + }); +}); + +export const make = Effect.gen(function* () { + const settings = yield* ServerSettings.ServerSettingsService; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const eventSink = yield* EventSink.EventSinkV2; + const ids = yield* IdAllocator.IdAllocatorV2; + const outbox = yield* EffectOutbox.EffectOutboxV2; + const reconcileProjection = Effect.fn("ProviderRuntimeRecoveryService.reconcileProjection")( + function* ( + projection: ProjectionStore.ProjectionRuntimeRecoveryState, + trigger: "startup" | "shutdown", + continueAfterRestart: boolean, + ) { + const now = yield* DateTime.now; + const commandId = CommandId.make( + `command:runtime-reconcile:${trigger}:${projection.thread.id}:${DateTime.formatIso(now)}`, + ); + const plan = yield* planThreadReconciliation({ + projection, + trigger, + continueAfterRestart, + commandId, + now, + ids, + outbox, + }); let retiredEffects: number; - if (events.length === 0) { + if (plan.events.length === 0) { const retiredEffectIds = yield* outbox .cancelUnsettled({ threadId: projection.thread.id, effectTypes: EffectOutbox.PROCESS_BOUND_EFFECT_TYPES, - reason: detail, + reason: plan.detail, }) .pipe( Effect.mapError( @@ -663,7 +825,7 @@ export const make = Effect.gen(function* () { new ProviderRuntimeRecoveryError({ operation: "reconcile", threadId: projection.thread.id, - cause: { detail, cause }, + cause: { detail: plan.detail, cause }, }), ), ); @@ -676,11 +838,11 @@ export const make = Effect.gen(function* () { threadId: projection.thread.id, commandType: "provider-runtime.reconcile", acceptedAt: now, - events, - effects, + events: plan.events, + effects: plan.effects, cancelUnsettledEffects: { effectTypes: EffectOutbox.PROCESS_BOUND_EFFECT_TYPES, - reason: detail, + reason: plan.detail, }, }) .pipe( @@ -696,9 +858,9 @@ export const make = Effect.gen(function* () { retiredEffects = result.cancelledEffectCount; } return { - terminalizedRuns: runs.length, - stoppedSessions, - closedRequests: requests.length, + terminalizedRuns: plan.terminalizedRuns, + stoppedSessions: plan.stoppedSessions, + closedRequests: plan.closedRequests, retiredEffects, }; }, diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index aab019166416..71211c126545 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -1433,7 +1433,7 @@ export const layer: Layer.Layer< }), ); -export function makeInterruptResultTurnItem(input: { +function makeInterruptResultTurnItem(input: { readonly idAllocator: IdAllocator.IdAllocatorV2Shape; readonly run: OrchestrationV2Run; readonly rootNode: OrchestrationV2ExecutionNode; diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 7659788b6399..0a57f57771eb 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -3001,7 +3001,9 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; const eventSink = yield* EventSink.EventSinkV2; + const outbox = yield* EffectOutbox.EffectOutboxV2; const threadId = ThreadId.make("runtime-layer-interrupt-dead-session"); + const childThreadId = ThreadId.make(`${threadId}:child`); yield* orchestrator.dispatch({ type: "thread.create", createdBy: "user", @@ -3030,6 +3032,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { const before = yield* orchestrator.getThreadProjection(threadId); const activeRun = before.runs[0]!; const providerThread = before.providerThreads[0]!; + const checkpointScope = before.checkpointScopes[0]!; const now = yield* DateTime.now; const providerTurn = { id: ProviderTurnId.make(`${threadId}:turn`), @@ -3044,7 +3047,85 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }; const streamingMessageId = MessageId.make(`${threadId}:assistant-message`); const streamingItemId = TurnItemId.make(`${threadId}:assistant-item`); + const dynamicToolItemId = TurnItemId.make(`${threadId}:dynamic-tool`); const subagentId = NodeId.make(`${threadId}:subagent`); + + // A delegated subagent child thread: a genuinely separate thread doing + // work on behalf of the active run, linked via subagent.childThreadId. + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "agent", + creationSource: "web", + commandId: CommandId.make(`${childThreadId}:create`), + threadId: childThreadId, + projectId: ProjectId.make(`${threadId}:project`), + title: "Delegated subtask", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: process.cwd(), + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${childThreadId}:message:0`), + threadId: childThreadId, + messageId: MessageId.make(`${childThreadId}:message:0`), + text: "Do the subtask", + attachments: [], + dispatchMode: { type: "start_immediately" }, + }); + const childBefore = yield* orchestrator.getThreadProjection(childThreadId); + const childRun = childBefore.runs[0]!; + const childProviderThread = childBefore.providerThreads[0]!; + const childItemId = TurnItemId.make(`${childThreadId}:dynamic-tool`); + yield* eventSink.write({ + commandId: CommandId.make(`${childThreadId}:force-running`), + events: [ + { + id: EventId.make(`${childThreadId}:run-running`), + type: "run.updated", + threadId: childThreadId, + runId: childRun.id, + occurredAt: now, + payload: { ...childRun, status: "running", startedAt: now }, + }, + { + id: EventId.make(`${childThreadId}:item-running`), + type: "turn-item.updated", + threadId: childThreadId, + runId: childRun.id, + nodeId: childRun.rootNodeId!, + occurredAt: now, + payload: { + id: childItemId, + type: "dynamic_tool", + threadId: childThreadId, + runId: childRun.id, + nodeId: childRun.rootNodeId!, + providerThreadId: childProviderThread.id, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 0, + status: "running", + title: "Subtask work", + startedAt: now, + completedAt: null, + updatedAt: now, + toolName: "subtask_tool", + input: {}, + }, + }, + ], + }); + + const checkpointCommandId = CommandId.make( + `command:effect:checkpoint.capture:${activeRun.id}`, + ); + // Drive the run and its provider turn into "running" the same way a // real provider session reaching that state would, without actually // opening one in ProviderSessionManagerV2. No server restart happens @@ -3054,11 +3135,12 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { // already gone: `sessions.get` for it returns none, exactly as it // would after a real release. // - // Also seed an in-flight turn item, a streaming message, and a - // running subagent on the same run, mirroring the in-flight work a - // real provider turn would leave open. None of these are reported - // terminal by a live process, so interrupt finalization must settle - // them itself. + // Also seed an in-flight turn item, a persistent dynamic_tool item, a + // streaming message, and a running subagent (linked to the child + // thread above) on the same run, mirroring the in-flight work a real + // provider turn would leave open. None of these are reported terminal + // by a live process, so interrupt finalization must settle them + // itself. yield* eventSink.write({ commandId: CommandId.make(`${threadId}:force-running`), events: [ @@ -3128,6 +3210,33 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { streaming: true, }, }, + { + id: EventId.make(`${threadId}:dynamic-tool-running`), + type: "turn-item.updated", + threadId, + runId: activeRun.id, + nodeId: activeRun.rootNodeId!, + occurredAt: now, + payload: { + id: dynamicToolItemId, + type: "dynamic_tool", + threadId, + runId: activeRun.id, + nodeId: activeRun.rootNodeId!, + providerThreadId: providerThread.id, + providerTurnId: providerTurn.id, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "running", + title: "Persistent monitor", + startedAt: now, + completedAt: null, + updatedAt: now, + toolName: "persistent_monitor", + input: { persistent: true }, + }, + }, { id: EventId.make(`${threadId}:subagent-running`), type: "subagent.updated", @@ -3147,7 +3256,7 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { driver: providerThread.driver, providerInstanceId: providerThread.providerInstanceId, providerThreadId: null, - childThreadId: null, + childThreadId, nativeTaskRef: null, prompt: "Do the subtask", title: null, @@ -3171,41 +3280,62 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { const after = yield* orchestrator.getThreadProjection(threadId); const interruptedRun = after.runs.find((run) => run.id === activeRun.id); - assert.equal(interruptedRun?.status, "interrupted"); + // The canonical reconciliation path (ProviderRuntimeRecoveryService) + // marks force-cancelled work "cancelled", not "interrupted", and does + // not add a separate run_interrupt_result marker: the run's own + // status already conveys it stopped. + assert.equal(interruptedRun?.status, "cancelled"); assert.equal( after.providerTurns.find( (providerTurn) => providerTurn.runAttemptId === activeRun.activeAttemptId, )?.status, - "interrupted", + "cancelled", ); assert.equal( after.attempts.find((attempt) => attempt.id === activeRun.activeAttemptId)?.status, - "interrupted", + "cancelled", ); assert.equal( after.nodes.find((node) => node.id === activeRun.rootNodeId)?.status, - "interrupted", + "cancelled", ); const streamingItem = after.turnItems.find((item) => item.id === streamingItemId); - assert.equal(streamingItem?.status, "interrupted"); + assert.equal(streamingItem?.status, "cancelled"); assert.equal( streamingItem !== undefined && "streaming" in streamingItem ? streamingItem.streaming : undefined, false, ); + const dynamicToolItem = after.turnItems.find((item) => item.id === dynamicToolItemId); + assert.equal(dynamicToolItem?.status, "cancelled"); assert.equal( after.messages.find((message) => message.id === streamingMessageId)?.streaming, false, ); assert.equal( after.subagents.find((subagent) => subagent.id === subagentId)?.status, - "interrupted", + "cancelled", ); const interruptResult = after.turnItems.find( (item) => item.type === "run_interrupt_result" && item.runId === activeRun.id, ); - assert.isDefined(interruptResult); + assert.isUndefined(interruptResult); + + const [checkpointEffect] = yield* outbox.listByCommandId(checkpointCommandId); + assert.equal(checkpointEffect?.status, "pending"); + assert.deepEqual(checkpointEffect?.request, { + type: "checkpoint.capture", + runId: activeRun.id, + scopeId: checkpointScope.id, + }); + + const childAfter = yield* orchestrator.getThreadProjection(childThreadId); + assert.equal(childAfter.runs.find((run) => run.id === childRun.id)?.status, "cancelled"); + assert.equal( + childAfter.turnItems.find((item) => item.id === childItemId)?.status, + "cancelled", + ); }), ); From 337fb7d2939346b37a4816f3e5172c0e0b7bd67e Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 17:18:17 -0400 Subject: [PATCH 5/8] fix(server): dead-session interrupts no longer touch work they do not own Signed-off-by: Yordis Prieto --- .../src/orchestration-v2/EffectOutbox.ts | 29 +++-- apps/server/src/orchestration-v2/EventSink.ts | 5 +- .../src/orchestration-v2/Orchestrator.ts | 97 +++-------------- .../ProviderRuntimeRecoveryService.ts | 2 +- .../ProviderSessionManager.test.ts | 101 ++++++++++++++++++ .../ProviderSessionManager.ts | 66 +++++++----- .../src/orchestration-v2/runtimeLayer.test.ts | 34 +++++- 7 files changed, 216 insertions(+), 118 deletions(-) diff --git a/apps/server/src/orchestration-v2/EffectOutbox.ts b/apps/server/src/orchestration-v2/EffectOutbox.ts index 4841a5762480..565be0580320 100644 --- a/apps/server/src/orchestration-v2/EffectOutbox.ts +++ b/apps/server/src/orchestration-v2/EffectOutbox.ts @@ -123,6 +123,24 @@ export const PROCESS_BOUND_EFFECT_TYPES = [ "runtime-request.respond", ] as const satisfies ReadonlyArray; +export const SESSION_BOUND_EFFECT_TYPES = [ + "provider-turn.interrupt", + "provider-turn.steer", + "provider-turn.restart", + "runtime-request.respond", +] as const satisfies ReadonlyArray; + +/** + * Which unsettled effects of a thread to cancel. `providerSessionId` narrows + * it to effects addressed to that provider session, leaving effects for other + * sessions on the same thread alone. + */ +export interface UnsettledEffectCancellation { + readonly effectTypes: ReadonlyArray; + readonly reason: string; + readonly providerSessionId?: ProviderSessionId; +} + export const OrchestrationEffectStatusV2 = Schema.Literals([ "pending", "running", @@ -184,11 +202,9 @@ export interface EffectOutboxV2Shape { readonly listByCommandId: ( commandId: CommandId, ) => Effect.Effect, EffectOutboxError>; - readonly cancelUnsettled: (input: { - readonly threadId: ThreadId; - readonly effectTypes: ReadonlyArray; - readonly reason: string; - }) => Effect.Effect, EffectOutboxError>; + readonly cancelUnsettled: ( + input: UnsettledEffectCancellation & { readonly threadId: ThreadId }, + ) => Effect.Effect, EffectOutboxError>; readonly signalCancellations: (effectIds: ReadonlyArray) => Effect.Effect; readonly awaitCancellation: (effectId: string) => Effect.Effect; readonly clearCancellation: (effectId: string) => Effect.Effect; @@ -391,7 +407,7 @@ export const layer: Layer.Layer = La : new EffectOutboxError({ operation: "list", cause }), ), ), - cancelUnsettled: ({ threadId, effectTypes, reason }) => + cancelUnsettled: ({ threadId, effectTypes, reason, providerSessionId }) => Effect.gen(function* () { if (effectTypes.length === 0) return []; const now = DateTime.formatIso(yield* DateTime.now); @@ -407,6 +423,7 @@ export const layer: Layer.Layer = La WHERE thread_id = ${threadId} AND status IN ('pending', 'running') AND effect_type IN ${sql.in(effectTypes)} + ${providerSessionId === undefined ? sql`` : sql`AND json_extract(payload_json, '$.providerSessionId') = ${providerSessionId}`} RETURNING effect_id `; return rows.map(({ effect_id }) => effect_id); diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index c61ce6ba8b7b..6296b922948f 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -124,10 +124,7 @@ export interface EventSinkV2Shape { readonly acceptedAt: DateTime.Utc; readonly events: ReadonlyArray; readonly effects: ReadonlyArray; - readonly cancelUnsettledEffects?: { - readonly effectTypes: ReadonlyArray; - readonly reason: string; - }; + readonly cancelUnsettledEffects?: EffectOutbox.UnsettledEffectCancellation; }) => Effect.Effect< { readonly receipt: CommandReceiptStore.CommandReceiptV2; diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 088c70f442b8..c88920d4e637 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -76,9 +76,9 @@ import { isUndeliveredMailboxSteer } from "./NotificationMailbox.ts"; import { EventSinkV2 } from "./EventSink.ts"; import { EffectOutboxV2, - PROCESS_BOUND_EFFECT_TYPES, - type OrchestrationEffectRequestV2, + SESSION_BOUND_EFFECT_TYPES, type PendingOrchestrationEffectV2, + type UnsettledEffectCancellation, } from "./EffectOutbox.ts"; import { IdAllocatorV2 } from "./IdAllocator.ts"; import { @@ -101,10 +101,7 @@ import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import { ProviderAdapterRegistryV2 } from "./ProviderAdapterRegistry.ts"; import { ProviderContinuationRequests } from "./ProviderContinuationRequests.ts"; import { makeProviderFailure } from "./ProviderFailure.ts"; -import { - planProcessLossReconciliation, - planThreadReconciliation, -} from "./ProviderRuntimeRecoveryService.ts"; +import { planProcessLossReconciliation } from "./ProviderRuntimeRecoveryService.ts"; import { ProviderSessionManagerV2 } from "./ProviderSessionManager.ts"; import { ProviderSwitchServiceV2 } from "./ProviderSwitchService.ts"; import { isAutomaticCompletionRun, queuedRunsInDeliveryOrder } from "./QueuedRunOrder.ts"; @@ -7929,10 +7926,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return { effectTypes: ["provider-turn.start", "provider-turn.restart"], reason: `Run ${run.id} was interrupted before its provider turn started.`, - } satisfies { - readonly effectTypes: ReadonlyArray; - readonly reason: string; - }; + } satisfies UnsettledEffectCancellation; } if (providerTurn === undefined) { @@ -8047,78 +8041,27 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }, ]); } - // A linked subagent child thread is a genuinely separate thread: its - // own open work needs its own reconciliation plan and its own effect - // cancellation, since cancelUnsettledEffects below only scopes the - // primary thread. - const linkedChildThreadIds = recoveryProjection.subagents - .filter((subagent) => subagent.runId === run.id) - .map((subagent) => subagent.childThreadId) - .filter((childThreadId): childThreadId is ThreadId => childThreadId !== null); - for (const childThreadId of linkedChildThreadIds) { - const childProjection = yield* projectionStore - .getRuntimeRecoveryProjection(childThreadId) - .pipe( - Effect.mapError( - (cause) => new OrchestratorProjectionError({ threadId: childThreadId, cause }), - ), - ); - const childPlan = yield* planThreadReconciliation({ - projection: childProjection, - trigger: "process-loss", - continueAfterRestart: false, - commandId: command.commandId, - now, - ids: idAllocator, - outbox, - }).pipe( - Effect.mapError( - (cause) => - new OrchestratorDispatchError({ - commandId: command.commandId, - commandType: command.type, - cause, - }), - ), - ); - yield* Ref.update(events, (existing) => [...existing, ...childPlan.events]); - yield* Ref.update(effects, (existing) => [...existing, ...childPlan.effects]); - if (childPlan.events.length === 0) continue; - const retiredEffectIds = yield* outbox - .cancelUnsettled({ - threadId: childThreadId, - effectTypes: PROCESS_BOUND_EFFECT_TYPES, - reason: childPlan.detail, - }) - .pipe( - Effect.mapError( - (cause) => - new OrchestratorDispatchError({ - commandId: command.commandId, - commandType: command.type, - cause, - }), - ), - ); - yield* outbox.signalCancellations(retiredEffectIds); - } if (command.holdQueue === true) yield* holdQueuedRuns; yield* stopCompletionCohort(); + // Read past the plan's events so a stale provider-thread row cannot + // overwrite the idle state the plan just wrote. yield* settleBackgroundWork({ command, events, - projection, + projection: yield* getProjectionWithPendingEvents(command.threadId, events), stoppedProviderThreadId: providerThread.id, throughRunOrdinal: run.ordinal, now, }); + // Only effects addressed to the dead session lost their process; a + // thread's other sessions keep theirs. return { - effectTypes: PROCESS_BOUND_EFFECT_TYPES, + effectTypes: providerThread.providerSessionId === null ? [] : SESSION_BOUND_EFFECT_TYPES, reason: plan.detail, - } satisfies { - readonly effectTypes: ReadonlyArray; - readonly reason: string; - }; + ...(providerThread.providerSessionId === null + ? {} + : { providerSessionId: providerThread.providerSessionId }), + } satisfies UnsettledEffectCancellation; } if (providerThread.providerSessionId === null) { return yield* new OrchestratorDispatchError({ @@ -9175,10 +9118,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio { readonly events: ReadonlyArray; readonly effects: ReadonlyArray; - readonly cancelUnsettledEffects?: { - readonly effectTypes: ReadonlyArray; - readonly reason: string; - }; + readonly cancelUnsettledEffects?: UnsettledEffectCancellation; }, OrchestratorV2Error > { @@ -9190,12 +9130,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const events = yield* Ref.make>([]); const effects = yield* Ref.make>([]); - let cancelUnsettledEffects: - | { - readonly effectTypes: ReadonlyArray; - readonly reason: string; - } - | undefined; + let cancelUnsettledEffects: UnsettledEffectCancellation | undefined; switch (command.type) { case "thread.create": yield* dispatchThreadCreate(command, events); diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index c2fa867f23ba..35736fa5128b 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -269,7 +269,7 @@ export function scopeProjectionToRun( * callers must scope `projection` first (see `scopeProjectionToRun`), since * this plans against everything the projection contains. */ -export const planThreadReconciliation = Effect.fn( +const planThreadReconciliation = Effect.fn( "ProviderRuntimeRecoveryService.planThreadReconciliation", )(function* (input: { readonly projection: ProjectionStore.ProjectionRuntimeRecoveryState; diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 55360e2f0d77..bda4bbc85e17 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -2813,6 +2813,107 @@ it.effect( }), ); +it.effect( + "ProviderSessionManagerV2 ignores a late turn.terminal from an earlier run on the same provider thread", + () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const projectId = yield* idAllocator.allocate.project({ + fixtureName: "provider-session-manager-late-terminal", + }); + const threadId = yield* idAllocator.allocate.thread({ + fixtureName: "provider-session-manager-late-terminal", + projectId, + }); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + const providerThread = makeProviderThread({ + idAllocator, + threadId, + providerSessionId, + now, + nativeThreadId: "native-thread-late-terminal", + }); + yield* eventSink.write({ + events: [yield* makeThreadCreatedEvent({ idAllocator, threadId, now })], + }); + const runtime = yield* manager.open({ + threadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* runtime.events.pipe(Stream.runDrain, Effect.forkScoped); + const appThread = (yield* projectionStore.getThreadProjection(threadId)).thread; + const startRun = (runOrdinal: number) => + Effect.gen(function* () { + const runId = idAllocator.derive.run({ threadId, ordinal: runOrdinal }); + yield* runtime.startTurn({ + appThread, + threadId, + runId, + runOrdinal, + providerTurnOrdinal: runOrdinal, + attemptId: idAllocator.derive.runAttempt({ runId, attemptOrdinal: 1 }), + rootNodeId: idAllocator.derive.rootNode({ runId }), + providerThread, + message: { + createdBy: "user", + creationSource: "web", + messageId: yield* idAllocator.allocate.message({ threadId, ordinal: runOrdinal }), + text: `run ${runOrdinal}`, + attachments: [], + }, + modelSelection, + runtimePolicy, + }); + }); + const queue = (yield* Ref.get(state)).eventQueues.get(String(providerSessionId)); + assert.isDefined(queue); + const terminal = (runOrdinal: number) => + Queue.offer(queue!, { + type: "turn.terminal", + driver: CODEX_DRIVER, + providerThreadId: providerThread.id, + providerTurnId: idAllocator.derive.providerTurn({ + driver: CODEX_DRIVER, + nativeTurnId: `native-turn-late-terminal-${runOrdinal}`, + }), + runOrdinal, + status: "completed", + failure: null, + threadDisposition: "reusable", + }); + + yield* startRun(1); + yield* terminal(1); + yield* Effect.yieldNow; + yield* startRun(2); + // The provider repeats run 1's terminal while run 2 is in flight. + yield* terminal(1); + yield* TestClock.adjust("2 seconds"); + yield* Effect.yieldNow; + assert.equal((yield* Ref.get(state)).closeCount, 0); + + yield* terminal(2); + yield* Effect.yieldNow; + yield* TestClock.adjust("2 seconds"); + yield* Effect.yieldNow; + assert.equal((yield* Ref.get(state)).closeCount, 1); + }); + + yield* effect.pipe(Effect.provide(makeTestLayer({ state, idleTimeoutMs: 1000 }))); + }), +); + it.effect( "ProviderSessionManagerV2 keeps a running turn busy when an overlapping startTurn on the same provider thread fails", () => diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 24127fb9bfea..29747ba30736 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -206,12 +206,12 @@ interface LiveSessionEntry { readonly scope: Scope.Closeable; readonly idleGeneration: number; /** - * Provider threads with a turn in flight, keyed by provider thread id - * rather than counted: a stray or duplicate `turn.terminal` for a thread - * that never started a turn here (or already settled one) is then a no-op - * instead of zeroing out another thread's genuinely running turn. + * Run ordinal of the turn in flight on each provider thread, keyed rather + * than counted: a stray, duplicate, or late `turn.terminal` (one for a + * thread that never started a turn here, or for an earlier run) is then a + * no-op instead of idling a genuinely running turn. */ - readonly busyProviderThreadIds: ReadonlySet; + readonly busyRunOrdinals: ReadonlyMap; readonly lastActivityAtMs: number; readonly idleFiber: Fiber.Fiber | null; /** Set when idle release is deferred for pending background work; bounds total deferral. */ @@ -732,7 +732,7 @@ export const layerWithOptions = ( } if ( input.onlyIfIdleGeneration !== undefined && - (existing.busyProviderThreadIds.size > 0 || + (existing.busyRunOrdinals.size > 0 || existing.idleGeneration !== input.onlyIfIdleGeneration) ) { return [Option.none(), current] as const; @@ -879,7 +879,7 @@ export const layerWithOptions = ( const entry = current.get(key); if ( entry === undefined || - entry.busyProviderThreadIds.size > 0 || + entry.busyRunOrdinals.size > 0 || entry.idleGeneration !== input.generation ) { return; @@ -901,7 +901,7 @@ export const layerWithOptions = ( const latestEntry = latest.get(key); if ( latestEntry === undefined || - latestEntry.busyProviderThreadIds.size > 0 || + latestEntry.busyRunOrdinals.size > 0 || latestEntry.idleGeneration !== input.generation || latestEntry.runtime !== probedRuntime ) { @@ -933,7 +933,7 @@ export const layerWithOptions = ( } // hasPendingBackgroundWork yields to the adapter, so the idle // decision above can go stale; the generation guard revalidates - // busyProviderThreadIds and idleGeneration inside releaseEntry's + // busyRunOrdinals and idleGeneration inside releaseEntry's // atomic entry removal. yield* releaseEntry({ providerSessionId: input.providerSessionId, @@ -970,7 +970,7 @@ export const layerWithOptions = ( const key = sessionKey(providerSessionId); const current = yield* Ref.get(sessions); const entry = current.get(key); - if (entry === undefined || entry.busyProviderThreadIds.size > 0) { + if (entry === undefined || entry.busyRunOrdinals.size > 0) { return; } @@ -983,7 +983,7 @@ export const layerWithOptions = ( const lastActivityAtMs = yield* Clock.currentTimeMillis; yield* Ref.update(sessions, (latest) => { const latestEntry = latest.get(key); - if (latestEntry === undefined || latestEntry.busyProviderThreadIds.size > 0) { + if (latestEntry === undefined || latestEntry.busyRunOrdinals.size > 0) { return latest; } const updated = new Map(latest); @@ -1168,7 +1168,11 @@ export const layerWithOptions = ( ); }); - const markBusy = (providerSessionId: ProviderSessionId, providerThreadId: ProviderThreadId) => + const markBusy = ( + providerSessionId: ProviderSessionId, + providerThreadId: ProviderThreadId, + runOrdinal: number, + ) => withActivityError( providerSessionId, Effect.gen(function* () { @@ -1180,17 +1184,20 @@ export const layerWithOptions = ( return [[null, false] as const, current] as const; } const updated = new Map(current); - const busyProviderThreadIds = new Set(entry.busyProviderThreadIds); - busyProviderThreadIds.add(providerThreadId); + const busyRunOrdinals = new Map(entry.busyRunOrdinals); + busyRunOrdinals.set( + providerThreadId, + Math.max(runOrdinal, entry.busyRunOrdinals.get(providerThreadId) ?? runOrdinal), + ); updated.set(key, { ...entry, - busyProviderThreadIds, + busyRunOrdinals, idleFiber: null, lastActivityAtMs: now, pinnedSinceMs: null, }); return [ - [entry.idleFiber, !entry.busyProviderThreadIds.has(providerThreadId)] as const, + [entry.idleFiber, !entry.busyRunOrdinals.has(providerThreadId)] as const, updated, ] as const; }); @@ -1199,7 +1206,11 @@ export const layerWithOptions = ( }), ); - const markIdle = (providerSessionId: ProviderSessionId, providerThreadId: ProviderThreadId) => + const markIdle = ( + providerSessionId: ProviderSessionId, + providerThreadId: ProviderThreadId, + runOrdinal: number, + ) => withActivityError( providerSessionId, Effect.gen(function* () { @@ -1211,11 +1222,14 @@ export const layerWithOptions = ( return current; } const updated = new Map(current); - const busyProviderThreadIds = new Set(entry.busyProviderThreadIds); - busyProviderThreadIds.delete(providerThreadId); + const busyRunOrdinals = new Map(entry.busyRunOrdinals); + const busyRunOrdinal = busyRunOrdinals.get(providerThreadId); + if (busyRunOrdinal !== undefined && busyRunOrdinal <= runOrdinal) { + busyRunOrdinals.delete(providerThreadId); + } updated.set(key, { ...entry, - busyProviderThreadIds, + busyRunOrdinals, lastActivityAtMs: now, }); return updated; @@ -1390,7 +1404,7 @@ export const layerWithOptions = ( }), ).pipe( Effect.andThen( - markBusy(providerSessionId, input.providerThread.id).pipe( + markBusy(providerSessionId, input.providerThread.id, input.runOrdinal).pipe( Effect.catchCause((cause) => Effect.logWarning("orchestration-v2.driver-session.activity-failed", { providerSessionId, @@ -1409,7 +1423,7 @@ export const layerWithOptions = ( (markedBusy ? observeActivity( providerSessionId, - markIdle(providerSessionId, input.providerThread.id), + markIdle(providerSessionId, input.providerThread.id, input.runOrdinal), ) : Effect.void ).pipe(Effect.andThen(Effect.fail(error))), @@ -1471,7 +1485,11 @@ export const layerWithOptions = ( return observeActivity( entry.runtime.providerSessionId, event.type === "turn.terminal" - ? markIdle(entry.runtime.providerSessionId, event.providerThreadId) + ? markIdle( + entry.runtime.providerSessionId, + event.providerThreadId, + event.runOrdinal, + ) : touchActivity(entry.runtime.providerSessionId), ).pipe( Effect.andThen( @@ -1714,7 +1732,7 @@ export const layerWithOptions = ( requestEventPermit: yield* Semaphore.make(1), scope: sessionScope, idleGeneration: 0, - busyProviderThreadIds: new Set(), + busyRunOrdinals: new Map(), lastActivityAtMs: now, idleFiber: null, pinnedSinceMs: null, diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 0a57f57771eb..d2aa88fb9210 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -12,6 +12,7 @@ import { EventId, MessageId, NodeId, + ProviderSessionId, RuntimeRequestId, TurnItemId, type ModelSelection, @@ -3125,6 +3126,27 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { const checkpointCommandId = CommandId.make( `command:effect:checkpoint.capture:${activeRun.id}`, ); + // Effects addressed to the dead session lose their process; one for + // another session on the same thread does not. + const sessionRespond = (name: string, providerSessionId: ProviderSessionId) => ({ + id: `effect:${threadId}:${name}`, + commandId: CommandId.make(`${threadId}:${name}`), + threadId, + request: { + type: "runtime-request.respond" as const, + providerSessionId, + requestId: RuntimeRequestId.make(`${threadId}:${name}`), + }, + }); + const deadSessionRespond = sessionRespond( + "dead-session", + providerThread.providerSessionId!, + ); + const otherSessionRespond = sessionRespond( + "other-session", + ProviderSessionId.make(`${threadId}:other-session`), + ); + yield* outbox.enqueue([deadSessionRespond, otherSessionRespond]); // Drive the run and its provider turn into "running" the same way a // real provider session reaching that state would, without actually @@ -3330,11 +3352,19 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { scopeId: checkpointScope.id, }); + const [deadSessionEffect] = yield* outbox.listByCommandId(deadSessionRespond.commandId); + assert.equal(deadSessionEffect?.status, "cancelled"); + const [otherSessionEffect] = yield* outbox.listByCommandId(otherSessionRespond.commandId); + assert.equal(otherSessionEffect?.status, "pending"); + + // Like a live interrupt, this one stays on its own thread: the child + // may run on a session that is still alive, and its own interrupt + // settles it if not. const childAfter = yield* orchestrator.getThreadProjection(childThreadId); - assert.equal(childAfter.runs.find((run) => run.id === childRun.id)?.status, "cancelled"); + assert.equal(childAfter.runs.find((run) => run.id === childRun.id)?.status, "running"); assert.equal( childAfter.turnItems.find((item) => item.id === childItemId)?.status, - "cancelled", + "running", ); }), ); From c2141a6b3716b1ae19e5f862a29f4a29b42cf488 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 17:32:19 -0400 Subject: [PATCH 6/8] fix(server): a failed overlapping turn start no longer pins the session busy Signed-off-by: Yordis Prieto --- .../ProviderSessionManager.test.ts | 22 +++++++ .../ProviderSessionManager.ts | 64 ++++++++++++++----- 2 files changed, 71 insertions(+), 15 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index bda4bbc85e17..fd26276da6c1 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -2986,6 +2986,28 @@ it.effect( yield* TestClock.adjust("2 seconds"); yield* Effect.yieldNow; assert.equal((yield* Ref.get(state)).closeCount, 0); + + // The failed run 2 must not leave its ordinal behind: run 1's + // terminal still idles the thread and lets the session release. + const queue = (yield* Ref.get(state)).eventQueues.get(String(providerSessionId)); + assert.isDefined(queue); + yield* Queue.offer(queue!, { + type: "turn.terminal", + driver: CODEX_DRIVER, + providerThreadId: providerThread.id, + providerTurnId: idAllocator.derive.providerTurn({ + driver: CODEX_DRIVER, + nativeTurnId: "native-turn-overlap-1", + }), + runOrdinal: 1, + status: "completed", + failure: null, + threadDisposition: "reusable", + }); + yield* Effect.yieldNow; + yield* TestClock.adjust("2 seconds"); + yield* Effect.yieldNow; + assert.equal((yield* Ref.get(state)).closeCount, 1); }); yield* effect.pipe( diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 29747ba30736..1015bb630672 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -185,6 +185,11 @@ export class ProviderSessionManagerV2 extends Context.Service< ProviderSessionManagerV2Shape >()("t3/orchestration-v2/ProviderSessionManager/ProviderSessionManagerV2") {} +interface BusyMark { + readonly previousRunOrdinal: number | null; + readonly runOrdinal: number; +} + interface LiveSessionEntry { readonly attachedThreadIds: ReadonlySet; readonly loadedProviderThreadKeyByThread: ReadonlyMap; @@ -1178,10 +1183,10 @@ export const layerWithOptions = ( Effect.gen(function* () { const key = sessionKey(providerSessionId); const now = yield* Clock.currentTimeMillis; - const [idleFiber, added] = yield* Ref.modify(sessions, (current) => { + const [idleFiber, mark] = yield* Ref.modify(sessions, (current) => { const entry = current.get(key); if (entry === undefined) { - return [[null, false] as const, current] as const; + return [[null, null] as const, current] as const; } const updated = new Map(current); const busyRunOrdinals = new Map(entry.busyRunOrdinals); @@ -1196,13 +1201,14 @@ export const layerWithOptions = ( lastActivityAtMs: now, pinnedSinceMs: null, }); - return [ - [entry.idleFiber, !entry.busyRunOrdinals.has(providerThreadId)] as const, - updated, - ] as const; + const mark: BusyMark = { + previousRunOrdinal: entry.busyRunOrdinals.get(providerThreadId) ?? null, + runOrdinal: busyRunOrdinals.get(providerThreadId)!, + }; + return [[entry.idleFiber, mark] as const, updated] as const; }); yield* cancelIdleFiber(idleFiber); - return added; + return mark; }), ); @@ -1238,6 +1244,34 @@ export const layerWithOptions = ( }), ); + // Undoes a markBusy whose turn failed to start, unless a terminal or a + // newer turn has since changed the thread's busy state. + const revertBusy = ( + providerSessionId: ProviderSessionId, + providerThreadId: ProviderThreadId, + mark: BusyMark, + ) => + withActivityError( + providerSessionId, + Effect.gen(function* () { + const key = sessionKey(providerSessionId); + yield* Ref.update(sessions, (current) => { + const entry = current.get(key); + if (entry?.busyRunOrdinals.get(providerThreadId) !== mark.runOrdinal) { + return current; + } + const busyRunOrdinals = new Map(entry.busyRunOrdinals); + if (mark.previousRunOrdinal === null) { + busyRunOrdinals.delete(providerThreadId); + } else { + busyRunOrdinals.set(providerThreadId, mark.previousRunOrdinal); + } + return new Map(current).set(key, { ...entry, busyRunOrdinals }); + }); + yield* scheduleIdleReleaseInternal(providerSessionId); + }), + ); + const observeActivity = ( providerSessionId: ProviderSessionId, activity: Effect.Effect, @@ -1409,23 +1443,23 @@ export const layerWithOptions = ( Effect.logWarning("orchestration-v2.driver-session.activity-failed", { providerSessionId, cause, - }).pipe(Effect.as(false)), + }).pipe(Effect.as(null)), ), ), ), - // Only the startTurn that marked the thread busy may clear it on - // failure; an overlapping attempt must not idle a running turn. - Effect.flatMap((markedBusy) => + // A failed start restores the busy state it found, so an + // overlapping attempt neither idles nor pins a running turn. + Effect.flatMap((mark) => runtime .startTurn(input) .pipe( Effect.catch((error) => - (markedBusy - ? observeActivity( + (mark === null + ? Effect.void + : observeActivity( providerSessionId, - markIdle(providerSessionId, input.providerThread.id, input.runOrdinal), + revertBusy(providerSessionId, input.providerThread.id, mark), ) - : Effect.void ).pipe(Effect.andThen(Effect.fail(error))), ), ), From 1c86144e8e401bb10a52331d5babb6358015d4dc Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 17:33:27 -0400 Subject: [PATCH 7/8] fix(server): busy marking typechecks against the session map Signed-off-by: Yordis Prieto --- .../ProviderSessionManager.ts | 56 +++++++++++-------- 1 file changed, 32 insertions(+), 24 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 1015bb630672..09adbdb684e6 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -1183,30 +1183,38 @@ export const layerWithOptions = ( Effect.gen(function* () { const key = sessionKey(providerSessionId); const now = yield* Clock.currentTimeMillis; - const [idleFiber, mark] = yield* Ref.modify(sessions, (current) => { - const entry = current.get(key); - if (entry === undefined) { - return [[null, null] as const, current] as const; - } - const updated = new Map(current); - const busyRunOrdinals = new Map(entry.busyRunOrdinals); - busyRunOrdinals.set( - providerThreadId, - Math.max(runOrdinal, entry.busyRunOrdinals.get(providerThreadId) ?? runOrdinal), - ); - updated.set(key, { - ...entry, - busyRunOrdinals, - idleFiber: null, - lastActivityAtMs: now, - pinnedSinceMs: null, - }); - const mark: BusyMark = { - previousRunOrdinal: entry.busyRunOrdinals.get(providerThreadId) ?? null, - runOrdinal: busyRunOrdinals.get(providerThreadId)!, - }; - return [[entry.idleFiber, mark] as const, updated] as const; - }); + const [idleFiber, mark] = yield* Ref.modify( + sessions, + ( + current, + ): readonly [ + readonly [LiveSessionEntry["idleFiber"], BusyMark | null], + Map, + ] => { + const entry = current.get(key); + if (entry === undefined) { + return [[null, null], current]; + } + const updated = new Map(current); + const busyRunOrdinals = new Map(entry.busyRunOrdinals); + busyRunOrdinals.set( + providerThreadId, + Math.max(runOrdinal, entry.busyRunOrdinals.get(providerThreadId) ?? runOrdinal), + ); + updated.set(key, { + ...entry, + busyRunOrdinals, + idleFiber: null, + lastActivityAtMs: now, + pinnedSinceMs: null, + }); + const mark: BusyMark = { + previousRunOrdinal: entry.busyRunOrdinals.get(providerThreadId) ?? null, + runOrdinal: busyRunOrdinals.get(providerThreadId)!, + }; + return [[entry.idleFiber, mark], updated]; + }, + ); yield* cancelIdleFiber(idleFiber); return mark; }), From 286eb4f0cb0b4a2adb048c674ba64dfc9afaa7ba Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Sat, 3 Oct 2026 17:43:13 -0400 Subject: [PATCH 8/8] fix(server): busy tracking settles only the runs each event is about Signed-off-by: Yordis Prieto --- .../ProviderRuntimeRecoveryService.ts | 2 +- .../ProviderSessionManager.test.ts | 116 ++++++++++++++ .../ProviderSessionManager.ts | 142 +++++++----------- 3 files changed, 168 insertions(+), 92 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index 35736fa5128b..27a52b33a114 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -222,7 +222,7 @@ function reconciliationRequestReason(trigger: ReconciliationTrigger): string { * out entirely, so a "process-loss" plan built from this view cannot cancel * work that is still alive elsewhere on the thread. */ -export function scopeProjectionToRun( +function scopeProjectionToRun( projection: ProjectionStore.ProjectionRuntimeRecoveryState, input: { readonly runId: RunId; readonly providerThreadId: ProviderThreadId }, ): ProjectionStore.ProjectionRuntimeRecoveryState { diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts index 352579787446..f9fbc18d2743 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.test.ts @@ -13,6 +13,7 @@ import { ProjectId, ProviderDriverKind, ProviderInstanceId, + ProviderTurnId, type ProviderSessionId, ThreadId, } from "@t3tools/contracts"; @@ -3079,6 +3080,121 @@ it.effect( }), ); +it.effect( + "ProviderSessionManagerV2 releases the session when the running turn ends while an overlapping startTurn fails", + () => + Effect.gen(function* () { + const state = yield* Ref.make(emptyState); + const startTurnCalls = yield* Ref.make(0); + const effect = Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + const manager = yield* ProviderSessionManager.ProviderSessionManagerV2; + const projectionStore = yield* ProjectionStore.ProjectionStoreV2; + const now = yield* DateTime.now; + const projectId = yield* idAllocator.allocate.project({ + fixtureName: "provider-session-manager-overlap-race", + }); + const threadId = yield* idAllocator.allocate.thread({ + fixtureName: "provider-session-manager-overlap-race", + projectId, + }); + const providerSessionId = yield* idAllocator.allocate.providerSession({ + providerInstanceId: modelSelection.instanceId, + threadId, + }); + const providerThread = makeProviderThread({ + idAllocator, + threadId, + providerSessionId, + now, + nativeThreadId: "native-thread-overlap-race", + }); + yield* eventSink.write({ + events: [yield* makeThreadCreatedEvent({ idAllocator, threadId, now })], + }); + const runtime = yield* manager.open({ + threadId, + providerSessionId, + modelSelection, + runtimePolicy, + }); + yield* runtime.events.pipe(Stream.runDrain, Effect.forkScoped); + const appThread = (yield* projectionStore.getThreadProjection(threadId)).thread; + const startInput = (ordinal: number) => + Effect.gen(function* () { + const runId = idAllocator.derive.run({ threadId, ordinal }); + return { + appThread, + threadId, + runId, + runOrdinal: ordinal, + providerTurnOrdinal: ordinal, + attemptId: idAllocator.derive.runAttempt({ runId, attemptOrdinal: 1 }), + rootNodeId: idAllocator.derive.rootNode({ runId }), + providerThread, + message: { + createdBy: "user" as const, + creationSource: "web" as const, + messageId: yield* idAllocator.allocate.message({ threadId, ordinal }), + text: `turn ${ordinal}`, + attachments: [], + }, + modelSelection, + runtimePolicy, + }; + }); + + yield* runtime.startTurn(yield* startInput(1)); + const overlapping = yield* runtime.startTurn(yield* startInput(2)).pipe(Effect.exit); + assert.isTrue(overlapping._tag === "Failure"); + + yield* TestClock.adjust("2 seconds"); + yield* Effect.yieldNow; + assert.equal((yield* Ref.get(state)).closeCount, 1); + }); + + yield* effect.pipe( + Effect.provide( + makeTestLayer({ + state, + idleTimeoutMs: 1000, + // Run 1's terminal lands while run 2's start is still failing. + startTurn: (input) => + Ref.getAndUpdate(startTurnCalls, (calls) => calls + 1).pipe( + Effect.flatMap((calls) => + calls === 0 + ? Effect.void + : Ref.get(state).pipe( + Effect.flatMap((current) => + Queue.offer( + current.eventQueues.get( + String(input.providerThread.providerSessionId), + )!, + { + type: "turn.terminal", + driver: CODEX_DRIVER, + providerThreadId: input.providerThread.id, + providerTurnId: ProviderTurnId.make("native-turn-overlap-race-1"), + runOrdinal: 1, + status: "completed", + failure: null, + threadDisposition: "reusable", + }, + ), + ), + Effect.andThen(Effect.yieldNow), + Effect.andThen(Effect.yieldNow), + Effect.andThen(unimplemented("overlapping turn rejected")), + ), + ), + ), + }), + ), + ); + }), +); + it.effect( "ProviderSessionManagerV2 opens one shared runtime, broadcasts events, and detaches threads independently", () => diff --git a/apps/server/src/orchestration-v2/ProviderSessionManager.ts b/apps/server/src/orchestration-v2/ProviderSessionManager.ts index 88a2939a8962..c66f5d31654d 100644 --- a/apps/server/src/orchestration-v2/ProviderSessionManager.ts +++ b/apps/server/src/orchestration-v2/ProviderSessionManager.ts @@ -187,11 +187,6 @@ export class ProviderSessionManagerV2 extends Context.Service< ProviderSessionManagerV2Shape >()("t3/orchestration-v2/ProviderSessionManager/ProviderSessionManagerV2") {} -interface BusyMark { - readonly previousRunOrdinal: number | null; - readonly runOrdinal: number; -} - interface LiveSessionEntry { readonly attachedThreadIds: ReadonlySet; readonly loadedProviderThreadKeyByThread: ReadonlyMap; @@ -213,12 +208,12 @@ interface LiveSessionEntry { readonly scope: Scope.Closeable; readonly idleGeneration: number; /** - * Run ordinal of the turn in flight on each provider thread, keyed rather - * than counted: a stray, duplicate, or late `turn.terminal` (one for a - * thread that never started a turn here, or for an earlier run) is then a - * no-op instead of idling a genuinely running turn. + * Run ordinals with a turn in flight, per provider thread; a thread with + * none is absent. Tracking runs rather than a count makes a stray, + * duplicate, or late `turn.terminal`, or a failed overlapping start, settle + * only the runs it is about instead of idling a genuinely running turn. */ - readonly busyRunOrdinals: ReadonlyMap; + readonly busyRunOrdinals: ReadonlyMap>; readonly lastActivityAtMs: number; readonly idleFiber: Fiber.Fiber | null; /** Set when idle release is deferred for pending background work; bounds total deferral. */ @@ -1193,98 +1188,59 @@ export const layerWithOptions = ( Effect.gen(function* () { const key = sessionKey(providerSessionId); const now = yield* Clock.currentTimeMillis; - const [idleFiber, mark] = yield* Ref.modify( - sessions, - ( - current, - ): readonly [ - readonly [LiveSessionEntry["idleFiber"], BusyMark | null], - Map, - ] => { - const entry = current.get(key); - if (entry === undefined) { - return [[null, null], current]; - } - const updated = new Map(current); - const busyRunOrdinals = new Map(entry.busyRunOrdinals); - busyRunOrdinals.set( - providerThreadId, - Math.max(runOrdinal, entry.busyRunOrdinals.get(providerThreadId) ?? runOrdinal), - ); - updated.set(key, { - ...entry, - busyRunOrdinals, - idleFiber: null, - lastActivityAtMs: now, - pinnedSinceMs: null, - }); - const mark: BusyMark = { - previousRunOrdinal: entry.busyRunOrdinals.get(providerThreadId) ?? null, - runOrdinal: busyRunOrdinals.get(providerThreadId)!, - }; - return [[entry.idleFiber, mark], updated]; - }, - ); - yield* cancelIdleFiber(idleFiber); - return mark; - }), - ); - - const markIdle = ( - providerSessionId: ProviderSessionId, - providerThreadId: ProviderThreadId, - runOrdinal: number, - ) => - withActivityError( - providerSessionId, - Effect.gen(function* () { - const key = sessionKey(providerSessionId); - const now = yield* Clock.currentTimeMillis; - yield* Ref.update(sessions, (current) => { + const [idleFiber, added] = yield* Ref.modify(sessions, (current) => { const entry = current.get(key); + const runs = entry?.busyRunOrdinals.get(providerThreadId); if (entry === undefined) { - return current; - } - const updated = new Map(current); - const busyRunOrdinals = new Map(entry.busyRunOrdinals); - const busyRunOrdinal = busyRunOrdinals.get(providerThreadId); - if (busyRunOrdinal !== undefined && busyRunOrdinal <= runOrdinal) { - busyRunOrdinals.delete(providerThreadId); + return [[null, false] as const, current] as const; } - updated.set(key, { + const busyRunOrdinals = new Map(entry.busyRunOrdinals).set( + providerThreadId, + new Set(runs).add(runOrdinal), + ); + const updated = new Map(current).set(key, { ...entry, busyRunOrdinals, + idleFiber: null, lastActivityAtMs: now, + pinnedSinceMs: null, }); - return updated; + return [[entry.idleFiber, runs?.has(runOrdinal) !== true] as const, updated] as const; }); - yield* scheduleIdleReleaseInternal(providerSessionId); + yield* cancelIdleFiber(idleFiber); + return added; }), ); - // Undoes a markBusy whose turn failed to start, unless a terminal or a - // newer turn has since changed the thread's busy state. - const revertBusy = ( + const markIdle = ( providerSessionId: ProviderSessionId, providerThreadId: ProviderThreadId, - mark: BusyMark, + settles: (runOrdinal: number) => boolean, ) => withActivityError( providerSessionId, Effect.gen(function* () { const key = sessionKey(providerSessionId); + const now = yield* Clock.currentTimeMillis; yield* Ref.update(sessions, (current) => { const entry = current.get(key); - if (entry?.busyRunOrdinals.get(providerThreadId) !== mark.runOrdinal) { + if (entry === undefined) { return current; } + const remaining = [...(entry.busyRunOrdinals.get(providerThreadId) ?? [])].filter( + (runOrdinal) => !settles(runOrdinal), + ); const busyRunOrdinals = new Map(entry.busyRunOrdinals); - if (mark.previousRunOrdinal === null) { + if (remaining.length === 0) { busyRunOrdinals.delete(providerThreadId); } else { - busyRunOrdinals.set(providerThreadId, mark.previousRunOrdinal); + busyRunOrdinals.set(providerThreadId, new Set(remaining)); } - return new Map(current).set(key, { ...entry, busyRunOrdinals }); + return new Map(current).set(key, { + ...entry, + busyRunOrdinals, + lastActivityAtMs: now, + }); }); yield* scheduleIdleReleaseInternal(providerSessionId); }), @@ -1461,26 +1417,28 @@ export const layerWithOptions = ( Effect.logWarning("orchestration-v2.driver-session.activity-failed", { providerSessionId, cause, - }).pipe(Effect.as(null)), + }).pipe(Effect.as(false)), ), ), ), - // A failed start restores the busy state it found, so an + // A failed start settles only the run it added, so an // overlapping attempt neither idles nor pins a running turn. - Effect.flatMap((mark) => - runtime - .startTurn(input) - .pipe( - Effect.catch((error) => - (mark === null - ? Effect.void - : observeActivity( + Effect.flatMap((added) => + runtime.startTurn(input).pipe( + Effect.catch((error) => + (added + ? observeActivity( + providerSessionId, + markIdle( providerSessionId, - revertBusy(providerSessionId, input.providerThread.id, mark), - ) - ).pipe(Effect.andThen(Effect.fail(error))), - ), + input.providerThread.id, + (runOrdinal) => runOrdinal === input.runOrdinal, + ), + ) + : Effect.void + ).pipe(Effect.andThen(Effect.fail(error))), ), + ), ), ), steerTurn: (input) => @@ -1540,7 +1498,9 @@ export const layerWithOptions = ( ? markIdle( entry.runtime.providerSessionId, event.providerThreadId, - event.runOrdinal, + // A run's terminal also settles any earlier run whose own + // terminal never arrived. + (runOrdinal) => runOrdinal <= event.runOrdinal, ) : touchActivity(entry.runtime.providerSessionId), ).pipe(