From 1f7c37e578eeef10a38f328a2e3811ef53b04bb0 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 13:09:26 -0700 Subject: [PATCH 01/11] fix(server): a restart that never starts settles the work its run inherited A steer that changes the selection restarts the run as a new attempt, and only the root turn is superseded. When the replacement session can't open or the provider thread can't load, the last start attempt failed just the run, its attempt and root node, leaving the approvals, streaming replies, background commands and provider-native subagents (and their child threads) the run inherited open forever. The failure path now builds the open work from the read model with openRunOwnedWorkFromProjection, following provider-native subagent links into their child threads, and settles it through RunExecutionService's cascadeTerminalizeRunOwnedSubagents inside the same guarded writeIfRunCurrent commit. App-owned delegate_task tasks keep running. Rebased onto main as one commit (originally 003de31898..76da7ae9fa). The root node exclusion moved from the shared cascade into openRunOwnedWorkFromProjection, because main's Stop recovery (#15442) relies on the cascade interrupting the root node. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../ProviderTurnStartService.test.ts | 1 + .../ProviderTurnStartService.ts | 102 +++- .../orchestration-v2/RunExecutionService.ts | 148 ++++- .../SelectionRestart.integration.test.ts | 570 +++++++++++++++++- 4 files changed, 817 insertions(+), 4 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index 06761c13b3fd..b4c064b1bd0c 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -513,6 +513,7 @@ function makeLocalCommandHarness(input: { (m.text.trim().toLowerCase() !== "/compact" || m.attachments.length > 0), ), }), + getThreadRecords: () => Effect.succeed(projection as never), getTurnStartHistory: () => Effect.fail( new ProjectionStore.ProjectionStoreReadError({ diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 086abc6c58eb..d1b005f18826 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -2,6 +2,7 @@ import { modelSelectionsEqual } from "@t3tools/shared/model"; import { projectComposerContextForProvider } from "@t3tools/shared/composerContextReferences"; import { CommandId, + isProviderNativeSubagentThread, latestProviderTurnForAttempt, type OrchestrationV2DomainEvent, type OrchestrationV2ExecutionNode, @@ -177,6 +178,99 @@ export const layer: Layer.Layer< }; }; + /** + * Ends the work a run that never reached its provider still shows. A + * restart supersedes only the root turn, so requests, streaming replies and + * background work an earlier attempt left open stay with the run, and once + * it fails nothing will report on them. The run's projection rows feed + * RunExecutionService's terminal cascade, the same one a started run uses. + */ + const inheritedWorkSettlement = Effect.fn( + "orchestrationV2.providerTurnStart.inheritedWorkSettlement", + )(function* (input: { readonly run: OrchestrationV2Run; readonly now: DateTime.Utc }) { + const { run } = input; + // Lifetime links, settled rows included: a subagent row can settle before + // its child thread does. App-owned tasks (delegate_task) run on their own. + const nativeLinks = (threadId: ThreadId) => + projectionStore + .getThreadRecords(threadId, ["subagents", "turnItems"], { + turnItemTypes: ["subagent"], + ...(threadId === run.threadId ? { turnItemRunIds: [run.id] } : {}), + }) + .pipe( + Effect.map((records) => + [...records.subagents, ...records.turnItems].flatMap((row) => + (threadId !== run.threadId || row.runId === run.id) && + "childThreadId" in row && + row.childThreadId !== null && + row.origin === "provider_native" + ? [row.childThreadId] + : [], + ), + ), + ); + const pending = [...(yield* nativeLinks(run.threadId))]; + const linkedChildThreadIds = new Set(); + const threads = [yield* projectionStore.getRuntimeRecoveryProjection(run.threadId)]; + for (let threadId = pending.shift(); threadId !== undefined; threadId = pending.shift()) { + if (threadId === run.threadId || linkedChildThreadIds.has(threadId)) continue; + const child = yield* projectionStore.getRuntimeRecoveryProjection(threadId).pipe( + Effect.map(Option.some), + Effect.catchTags({ + ProjectionStoreThreadNotFoundError: () => Effect.succeed(Option.none()), + }), + ); + if (Option.isNone(child) || !isProviderNativeSubagentThread(child.value.thread)) continue; + linkedChildThreadIds.add(threadId); + // A child thread has no runs, so the recovery read leaves out its + // streaming replies and its pending requests' nodes and items; read + // them by thread instead. + const records = yield* projectionStore.getThreadRecords( + threadId, + ["messages", "nodes", "turnItems"], + { + messageRoles: ["assistant"], + turnItemTypes: ["approval_request", "user_input_request"], + turnItemStatuses: ["pending", "running", "waiting"], + }, + ); + const requestNodeIds = new Set( + child.value.runtimeRequests.flatMap((request) => + request.status === "pending" ? [request.nodeId] : [], + ), + ); + const knownNodeIds = new Set(child.value.nodes.map((node) => node.id)); + const knownItemIds = new Set(child.value.turnItems.map((item) => item.id)); + threads.push({ + ...child.value, + nodes: [ + ...child.value.nodes, + ...records.nodes.filter( + (node) => requestNodeIds.has(node.id) && !knownNodeIds.has(node.id), + ), + ], + turnItems: [ + ...child.value.turnItems, + ...records.turnItems.filter((item) => !knownItemIds.has(item.id)), + ], + messages: records.messages.filter((message) => message.streaming), + }); + pending.push(...(yield* nativeLinks(threadId))); + } + const linked = RunExecutionService.openRunOwnedWorkFromProjection({ + run, + threads, + linkedChildThreadIds, + }); + return yield* RunExecutionService.cascadeTerminalizeRunOwnedSubagents({ + run, + open: linked, + status: "cancelled", + completedAt: input.now, + allocateEventId: () => idAllocator.allocate.event({ threadId: run.threadId }), + }); + }); + const makeDeliverySession = ( session: ProviderAdapterV2SessionRuntime, startWithHandoffs: ( @@ -365,7 +459,13 @@ export const layer: Layer.Layer< runId, activeAttemptId: attempt.id, expectedStatus: "starting", - events, + // A user answer that lands while the failure is written wins over + // cancelling the request it answered. + guardPendingUserInputCancellations: true, + events: + status === "failed" + ? [...(yield* inheritedWorkSettlement({ run, now })), ...events] + : events, }); }, ); diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index 254b810c6daf..33e6d92b09f8 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -8,6 +8,7 @@ import { type NodeId, type OrchestrationV2AppThread, type OrchestrationV2CheckpointScope, + type OrchestrationV2ConversationMessage, type OrchestrationV2DomainEvent, type OrchestrationV2ExecutionNode, type OrchestrationV2ProviderFailure, @@ -15,6 +16,7 @@ import { type OrchestrationV2ProviderTurn, type OrchestrationV2Run, type OrchestrationV2RunAttempt, + type OrchestrationV2RuntimeRequest, type OrchestrationV2Subagent, type OrchestrationV2TurnItem, type ProviderSessionId, @@ -40,6 +42,7 @@ import * as ServerSettings from "../serverSettings.ts"; import * as CheckpointService from "./CheckpointService.ts"; import * as EventSink from "./EventSink.ts"; import * as IdAllocator from "./IdAllocator.ts"; +import type * as ProjectionStore from "./ProjectionStore.ts"; import type { ProviderAdapterV2Event, ProviderAdapterV2RuntimePolicy, @@ -146,13 +149,21 @@ export function selectInheritedBackgroundTurnItems(input: { type SubagentTurnItem = Extract; -type OpenRunOwnedSubagentProjection = { +export type OpenRunOwnedSubagentProjection = { readonly subagents: ReadonlyMap; readonly turnItems: ReadonlyMap; + /** Open items on linked child threads, and on the run's own thread for the run. */ readonly childTurnItems: ReadonlyMap; readonly nodes: ReadonlyMap; /** Child threads once linked by a root-run subagent row; kept for cascade. */ readonly linkedChildThreadIds: ReadonlySet; + // A live run's provider reports these itself. A run that never reached its + // provider has nothing to report them, so its read-model rows are passed in. + /** Pending requests; each settles with the node it waits on. */ + readonly runtimeRequests?: ReadonlyArray; + readonly streamingMessages?: ReadonlyArray; + /** Open provider turns of the run's own attempts. */ + readonly providerTurns?: ReadonlyArray; }; type RunOwnedSubagentTerminalStatus = Extract< @@ -202,6 +213,10 @@ function withLinkedChildThreadId( return { ...current, linkedChildThreadIds }; } +/** + * What a run that ends interrupted, failed or cancelled settles besides its + * root: the run's open work and the provider-native child threads it linked. + */ export function cascadeTerminalizeRunOwnedSubagents(input: { readonly run: OrchestrationV2Run; readonly open: OpenRunOwnedSubagentProjection; @@ -222,6 +237,7 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { childThreadIds.add(item.childThreadId); } } + const settledNodeIds = new Set(); const keys = new Set([ ...input.open.subagents.keys(), ...input.open.turnItems.keys(), @@ -254,6 +270,7 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { childThreadIds.has(node.threadId)) && isOpenExecutionNodeStatus(node.status) ) { + settledNodeIds.add(node.id); events.push({ id: yield* input.allocateEventId(), type: "node.updated", @@ -292,8 +309,11 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { }); } } + const ownsThreadRow = (row: { readonly threadId: ThreadId; readonly runId: string | null }) => + childThreadIds.has(row.threadId) || + (row.threadId === input.run.threadId && row.runId === input.run.id); for (const turnItem of input.open.childTurnItems.values()) { - if (!childThreadIds.has(turnItem.threadId) || isSettledTurnItemStatus(turnItem.status)) { + if (!ownsThreadRow(turnItem) || isSettledTurnItemStatus(turnItem.status)) { continue; } events.push({ @@ -313,10 +333,134 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { }, }); } + for (const request of input.open.runtimeRequests ?? []) { + if (request.status !== "pending" || !settledNodeIds.has(request.nodeId)) continue; + const node = input.open.nodes.get(request.nodeId); + events.push({ + id: yield* input.allocateEventId(), + type: "runtime-request.updated", + threadId: node?.threadId ?? input.run.threadId, + nodeId: request.nodeId, + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { + ...request, + status: "cancelled", + responseCapability: { + type: "not_resumable", + reason: "The run ended before this request was resolved.", + }, + resolvedAt: input.completedAt, + }, + }); + } + for (const message of input.open.streamingMessages ?? []) { + if (!message.streaming || !ownsThreadRow(message)) continue; + events.push({ + id: yield* input.allocateEventId(), + type: "message.updated", + threadId: message.threadId, + runId: message.runId ?? input.run.id, + ...(message.nodeId === null ? {} : { nodeId: message.nodeId }), + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { ...message, streaming: false, updatedAt: input.completedAt }, + }); + } + for (const providerTurn of input.open.providerTurns ?? []) { + if (isTerminalProviderTurnStatus(providerTurn.status)) continue; + events.push({ + id: yield* input.allocateEventId(), + type: "provider-turn.updated", + threadId: input.run.threadId, + runId: input.run.id, + nodeId: providerTurn.nodeId, + providerInstanceId: input.run.providerInstanceId, + occurredAt: input.completedAt, + payload: { ...providerTurn, status: input.status, completedAt: input.completedAt }, + }); + } return events; }); } +/** + * The cascade's input for a run that never reached its provider, read from the + * projection instead of tracked from a live event stream. `threads` holds the + * run's own thread and the provider-native child threads it linked. + */ +export function openRunOwnedWorkFromProjection(input: { + readonly run: OrchestrationV2Run; + readonly threads: ReadonlyArray; + readonly linkedChildThreadIds: ReadonlySet; +}): OpenRunOwnedSubagentProjection { + const { run } = input; + const owns = (row: { readonly threadId: ThreadId; readonly runId: string | null }) => + row.threadId === run.threadId + ? row.runId === run.id + : input.linkedChildThreadIds.has(row.threadId); + const rows = ( + select: (thread: ProjectionStore.ProjectionRuntimeRecoveryState) => ReadonlyArray, + ) => input.threads.flatMap(select); + const runAttemptIds = new Set( + rows((thread) => thread.attempts) + .filter((attempt) => attempt.runId === run.id) + .map((attempt) => attempt.id), + ); + // App-owned tasks (delegate_task) run in their own threads and outlive it. + const appOwnedTaskIds = new Set( + rows((thread) => thread.subagents) + .filter((subagent) => subagent.origin === "app_owned") + .map((subagent) => subagent.id), + ); + const ownedTurnItems = rows((thread) => thread.turnItems).filter( + (item) => + owns(item) && + !isSettledTurnItemStatus(item.status) && + !(item.type === "subagent" && item.origin === "app_owned"), + ); + return { + subagents: new Map( + rows((thread) => thread.subagents) + .filter( + (subagent) => + subagent.runId === run.id && + subagent.origin === "provider_native" && + !isSettledSubagentStatus(subagent.status), + ) + .map((subagent) => [subagent.id, subagent]), + ), + turnItems: new Map( + ownedTurnItems.flatMap((item) => + item.type === "subagent" && item.runId === run.id ? [[item.subagentId, item] as const] : [], + ), + ), + childTurnItems: new Map( + ownedTurnItems + .filter((item) => !(item.type === "subagent" && item.runId === run.id)) + .map((item) => [item.id, item]), + ), + nodes: new Map( + rows((thread) => thread.nodes) + .filter( + (node) => + owns(node) && + // The never-started path writes the root with the run's own status. + node.id !== run.rootNodeId && + !appOwnedTaskIds.has(node.id) && + isOpenExecutionNodeStatus(node.status), + ) + .map((node) => [node.id, node]), + ), + linkedChildThreadIds: input.linkedChildThreadIds, + runtimeRequests: rows((thread) => thread.runtimeRequests), + streamingMessages: rows((thread) => thread.messages).filter((message) => message.streaming), + providerTurns: rows((thread) => thread.providerTurns).filter( + (turn) => turn.runAttemptId !== null && runAttemptIds.has(turn.runAttemptId), + ), + }; +} + export function finalProviderThreadStatus( disposition: ProviderTerminalEvent["threadDisposition"], ): OrchestrationV2ProviderThread["status"] { diff --git a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts index 0737be290afd..f1b3a0d64ba4 100644 --- a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts +++ b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts @@ -4,9 +4,12 @@ import { EventId, MessageId, type ModelSelection, + NodeId, + type OrchestrationV2DomainEvent, type OrchestrationV2ProviderCapabilities, type OrchestrationV2ProviderSession, type OrchestrationV2ProviderThread, + type OrchestrationV2ThreadProjection, ProjectId, ProviderDriverKind, ProviderInstanceId, @@ -14,7 +17,9 @@ import { ProviderThreadId, ProviderTurnId, type RunId, + RuntimeRequestId, ThreadId, + TurnItemId, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -86,6 +91,313 @@ interface RestartAdapterState { }>; readonly closedSessionCount: number; readonly failedReplacementOpen: boolean; + /** Every replacement open fails, not just the first. */ + readonly replacementOpenAlwaysFails?: boolean; + /** The first turn leaves an approval, a streaming reply, a subagent and a command open. */ + readonly firstTurnLeavesWorkOpen?: boolean; +} + +const providerNativeChildThreadId = (threadId: ThreadId) => + ThreadId.make(`${threadId}:native-subagent`); + +// Work a turn leaves open when it is superseded: an approval it is waiting on, +// a reply still streaming, a provider-native subagent whose own thread is still +// running a command, and a background command. +function openTurnWork( + providerSessionId: ProviderSessionId, + active: ActiveTurn, + now: DateTime.Utc, +): ReadonlyArray { + const { input, providerTurnId } = active; + const base = { + threadId: input.threadId, + runId: input.runId, + providerThreadId: input.providerThread.id, + providerTurnId, + nativeItemRef: null, + parentItemId: null, + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }; + const approvalNodeId = NodeId.make(`node:approval:${input.attemptId}`); + const subagentNodeId = NodeId.make(`node:subagent:${input.attemptId}`); + const childThreadId = providerNativeChildThreadId(input.threadId); + const grandchildThreadId = providerNativeChildThreadId(childThreadId); + const childApprovalNodeId = NodeId.make(`node:child-approval:${input.attemptId}`); + const childRequestId = RuntimeRequestId.make(`request:child-approval:${input.attemptId}`); + const requestId = RuntimeRequestId.make(`request:approval:${input.attemptId}`); + return [ + { + type: "node.updated", + driver, + node: { + id: approvalNodeId, + threadId: input.threadId, + runId: input.runId, + parentNodeId: input.rootNodeId, + rootNodeId: input.rootNodeId, + kind: "approval_request", + status: "waiting", + countsForRun: false, + providerThreadId: input.providerThread.id, + providerTurnId, + nativeItemRef: null, + runtimeRequestId: requestId, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + type: "runtime_request.updated", + driver, + threadId: input.threadId, + runtimeRequest: { + id: requestId, + nodeId: approvalNodeId, + providerTurnId, + nativeRequestRef: null, + kind: "command", + status: "pending", + responseCapability: { type: "live", providerSessionId }, + createdAt: now, + resolvedAt: null, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:approval:${input.attemptId}`), + nodeId: approvalNodeId, + ordinal: 1, + status: "waiting", + type: "approval_request", + requestId, + requestKind: "command", + }, + }, + { + type: "message.updated", + driver, + message: { + id: MessageId.make(`message:reply:${input.attemptId}`), + threadId: input.threadId, + runId: input.runId, + nodeId: input.rootNodeId, + role: "assistant", + createdBy: "agent", + creationSource: "provider", + text: "Working on it", + attachments: [], + streaming: true, + createdAt: now, + updatedAt: now, + }, + }, + { + type: "app_thread.created", + driver, + appThread: { + ...input.appThread, + createdBy: "agent", + creationSource: "provider", + id: childThreadId, + title: "Explore", + activeProviderThreadId: null, + lineage: { + parentThreadId: input.threadId, + relationshipToParent: "subagent", + rootThreadId: input.threadId, + }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:subagent:${input.attemptId}`), + nodeId: subagentNodeId, + ordinal: 2, + status: "running", + type: "subagent", + subagentId: subagentNodeId, + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId, + prompt: "Explore the repo", + result: null, + }, + }, + // The subagent is waiting on an approval in its own thread. Its rows carry + // no run id; the child thread has no runs. + { + type: "node.updated", + driver, + node: { + id: childApprovalNodeId, + threadId: childThreadId, + runId: null, + parentNodeId: null, + rootNodeId: childApprovalNodeId, + kind: "approval_request", + status: "waiting", + countsForRun: false, + providerThreadId: input.providerThread.id, + providerTurnId, + nativeItemRef: null, + runtimeRequestId: childRequestId, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + type: "runtime_request.updated", + driver, + threadId: childThreadId, + runtimeRequest: { + id: childRequestId, + nodeId: childApprovalNodeId, + providerTurnId, + nativeRequestRef: null, + kind: "command", + status: "pending", + responseCapability: { type: "live", providerSessionId }, + createdAt: now, + resolvedAt: null, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:child-approval:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: childApprovalNodeId, + ordinal: 3, + status: "waiting", + type: "approval_request", + requestId: childRequestId, + requestKind: "command", + }, + }, + { + type: "message.updated", + driver, + message: { + id: MessageId.make(`message:child-reply:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + role: "assistant", + createdBy: "agent", + creationSource: "provider", + text: "Looking", + attachments: [], + streaming: true, + createdAt: now, + updatedAt: now, + }, + }, + // The subagent's own subagent already reported done, but its thread is + // still running a command. + { + type: "app_thread.created", + driver, + appThread: { + ...input.appThread, + createdBy: "agent", + creationSource: "provider", + id: grandchildThreadId, + title: "Search", + activeProviderThreadId: null, + lineage: { + parentThreadId: childThreadId, + relationshipToParent: "subagent", + rootThreadId: input.threadId, + }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:nested-subagent:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 2, + status: "completed", + completedAt: now, + type: "subagent", + subagentId: NodeId.make(`node:nested-subagent:${input.attemptId}`), + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId: grandchildThreadId, + prompt: "Search for TODOs", + result: "done", + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:grandchild-command:${input.attemptId}`), + threadId: grandchildThreadId, + runId: null, + nodeId: null, + ordinal: 1, + status: "running", + type: "command_execution", + input: "rg FIXME", + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:child-command:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 1, + status: "running", + type: "command_execution", + input: "rg TODO", + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:command:${input.attemptId}`), + nodeId: input.rootNodeId, + ordinal: 3, + status: "running", + type: "command_execution", + input: "npm run dev", + }, + }, + ]; } function makeRestartAdapter( @@ -108,7 +420,7 @@ function makeRestartAdapter( const failThisOpen = yield* Ref.modify(state, (current) => { const shouldFail = sessionInput.modelSelection.model === replacementSelection.model && - !current.failedReplacementOpen; + (current.replacementOpenAlwaysFails === true || !current.failedReplacementOpen); return [ shouldFail, { @@ -258,6 +570,15 @@ function makeRestartAdapter( completedAt: null, }, }); + if ((yield* Ref.get(state)).firstTurnLeavesWorkOpen === true) { + for (const event of openTurnWork( + sessionInput.providerSessionId, + active, + occurredAt, + )) { + yield* Queue.offer(events, event); + } + } return; } yield* publishTerminal(active, "completed"); @@ -537,6 +858,253 @@ it.live("restarts selection as a new attempt and retries after old-session clean ), ); +// A restart only replaces the root turn: the request, reply and command the +// superseded attempt left open become the restarted run's to settle. When the +// replacement never starts, failing the run must end them too, or the thread +// waits on work nothing is running. +it.live("settles the work a restarted run inherited when its replacement never opens", () => + Effect.scoped( + Effect.gen(function* () { + const name = "selection-restart-failed-start"; + const cwd = yield* checkpointWorkspace(name); + const threadId = ThreadId.make(`thread:${name}`); + const state = yield* Ref.make({ + activeTurn: null, + opened: [], + started: [], + closedSessionCount: 0, + failedReplacementOpen: false, + replacementOpenAlwaysFails: true, + firstTurnLeavesWorkOpen: true, + }); + const layerRegistry = ProviderAdapterRegistry.layerSingle(makeRestartAdapter(state)); + + const { + projection, + nativeChildItems, + nativeChildStreaming, + nativeChildRequests, + nativeGrandchildItems, + delegatedItems, + } = yield* Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const eventSink = yield* EventSink.EventSinkV2; + const awaitDomainEvent = (matches: (event: OrchestrationV2DomainEvent) => boolean) => + orchestrator.streamDomainEvents.pipe( + Stream.filter(matches), + Stream.take(1), + Stream.runDrain, + Effect.forkScoped({ startImmediately: true }), + ); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:create`), + threadId, + projectId: ProjectId.make(`project:${name}`), + title: name, + modelSelection: initialSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: cwd, + }); + // The parent's command item is the last open work the first turn reports. + const workOpen = yield* awaitDomainEvent( + (event) => + event.type === "turn-item.updated" && + event.payload.type === "command_execution" && + event.payload.threadId === threadId, + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:first`), + threadId, + messageId: MessageId.make(`${name}:first`), + text: "first", + attachments: [], + modelSelection: initialSelection, + dispatchMode: { type: "start_immediately" }, + }); + yield* Fiber.join(workOpen); + const started = yield* orchestrator.getThreadProjection(threadId); + const runId = started.runs[0]?.id; + if (runId === undefined) return yield* Effect.die("the first run is missing"); + // The run also delegated a task (delegate_task). Its thread runs on its + // own, so the failed restart must leave it running. + const now = yield* DateTime.now; + const delegatedThreadId = ThreadId.make(`${threadId}:delegated`); + const delegatedTaskId = NodeId.make(`node:delegated:${runId}`); + yield* eventSink.write({ + events: [ + { + id: EventId.make(`event:${name}:delegated-thread`), + type: "thread.created", + threadId: delegatedThreadId, + occurredAt: now, + payload: { + ...started.thread, + createdBy: "agent", + creationSource: "mcp", + id: delegatedThreadId, + activeProviderThreadId: null, + lineage: { + parentThreadId: threadId, + relationshipToParent: "subagent", + rootThreadId: threadId, + }, + }, + }, + { + id: EventId.make(`event:${name}:delegated-task`), + type: "subagent.updated", + threadId, + runId, + nodeId: delegatedTaskId, + driver, + providerInstanceId, + occurredAt: now, + payload: { + id: delegatedTaskId, + threadId, + runId, + parentNodeId: started.runs[0]!.rootNodeId!, + origin: "app_owned", + createdBy: "agent", + driver, + providerInstanceId, + providerThreadId: null, + childThreadId: delegatedThreadId, + nativeTaskRef: null, + prompt: "Write the tests", + title: null, + model: initialSelection.model, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + { + id: EventId.make(`event:${name}:delegated-command`), + type: "turn-item.updated", + threadId: delegatedThreadId, + occurredAt: now, + payload: { + id: TurnItemId.make(`turn-item:${name}:delegated-command`), + threadId: delegatedThreadId, + runId: null, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + type: "command_execution", + input: "vp test run", + }, + }, + ], + }); + + const runFailed = yield* awaitDomainEvent( + (event) => event.type === "run.updated" && event.payload.status === "failed", + ); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make(`${name}:second`), + threadId, + messageId: MessageId.make(`${name}:second`), + text: "second", + attachments: [], + modelSelection: replacementSelection, + dispatchMode: { type: "restart_active", targetRunId: runId }, + }); + yield* Fiber.join(runFailed); + const openItems = (projection: OrchestrationV2ThreadProjection) => + projection.turnItems.flatMap((item) => + item.type === "approval_request" || + item.type === "command_execution" || + item.type === "subagent" + ? [[item.type, item.status]] + : [], + ); + return { + projection: yield* orchestrator.getThreadProjection(threadId), + nativeChildItems: openItems( + yield* orchestrator.getThreadProjection(providerNativeChildThreadId(threadId)), + ), + nativeChildStreaming: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).messages.some((message) => message.streaming), + nativeChildRequests: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).runtimeRequests.map((request) => request.status), + nativeGrandchildItems: openItems( + yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(providerNativeChildThreadId(threadId)), + ), + ), + delegatedItems: openItems(yield* orchestrator.getThreadProjection(delegatedThreadId)), + }; + }).pipe(Effect.provide(ProviderReplayHarness.layerWithRegistry({ name }, layerRegistry))); + + assert.deepEqual( + projection.runs.map((run) => run.status), + ["failed"], + ); + assert.deepEqual( + projection.runtimeRequests.map((request) => request.status), + ["cancelled"], + ); + assert.isFalse(projection.messages.some((message) => message.streaming)); + assert.deepEqual( + projection.turnItems.flatMap((item) => + item.type === "approval_request" || + item.type === "command_execution" || + item.type === "subagent" + ? [[item.type, item.status]] + : [], + ), + [ + ["approval_request", "cancelled"], + ["subagent", "cancelled"], + ["command_execution", "cancelled"], + ], + ); + // The provider-native subagent's own thread ends with the run, and so does + // the thread of a nested subagent whose row already settled. + assert.deepEqual(nativeChildItems, [ + ["approval_request", "cancelled"], + ["subagent", "completed"], + ["command_execution", "cancelled"], + ]); + assert.isFalse(nativeChildStreaming); + assert.deepEqual(nativeChildRequests, ["cancelled"]); + assert.deepEqual(nativeGrandchildItems, [["command_execution", "cancelled"]]); + // The delegated task and its thread keep running. + assert.deepEqual( + projection.subagents.flatMap((subagent) => + subagent.origin === "app_owned" ? [subagent.status] : [], + ), + ["running"], + ); + assert.deepEqual(delegatedItems, [["command_execution", "running"]]); + }), + ), +); + it.live.each(["stopped", "error"] as const)( "restarts the live session on a model change when a newer %s session record exists", (deadStatus) => From 93082c0a79e2c219358aeff27c3386e6214515a6 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:34:50 -0700 Subject: [PATCH 02/11] fix(server): an approval accepted while a failed start is written stays accepted The start-failure write cancels the requests the run inherited, guarded with guardPendingUserInputCancellations. That guard only rechecks user_input requests, so an approval the user accepted between the projection read and the commit was overwritten to cancelled (request, node and item), and the response executor then refused to deliver it. writeIfRunCurrent gains guardPendingRequestCancellations, which rechecks every request kind and also keeps approval_request items. It is opt-in: the provider event ingestor keeps the user-input guard, because a provider's own approval cancellation is authoritative. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/orchestration-v2/EventSink.ts | 47 ++- .../ProviderTurnStartService.test.ts | 372 ++++++++++++++++++ .../ProviderTurnStartService.ts | 6 +- 3 files changed, 403 insertions(+), 22 deletions(-) diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index 039d1872da6f..1cc20e281afc 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -84,6 +84,12 @@ export interface EventSinkV2Shape { }) => Effect.Effect, EventSinkV2Error>; readonly writeIfRunCurrent: (input: { readonly guardPendingUserInputCancellations?: boolean; + /** + * Guards approvals and every other request kind too. For a writer that + * cancels requests an earlier projection read showed pending; a + * provider's own approval cancellation stays authoritative. + */ + readonly guardPendingRequestCancellations?: boolean; readonly commandId?: CommandId; readonly threadId: ThreadId; readonly runId: RunId; @@ -263,14 +269,18 @@ const layerBase: Layer.Layer< // A user can answer after terminal normalization reads the pending request. // Recheck inside the write transaction so stale cleanup cannot erase answers. - const guardUserInputCancellations = (events: ReadonlyArray) => + // `allKinds` extends this from questions to approvals and every other kind. + const guardRequestCancellations = ( + events: ReadonlyArray, + allKinds: boolean, + ) => Effect.gen(function* () { const staleRequests = new Set(); const staleNodes = new Set(); for (const event of events) { if ( event.type !== "runtime-request.updated" || - event.payload.kind !== "user_input" || + (!allKinds && event.payload.kind !== "user_input") || event.payload.status !== "cancelled" ) continue; @@ -280,7 +290,7 @@ const layerBase: Layer.Layer< ); if ( current?.status !== "pending" || - current.kind !== "user_input" || + current.kind !== event.payload.kind || current.providerTurnId !== event.payload.providerTurnId || current.responseCapability.type === "message" ) { @@ -296,7 +306,8 @@ const layerBase: Layer.Layer< return event.payload.status !== "cancelled" || !staleNodes.has(event.payload.id); case "turn-item.updated": return ( - event.payload.type !== "user_input_request" || + (event.payload.type !== "user_input_request" && + event.payload.type !== "approval_request") || event.payload.status !== "cancelled" || !staleRequests.has(event.payload.requestId) ); @@ -305,6 +316,16 @@ const layerBase: Layer.Layer< } }); }); + const guardCancellations = (input: { + readonly guardPendingUserInputCancellations?: boolean; + readonly guardPendingRequestCancellations?: boolean; + readonly events: ReadonlyArray; + }) => + input.guardPendingRequestCancellations === true + ? guardRequestCancellations(input.events, true) + : input.guardPendingUserInputCancellations === true + ? guardRequestCancellations(input.events, false) + : Effect.succeed(input.events); const normalizeEvents = (events: ReadonlyArray) => { const runOrdinals = new Map( @@ -370,11 +391,7 @@ const layerBase: Layer.Layer< return yield* commitThenPublish( Effect.gen(function* () { - const normalized = yield* normalizeEvents( - input.guardPendingUserInputCancellations === true - ? yield* guardUserInputCancellations(input.events) - : input.events, - ); + const normalized = yield* normalizeEvents(yield* guardCancellations(input)); const committed = yield* eventStore.append({ ...(input.commandId === undefined ? {} : { commandId: input.commandId }), events: normalized, @@ -428,11 +445,7 @@ const layerBase: Layer.Layer< }; } - const normalized = yield* normalizeEvents( - input.guardPendingUserInputCancellations === true - ? yield* guardUserInputCancellations(input.events) - : input.events, - ); + const normalized = yield* normalizeEvents(yield* guardCancellations(input)); const storedEvents = yield* eventStore.append({ ...(input.commandId === undefined ? {} : { commandId: input.commandId }), events: normalized, @@ -493,11 +506,7 @@ const layerBase: Layer.Layer< }; } - const normalized = yield* normalizeEvents( - input.guardPendingUserInputCancellations === true - ? yield* guardUserInputCancellations(input.events) - : input.events, - ); + const normalized = yield* normalizeEvents(yield* guardCancellations(input)); const storedEvents = yield* eventStore.append({ ...(input.commandId === undefined ? {} : { commandId: input.commandId }), events: normalized, diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index b4c064b1bd0c..e14ad41407aa 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -2,18 +2,27 @@ import { expect, it, vi } from "vite-plus/test"; import { it as effectIt } from "@effect/vitest"; import { CheckpointScopeId, + EventId, MessageId, NodeId, ProviderSessionId, ProviderThreadId, + ProviderTurnId, ProviderDriverKind, ProviderInstanceId, ProviderSetupError, RunAttemptId, RunId, + RuntimeRequestId, ThreadId, + TurnItemId, ProjectId, + type OrchestrationV2AppThread, + type OrchestrationV2ExecutionNode, + type OrchestrationV2Run, + type OrchestrationV2RuntimeRequest, type OrchestrationV2ThreadProjection, + type OrchestrationV2TurnItem, OrchestrationV2DomainEvent, } from "@t3tools/contracts"; import * as Cause from "effect/Cause"; @@ -26,9 +35,11 @@ import * as Schema from "effect/Schema"; import * as GitWorkflow from "../git/GitWorkflowService.ts"; import * as ProjectService from "../project/ProjectService.ts"; +import * as SqlitePersistence from "../persistence/Sqlite.ts"; import * as ProviderAuthService from "../provider/ProviderAuthService.ts"; import * as ContextHandoffService from "./ContextHandoffService.ts"; import * as EventSink from "./EventSink.ts"; +import * as EventStore from "./EventStore.ts"; import * as IdAllocator from "./IdAllocator.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; @@ -856,3 +867,364 @@ for (const previousMessages of [[], ["/compact", " /COMPACT "]]) { }), ); } + +// Runs the start path against the real event sink and projection store, so the +// failure write goes through the same commit guard production uses. +function makePersistedStartFailureHarness() { + const now = DateTime.makeUnsafe("2026-10-06T12:00:00Z"); + const driver = ProviderDriverKind.make("codex"); + const instanceId = ProviderInstanceId.make("codex-start-failure"); + const threadId = ThreadId.make("thread-start-failure"); + const runId = RunId.make("run-start-failure"); + const attemptId = RunAttemptId.make("attempt-start-failure"); + const rootNodeId = NodeId.make("root-start-failure"); + const providerThreadId = ProviderThreadId.make("provider-thread-start-failure"); + const providerSessionId = ProviderSessionId.make("provider-session-start-failure"); + const providerTurnId = ProviderTurnId.make("provider-turn-start-failure"); + const checkpointScopeId = CheckpointScopeId.make("scope-start-failure"); + const messageId = MessageId.make("message-start-failure"); + const approvalNodeId = NodeId.make("node-approval-start-failure"); + const approvalRequestId = RuntimeRequestId.make("request-approval-start-failure"); + const approvalItemId = TurnItemId.make("item-approval-start-failure"); + const thread: OrchestrationV2AppThread = { + createdBy: "user", + creationSource: "web", + id: threadId, + projectId: ProjectId.make("project-start-failure"), + title: "Start failure", + providerInstanceId: instanceId, + modelSelection: { instanceId, model: "gpt-5.4" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + activeProviderThreadId: providerThreadId, + lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: threadId }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + archivedAt: null, + settledOverride: null, + settledAt: null, + lastVisitedAt: null, + deletedAt: null, + }; + const run: OrchestrationV2Run = { + id: runId, + threadId, + ordinal: 1, + providerInstanceId: instanceId, + modelSelection: thread.modelSelection, + providerThreadId, + userMessageId: messageId, + rootNodeId, + activeAttemptId: attemptId, + status: "starting", + requestedAt: now, + startedAt: null, + completedAt: null, + checkpointId: null, + contextHandoffId: null, + }; + const approvalRequest: OrchestrationV2RuntimeRequest = { + id: approvalRequestId, + nodeId: approvalNodeId, + providerTurnId, + nativeRequestRef: null, + kind: "command", + status: "pending", + responseCapability: { type: "live", providerSessionId }, + createdAt: now, + resolvedAt: null, + }; + const approvalNode: OrchestrationV2ExecutionNode = { + id: approvalNodeId, + threadId, + runId, + parentNodeId: rootNodeId, + rootNodeId, + kind: "approval_request", + status: "waiting", + countsForRun: false, + providerThreadId, + providerTurnId, + nativeItemRef: null, + runtimeRequestId: approvalRequestId, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }; + const approvalItem: OrchestrationV2TurnItem = { + id: approvalItemId, + threadId, + runId, + nodeId: approvalNodeId, + providerThreadId, + providerTurnId, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "waiting", + title: "Run npm test?", + startedAt: now, + completedAt: null, + updatedAt: now, + type: "approval_request", + requestId: approvalRequestId, + requestKind: "command", + }; + // The run as a restarted attempt inherits it: the earlier attempt left an + // approval waiting under the run. + const seedPayloads: ReadonlyArray< + Pick & { readonly runId?: RunId } + > = [ + { type: "thread.created", payload: thread }, + { + type: "message.updated", + payload: { + id: messageId, + threadId, + runId, + nodeId: rootNodeId, + role: "user", + createdBy: "user", + creationSource: "web", + text: "Continue", + attachments: [], + streaming: false, + createdAt: now, + updatedAt: now, + }, + }, + { type: "run.created", payload: run }, + { + type: "run-attempt.created", + payload: { + id: attemptId, + runId, + attemptOrdinal: 2, + rootNodeId, + providerInstanceId: instanceId, + providerThreadId, + providerTurnId: null, + reason: "steering_restart", + status: "pending", + startedAt: null, + completedAt: null, + }, + }, + { + type: "checkpoint-scope.created", + payload: { + id: checkpointScopeId, + threadId, + runId, + nodeId: rootNodeId, + parentScopeId: null, + providerThreadId, + kind: "root_run", + ordinalWithinParent: 0, + advancesAppRunCount: true, + cwd: "/tmp/start-failure", + createdAt: now, + }, + }, + { + type: "node.updated", + payload: { + id: rootNodeId, + threadId, + runId, + parentNodeId: null, + rootNodeId, + kind: "root_turn", + status: "running", + countsForRun: true, + providerThreadId, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId, + startedAt: now, + completedAt: null, + }, + }, + { + type: "provider-thread.updated", + payload: { + id: providerThreadId, + driver, + providerInstanceId: instanceId, + providerSessionId, + appThreadId: threadId, + ownerNodeId: null, + nativeThreadRef: null, + nativeConversationHeadRef: null, + status: "not_loaded", + firstRunOrdinal: 1, + lastRunOrdinal: 1, + handoffIds: [], + forkedFrom: null, + createdAt: now, + updatedAt: now, + }, + }, + { type: "node.updated", payload: approvalNode }, + { type: "runtime-request.updated", payload: approvalRequest }, + { type: "turn-item.updated", payload: approvalItem }, + ]; + const layerDatabase = SqlitePersistence.layerMemory; + const layerStores = Layer.mergeAll(EventStore.layer, ProjectionStore.layer).pipe( + Layer.provideMerge(layerDatabase), + ); + const layerSink = EventSink.layer.pipe(Layer.provide(layerStores)); + const layerPersistence = Layer.mergeAll(layerStores, layerSink, IdAllocator.layer); + // `beforeFailureWrite` runs once the start has read its projection and + // given up on the provider, right before the failure is committed. + const run_ = (options: { + readonly beforeFailureWrite?: Effect.Effect< + void, + EventSink.EventSinkV2Error, + EventSink.EventSinkV2 + >; + readonly projectionStore?: ( + store: ProjectionStore.ProjectionStoreV2Shape, + ) => ProjectionStore.ProjectionStoreV2Shape; + }) => + Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const store = yield* ProjectionStore.ProjectionStoreV2; + const idAllocator = yield* IdAllocator.IdAllocatorV2; + yield* eventSink.write({ + events: yield* Effect.forEach(seedPayloads, (event) => + Effect.gen(function* () { + return { + ...event, + id: yield* idAllocator.allocate.event({ threadId }), + threadId, + occurredAt: now, + } as OrchestrationV2DomainEvent; + }), + ), + }); + const gatedSink = EventSink.EventSinkV2.of({ + ...eventSink, + writeIfRunCurrent: (input) => + (options.beforeFailureWrite ?? Effect.void).pipe( + Effect.provideService(EventSink.EventSinkV2, eventSink), + Effect.orDie, + Effect.andThen(eventSink.writeIfRunCurrent(input)), + ), + }); + const service = yield* ProviderTurnStart.ProviderTurnStartServiceV2.pipe( + Effect.provide( + ProviderTurnStart.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.succeed(EventSink.EventSinkV2, gatedSink), + Layer.succeed( + ProjectionStore.ProjectionStoreV2, + options.projectionStore?.(store) ?? store, + ), + IdAllocator.layer, + Layer.mock(ContextHandoffService.ContextHandoffServiceV2)({}), + FileSystem.layerNoop({}), + Layer.mock(GitWorkflow.GitWorkflowService)({}), + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(ProviderSessionManager.ProviderSessionManagerV2)({ + open: () => + Effect.fail( + new ProviderSessionManager.ProviderSessionOpenError({ + instanceId, + providerSessionId, + cause: "replacement session failed to open", + }), + ), + }), + Layer.mock(ProviderAuthService.ProviderAuthService)({ + tryHandlePromptCommand: () => Effect.succeed(false), + }), + Layer.mock(RunExecutionService.RunExecutionServiceV2)({ + startRootRun: () => Effect.die("a failed open must not start the run"), + }), + Layer.mock(RuntimePolicy.RuntimePolicyV2)({ + resolve: () => Effect.succeed({} as never), + }), + ), + ), + ), + ), + ); + yield* service.start({ threadId, runId }); + return yield* store.getThreadProjection(threadId); + }).pipe(Effect.provide(Layer.fresh(layerPersistence))); + return { + now, + threadId, + runId, + approvalRequest, + approvalNode, + approvalItem, + run: run_, + }; +} + +effectIt.effect("keeps an approval accepted while a failed start is written", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const projection = yield* harness.run({ + // The user approves after the failure read the approval as pending. + beforeFailureWrite: Effect.gen(function* () { + const eventSink = yield* EventSink.EventSinkV2; + const { now, threadId } = harness; + yield* eventSink.write({ + events: [ + { + id: EventId.make("event-approval-accepted"), + type: "runtime-request.updated", + threadId, + nodeId: harness.approvalNode.id, + occurredAt: now, + payload: { + ...harness.approvalRequest, + status: "resolved", + decision: "accept", + resolvedAt: now, + }, + }, + { + id: EventId.make("event-approval-node-completed"), + type: "node.updated", + threadId, + nodeId: harness.approvalNode.id, + occurredAt: now, + payload: { ...harness.approvalNode, status: "completed", completedAt: now }, + }, + { + id: EventId.make("event-approval-item-completed"), + type: "turn-item.updated", + threadId, + nodeId: harness.approvalNode.id, + occurredAt: now, + payload: { + ...harness.approvalItem, + status: "completed", + completedAt: now, + updatedAt: now, + }, + }, + ], + }); + }), + }); + + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.runtimeRequests).toMatchObject([{ status: "resolved", decision: "accept" }]); + expect(projection.nodes.find((node) => node.id === harness.approvalNode.id)?.status).toBe( + "completed", + ); + expect(projection.turnItems.find((item) => item.id === harness.approvalItem.id)?.status).toBe( + "completed", + ); + }), +); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index d1b005f18826..7d787d069422 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -459,9 +459,9 @@ export const layer: Layer.Layer< runId, activeAttemptId: attempt.id, expectedStatus: "starting", - // A user answer that lands while the failure is written wins over - // cancelling the request it answered. - guardPendingUserInputCancellations: true, + // An answer or approval that lands while the failure is written wins + // over cancelling the request it resolved. + guardPendingRequestCancellations: true, events: status === "failed" ? [...(yield* inheritedWorkSettlement({ run, now })), ...events] From 0c4a21462047f1a2a5b80e6a899183d990219ec1 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:34:58 -0700 Subject: [PATCH 03/11] fix(server): a failed start still fails its run when inherited work can't be read Only ProjectionStoreThreadNotFoundError was handled while collecting the work a failed run inherited. Any other store or decode error propagated out of settleStartFailure, so the run, its attempt and root node were never written failed and the run stayed starting, which is the symptom the failure path exists to prevent. A failure to collect inherited work is now logged with its cause, and the run, attempt and root failure is still written without the inherited events. Interruption still propagates. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../ProviderTurnStartService.test.ts | 28 +++++++++++++++++++ .../ProviderTurnStartService.ts | 17 ++++++++++- 2 files changed, 44 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index e14ad41407aa..f2360629799c 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -1228,3 +1228,31 @@ effectIt.effect("keeps an approval accepted while a failed start is written", () ); }), ); + +effectIt.effect("fails the run when the work it inherited cannot be read", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const projection = yield* harness.run({ + // The start's own reads succeed; only the inherited-work read fails. + projectionStore: (store) => ({ + ...store, + getThreadRecords: (threadId, fields, filter) => + fields.some((field) => field === "subagents") + ? Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ + threadId, + cause: "database unavailable", + }), + ) + : store.getThreadRecords(threadId, fields, filter), + }), + }); + + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.attempts.map((attempt) => attempt.status)).toEqual(["failed"]); + expect(projection.turnItems).toMatchObject([ + { type: "approval_request", status: "waiting" }, + { type: "error", title: "Provider session failed to open" }, + ]); + }), +); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 7d787d069422..24c7c5c0c259 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -464,7 +464,22 @@ export const layer: Layer.Layer< guardPendingRequestCancellations: true, events: status === "failed" - ? [...(yield* inheritedWorkSettlement({ run, now })), ...events] + ? [ + // Failing the run matters more than settling what it + // inherited: a read that fails here must not leave it + // `starting`. + ...(yield* inheritedWorkSettlement({ run, now }).pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning( + "provider turn start could not settle the work a failed run inherited", + { threadId: projection.thread.id, runId, cause: Cause.pretty(cause) }, + ).pipe(Effect.as([])), + ), + )), + ...events, + ] : events, }); }, From fb87524701ab12c530ca629e560be6e07a7bb7be Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:43:26 -0700 Subject: [PATCH 04/11] fix(server): a failed start settles subagents its native subagents launched A provider-native subagent can launch its own subagent; that row lives on the child thread with runId null (Codex does this). openRunOwnedWorkFromProjection only took subagent rows carrying the root run's id, so the nested subagent's item was cancelled while its row, which clients render, stayed running. Subagent rows are now taken by the same ownership as every other row: the run's own thread by run id, a linked child thread by thread. App-owned tasks stay excluded. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../orchestration-v2/RunExecutionService.ts | 4 +- .../SelectionRestart.integration.test.ts | 56 +++++++++++++++++++ 2 files changed, 59 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index 33e6d92b09f8..36f4c3ea8f7d 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -420,11 +420,13 @@ export function openRunOwnedWorkFromProjection(input: { !(item.type === "subagent" && item.origin === "app_owned"), ); return { + // A linked child's own subagents carry no run id, so they are taken by + // the thread they run on. subagents: new Map( rows((thread) => thread.subagents) .filter( (subagent) => - subagent.runId === run.id && + owns(subagent) && subagent.origin === "provider_native" && !isSettledSubagentStatus(subagent.status), ) diff --git a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts index f1b3a0d64ba4..034f8d89668a 100644 --- a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts +++ b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts @@ -127,6 +127,7 @@ function openTurnWork( const grandchildThreadId = providerNativeChildThreadId(childThreadId); const childApprovalNodeId = NodeId.make(`node:child-approval:${input.attemptId}`); const childRequestId = RuntimeRequestId.make(`request:child-approval:${input.attemptId}`); + const nestedRunningSubagentId = NodeId.make(`node:nested-running-subagent:${input.attemptId}`); const requestId = RuntimeRequestId.make(`request:approval:${input.attemptId}`); return [ { @@ -354,6 +355,54 @@ function openTurnWork( result: "done", }, }, + // Another subagent it launched is still running. Like every row on the + // subagent's thread, it carries no run id. + { + type: "subagent.updated", + driver, + subagent: { + id: nestedRunningSubagentId, + threadId: childThreadId, + runId: null, + parentNodeId: subagentNodeId, + origin: "provider_native", + createdBy: "agent", + driver, + providerInstanceId: input.modelSelection.instanceId, + providerThreadId: null, + childThreadId: null, + nativeTaskRef: null, + prompt: "Read the docs", + title: null, + model: null, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:nested-running-subagent:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 4, + status: "running", + type: "subagent", + subagentId: nestedRunningSubagentId, + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId: null, + prompt: "Read the docs", + result: null, + }, + }, { type: "turn_item.updated", driver, @@ -884,6 +933,7 @@ it.live("settles the work a restarted run inherited when its replacement never o nativeChildItems, nativeChildStreaming, nativeChildRequests, + nativeChildSubagents, nativeGrandchildItems, delegatedItems, } = yield* Effect.gen(function* () { @@ -1051,6 +1101,9 @@ it.live("settles the work a restarted run inherited when its replacement never o nativeChildRequests: (yield* orchestrator.getThreadProjection( providerNativeChildThreadId(threadId), )).runtimeRequests.map((request) => request.status), + nativeChildSubagents: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).subagents.map((subagent) => subagent.status), nativeGrandchildItems: openItems( yield* orchestrator.getThreadProjection( providerNativeChildThreadId(providerNativeChildThreadId(threadId)), @@ -1088,10 +1141,13 @@ it.live("settles the work a restarted run inherited when its replacement never o assert.deepEqual(nativeChildItems, [ ["approval_request", "cancelled"], ["subagent", "completed"], + ["subagent", "cancelled"], ["command_execution", "cancelled"], ]); assert.isFalse(nativeChildStreaming); assert.deepEqual(nativeChildRequests, ["cancelled"]); + // A subagent the provider-native subagent launched ends with it. + assert.deepEqual(nativeChildSubagents, ["cancelled"]); assert.deepEqual(nativeGrandchildItems, [["command_execution", "cancelled"]]); // The delegated task and its thread keep running. assert.deepEqual( From c4ba9d6b01d581281dd91c2008d250ec5d10644b Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:46:10 -0700 Subject: [PATCH 05/11] fix(server): a failed start settles its native subagents' provider turns A provider-native child thread's provider turns have runAttemptId null. The recovery read omits them, the supplementary child read did not load them, and openRunOwnedWorkFromProjection only took turns of the run's own attempts, so a reloaded child still showed a running provider turn under a cancelled root. The child read now loads provider turns, the builder takes a linked child thread's open turns by thread, and the cascade writes each turn's event on the thread it runs on instead of always on the parent. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../ProviderTurnStartService.ts | 7 ++--- .../orchestration-v2/RunExecutionService.ts | 26 ++++++++++++++----- .../SelectionRestart.integration.test.ts | 24 +++++++++++++++++ 3 files changed, 48 insertions(+), 9 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 24c7c5c0c259..b3156977b46b 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -223,11 +223,11 @@ export const layer: Layer.Layer< if (Option.isNone(child) || !isProviderNativeSubagentThread(child.value.thread)) continue; linkedChildThreadIds.add(threadId); // A child thread has no runs, so the recovery read leaves out its - // streaming replies and its pending requests' nodes and items; read - // them by thread instead. + // streaming replies, provider turns and its pending requests' nodes + // and items; read them by thread instead. const records = yield* projectionStore.getThreadRecords( threadId, - ["messages", "nodes", "turnItems"], + ["messages", "nodes", "turnItems", "providerTurns"], { messageRoles: ["assistant"], turnItemTypes: ["approval_request", "user_input_request"], @@ -243,6 +243,7 @@ export const layer: Layer.Layer< const knownItemIds = new Set(child.value.turnItems.map((item) => item.id)); threads.push({ ...child.value, + providerTurns: records.providerTurns, nodes: [ ...child.value.nodes, ...records.nodes.filter( diff --git a/apps/server/src/orchestration-v2/RunExecutionService.ts b/apps/server/src/orchestration-v2/RunExecutionService.ts index 36f4c3ea8f7d..c611db246810 100644 --- a/apps/server/src/orchestration-v2/RunExecutionService.ts +++ b/apps/server/src/orchestration-v2/RunExecutionService.ts @@ -162,8 +162,14 @@ export type OpenRunOwnedSubagentProjection = { /** Pending requests; each settles with the node it waits on. */ readonly runtimeRequests?: ReadonlyArray; readonly streamingMessages?: ReadonlyArray; - /** Open provider turns of the run's own attempts. */ - readonly providerTurns?: ReadonlyArray; + /** + * Open provider turns of the run's own attempts and of linked child threads, + * with the thread each one runs on (the turn row does not carry it). + */ + readonly providerTurns?: ReadonlyArray<{ + readonly threadId: ThreadId; + readonly providerTurn: OrchestrationV2ProviderTurn; + }>; }; type RunOwnedSubagentTerminalStatus = Extract< @@ -367,12 +373,12 @@ export function cascadeTerminalizeRunOwnedSubagents(input: { payload: { ...message, streaming: false, updatedAt: input.completedAt }, }); } - for (const providerTurn of input.open.providerTurns ?? []) { + for (const { threadId, providerTurn } of input.open.providerTurns ?? []) { if (isTerminalProviderTurnStatus(providerTurn.status)) continue; events.push({ id: yield* input.allocateEventId(), type: "provider-turn.updated", - threadId: input.run.threadId, + threadId, runId: input.run.id, nodeId: providerTurn.nodeId, providerInstanceId: input.run.providerInstanceId, @@ -457,8 +463,16 @@ export function openRunOwnedWorkFromProjection(input: { linkedChildThreadIds: input.linkedChildThreadIds, runtimeRequests: rows((thread) => thread.runtimeRequests), streamingMessages: rows((thread) => thread.messages).filter((message) => message.streaming), - providerTurns: rows((thread) => thread.providerTurns).filter( - (turn) => turn.runAttemptId !== null && runAttemptIds.has(turn.runAttemptId), + // A linked child's provider turns belong to no run attempt. + providerTurns: input.threads.flatMap((thread) => + thread.providerTurns.flatMap((providerTurn) => + (thread.thread.id === run.threadId + ? providerTurn.runAttemptId !== null && runAttemptIds.has(providerTurn.runAttemptId) + : input.linkedChildThreadIds.has(thread.thread.id)) && + !isTerminalProviderTurnStatus(providerTurn.status) + ? [{ threadId: thread.thread.id, providerTurn }] + : [], + ), ), }; } diff --git a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts index 034f8d89668a..fa1a6825fb98 100644 --- a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts +++ b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts @@ -238,6 +238,24 @@ function openTurnWork( result: null, }, }, + // The subagent's own provider turn is running. It belongs to no run + // attempt. + { + type: "provider_turn.updated", + driver, + threadId: childThreadId, + providerTurn: { + id: ProviderTurnId.make(`provider-turn:child:${input.attemptId}`), + providerThreadId: ProviderThreadId.make(`provider-thread:${childThreadId}`), + nodeId: subagentNodeId, + runAttemptId: null, + nativeTurnRef: null, + ordinal: 1, + status: "running", + startedAt: now, + completedAt: null, + }, + }, // The subagent is waiting on an approval in its own thread. Its rows carry // no run id; the child thread has no runs. { @@ -934,6 +952,7 @@ it.live("settles the work a restarted run inherited when its replacement never o nativeChildStreaming, nativeChildRequests, nativeChildSubagents, + nativeChildProviderTurns, nativeGrandchildItems, delegatedItems, } = yield* Effect.gen(function* () { @@ -1104,6 +1123,9 @@ it.live("settles the work a restarted run inherited when its replacement never o nativeChildSubagents: (yield* orchestrator.getThreadProjection( providerNativeChildThreadId(threadId), )).subagents.map((subagent) => subagent.status), + nativeChildProviderTurns: (yield* orchestrator.getThreadProjection( + providerNativeChildThreadId(threadId), + )).providerTurns.map((turn) => turn.status), nativeGrandchildItems: openItems( yield* orchestrator.getThreadProjection( providerNativeChildThreadId(providerNativeChildThreadId(threadId)), @@ -1148,6 +1170,8 @@ it.live("settles the work a restarted run inherited when its replacement never o assert.deepEqual(nativeChildRequests, ["cancelled"]); // A subagent the provider-native subagent launched ends with it. assert.deepEqual(nativeChildSubagents, ["cancelled"]); + // So does the subagent's own provider turn, which no run attempt owns. + assert.deepEqual(nativeChildProviderTurns, ["cancelled"]); assert.deepEqual(nativeGrandchildItems, [["command_execution", "cancelled"]]); // The delegated task and its thread keep running. assert.deepEqual( From 377294a7919e6fb157207372f91f6c669e694321 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 15:52:52 -0700 Subject: [PATCH 06/11] fix(server): work a failed start inherited ends failed, like a started run's A started run that fails cascades failed to the work it owned (writeFinalRunEvents passes the terminal status). A run that never started cascaded cancelled under a failed run, so the same failure looked different depending on whether the provider had started. Nothing supersedes this work: the run it belongs to failed, so it now ends failed too. Runtime requests stay cancelled, as their schema has no failed status. The request guard now also keeps a resolved request's node and item when the cleanup would have settled them failed or interrupted, not only cancelled. That widening applies only with guardPendingRequestCancellations; the ingestor's user-input guard is unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/orchestration-v2/EventSink.ts | 8 ++++++-- .../ProviderTurnStartService.ts | 3 ++- .../SelectionRestart.integration.test.ts | 18 +++++++++--------- 3 files changed, 17 insertions(+), 12 deletions(-) diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index 1cc20e281afc..ddfde6fe8cce 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -298,17 +298,21 @@ const layerBase: Layer.Layer< staleNodes.add(event.payload.nodeId); } } + // A run's cleanup settles a request's node and item with the run's own + // status, which is not always `cancelled`. + const settles = (status: string) => + status === "cancelled" || (allKinds && (status === "failed" || status === "interrupted")); return events.filter((event) => { switch (event.type) { case "runtime-request.updated": return event.payload.status !== "cancelled" || !staleRequests.has(event.payload.id); case "node.updated": - return event.payload.status !== "cancelled" || !staleNodes.has(event.payload.id); + return !settles(event.payload.status) || !staleNodes.has(event.payload.id); case "turn-item.updated": return ( (event.payload.type !== "user_input_request" && event.payload.type !== "approval_request") || - event.payload.status !== "cancelled" || + !settles(event.payload.status) || !staleRequests.has(event.payload.requestId) ); default: diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index b3156977b46b..94d41e98e293 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -266,7 +266,8 @@ export const layer: Layer.Layer< return yield* RunExecutionService.cascadeTerminalizeRunOwnedSubagents({ run, open: linked, - status: "cancelled", + // The same status a started run's failure gives the work it owned. + status: "failed", completedAt: input.now, allocateEventId: () => idAllocator.allocate.event({ threadId: run.threadId }), }); diff --git a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts index fa1a6825fb98..bf747546d9f1 100644 --- a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts +++ b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts @@ -1153,26 +1153,26 @@ it.live("settles the work a restarted run inherited when its replacement never o : [], ), [ - ["approval_request", "cancelled"], - ["subagent", "cancelled"], - ["command_execution", "cancelled"], + ["approval_request", "failed"], + ["subagent", "failed"], + ["command_execution", "failed"], ], ); // The provider-native subagent's own thread ends with the run, and so does // the thread of a nested subagent whose row already settled. assert.deepEqual(nativeChildItems, [ - ["approval_request", "cancelled"], + ["approval_request", "failed"], ["subagent", "completed"], - ["subagent", "cancelled"], - ["command_execution", "cancelled"], + ["subagent", "failed"], + ["command_execution", "failed"], ]); assert.isFalse(nativeChildStreaming); assert.deepEqual(nativeChildRequests, ["cancelled"]); // A subagent the provider-native subagent launched ends with it. - assert.deepEqual(nativeChildSubagents, ["cancelled"]); + assert.deepEqual(nativeChildSubagents, ["failed"]); // So does the subagent's own provider turn, which no run attempt owns. - assert.deepEqual(nativeChildProviderTurns, ["cancelled"]); - assert.deepEqual(nativeGrandchildItems, [["command_execution", "cancelled"]]); + assert.deepEqual(nativeChildProviderTurns, ["failed"]); + assert.deepEqual(nativeGrandchildItems, [["command_execution", "failed"]]); // The delegated task and its thread keep running. assert.deepEqual( projection.subagents.flatMap((subagent) => From a392118bd1fee8d99bd875b6c1d3331a6a99daec Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 16:18:35 -0700 Subject: [PATCH 07/11] fix(server): an approval answered mid-read keeps its node when a start fails The failure path reads nodes, pending requests and items at different moments. If the user approves after the waiting approval node is read but before the pending-requests query runs, the cascade settles the node (and item) failed with no runtime-request cancellation beside it. The request guard only found protected nodes and items through cancellation events, so the stale node overwrote the approved one: run failed, request resolved, node failed. With guardPendingRequestCancellations, the guard now checks each request node and request item a write settles against its own request inside the transaction, and drops the settlement if the request was answered. The ingestor's guardPendingUserInputCancellations mode is unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/orchestration-v2/EventSink.ts | 51 +++++- .../ProviderTurnStartService.test.ts | 150 ++++++++++++------ 2 files changed, 148 insertions(+), 53 deletions(-) diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index ddfde6fe8cce..ea513924fad7 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -277,7 +277,54 @@ const layerBase: Layer.Layer< Effect.gen(function* () { const staleRequests = new Set(); const staleNodes = new Set(); + // A run's cleanup settles a request's node and item with the run's own + // status, which is not always `cancelled`. + const settles = (status: string) => + status === "cancelled" || (allKinds && (status === "failed" || status === "interrupted")); + // Whether a request was answered after the cleanup read it. A missing + // request leaves nothing to protect. + const answered = new Map(); + const wasAnswered = (threadId: ThreadId, requestId: RuntimeRequestId) => + Effect.gen(function* () { + const known = answered.get(requestId); + if (known !== undefined) return known; + const current = yield* projectionStore.getRuntimeRequest(threadId, requestId); + const result = + current !== undefined && + (current.status !== "pending" || current.responseCapability.type === "message"); + answered.set(requestId, result); + return result; + }); for (const event of events) { + // The cleanup reads nodes, requests and items at different moments, + // so it can settle a request's node or item after the request was + // answered and left its pending read. Each one is checked against + // its own request, with or without a cancellation beside it. + if (allKinds && event.type === "node.updated") { + const requestId = event.payload.runtimeRequestId; + if ( + requestId !== null && + settles(event.payload.status) && + (yield* wasAnswered(event.payload.threadId, requestId)) + ) { + staleNodes.add(event.payload.id); + } + continue; + } + if ( + allKinds && + event.type === "turn-item.updated" && + (event.payload.type === "user_input_request" || + event.payload.type === "approval_request") + ) { + if ( + settles(event.payload.status) && + (yield* wasAnswered(event.payload.threadId, event.payload.requestId)) + ) { + staleRequests.add(event.payload.requestId); + } + continue; + } if ( event.type !== "runtime-request.updated" || (!allKinds && event.payload.kind !== "user_input") || @@ -298,10 +345,6 @@ const layerBase: Layer.Layer< staleNodes.add(event.payload.nodeId); } } - // A run's cleanup settles a request's node and item with the run's own - // status, which is not always `cancelled`. - const settles = (status: string) => - status === "cancelled" || (allKinds && (status === "failed" || status === "interrupted")); return events.filter((event) => { switch (event.type) { case "runtime-request.updated": diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index f2360629799c..ebbc85283bfb 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -1089,6 +1089,7 @@ function makePersistedStartFailureHarness() { >; readonly projectionStore?: ( store: ProjectionStore.ProjectionStoreV2Shape, + eventSink: EventSink.EventSinkV2Shape, ) => ProjectionStore.ProjectionStoreV2Shape; }) => Effect.gen(function* () { @@ -1124,7 +1125,7 @@ function makePersistedStartFailureHarness() { Layer.succeed(EventSink.EventSinkV2, gatedSink), Layer.succeed( ProjectionStore.ProjectionStoreV2, - options.projectionStore?.(store) ?? store, + options.projectionStore?.(store, eventSink) ?? store, ), IdAllocator.layer, Layer.mock(ContextHandoffService.ContextHandoffServiceV2)({}), @@ -1169,63 +1170,114 @@ function makePersistedStartFailureHarness() { }; } +// The user accepts the approval the run inherited. +const acceptApproval = ( + harness: ReturnType, + eventSink: EventSink.EventSinkV2Shape, +) => + eventSink.write({ + events: [ + { + id: EventId.make("event-approval-accepted"), + type: "runtime-request.updated", + threadId: harness.threadId, + nodeId: harness.approvalNode.id, + occurredAt: harness.now, + payload: { + ...harness.approvalRequest, + status: "resolved", + decision: "accept", + resolvedAt: harness.now, + }, + }, + { + id: EventId.make("event-approval-node-completed"), + type: "node.updated", + threadId: harness.threadId, + nodeId: harness.approvalNode.id, + occurredAt: harness.now, + payload: { ...harness.approvalNode, status: "completed", completedAt: harness.now }, + }, + { + id: EventId.make("event-approval-item-completed"), + type: "turn-item.updated", + threadId: harness.threadId, + nodeId: harness.approvalNode.id, + occurredAt: harness.now, + payload: { + ...harness.approvalItem, + status: "completed", + completedAt: harness.now, + updatedAt: harness.now, + }, + }, + ], + }); + +const expectApprovalAccepted = (projection: OrchestrationV2ThreadProjection) => { + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.runtimeRequests).toMatchObject([{ status: "resolved", decision: "accept" }]); + const approvalNode = projection.nodes.find((node) => node.kind === "approval_request"); + const approvalItem = projection.turnItems.find((item) => item.type === "approval_request"); + expect(approvalNode?.status).toBe("completed"); + expect(approvalItem?.status).toBe("completed"); +}; + effectIt.effect("keeps an approval accepted while a failed start is written", () => Effect.gen(function* () { const harness = makePersistedStartFailureHarness(); const projection = yield* harness.run({ // The user approves after the failure read the approval as pending. beforeFailureWrite: Effect.gen(function* () { - const eventSink = yield* EventSink.EventSinkV2; - const { now, threadId } = harness; - yield* eventSink.write({ - events: [ - { - id: EventId.make("event-approval-accepted"), - type: "runtime-request.updated", - threadId, - nodeId: harness.approvalNode.id, - occurredAt: now, - payload: { - ...harness.approvalRequest, - status: "resolved", - decision: "accept", - resolvedAt: now, - }, - }, - { - id: EventId.make("event-approval-node-completed"), - type: "node.updated", - threadId, - nodeId: harness.approvalNode.id, - occurredAt: now, - payload: { ...harness.approvalNode, status: "completed", completedAt: now }, - }, - { - id: EventId.make("event-approval-item-completed"), - type: "turn-item.updated", - threadId, - nodeId: harness.approvalNode.id, - occurredAt: now, - payload: { - ...harness.approvalItem, - status: "completed", - completedAt: now, - updatedAt: now, - }, - }, - ], - }); + yield* acceptApproval(harness, yield* EventSink.EventSinkV2); }), }); - expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); - expect(projection.runtimeRequests).toMatchObject([{ status: "resolved", decision: "accept" }]); - expect(projection.nodes.find((node) => node.id === harness.approvalNode.id)?.status).toBe( - "completed", - ); - expect(projection.turnItems.find((item) => item.id === harness.approvalItem.id)?.status).toBe( - "completed", - ); + expectApprovalAccepted(projection); + }), +); + +effectIt.effect("keeps an approval accepted while the failure reads what it inherited", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + let settling = false; + let accepted = false; + const projection = yield* harness.run({ + // The user approves while the failure's recovery read runs: it loaded + // the approval's node and item as waiting, but the request had left the + // pending-requests read, so no cancellation accompanies their settlement. + projectionStore: (store, eventSink) => ({ + ...store, + // Only the inherited-work read asks for subagent links. + getThreadRecords: (threadId, fields, filter) => + Effect.sync(() => { + if (fields.some((field) => field === "subagents")) settling = true; + }).pipe(Effect.andThen(store.getThreadRecords(threadId, fields, filter))), + getRuntimeRecoveryProjection: (threadId) => + store.getRuntimeRecoveryProjection(threadId).pipe( + Effect.flatMap((recovery) => + !settling || accepted + ? Effect.succeed(recovery) + : Effect.sync(() => { + accepted = true; + }).pipe( + Effect.andThen(acceptApproval(harness, eventSink)), + Effect.orDie, + Effect.as({ + ...recovery, + runtimeRequests: recovery.runtimeRequests.filter( + (request) => request.id !== harness.approvalRequest.id, + ), + }), + ), + ), + ), + }), + }); + + expect(accepted).toBe(true); + + expectApprovalAccepted(projection); }), ); From 9da787556f8f15451145e9c569a60c8f54e2ecd1 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Tue, 6 Oct 2026 16:22:24 -0700 Subject: [PATCH 08/11] fix(server): a failed start settles native subagents whose item already ended On a linked child thread, the recovery read returns a subagent row only while its turn item is still open: its run-id clause never matches on a thread with no runs. A nested subagent whose item reported done while its row was still running was therefore missing from the cascade's input and stayed running. The child-thread read now loads the thread's subagents by thread, the same way it loads provider turns. Co-Authored-By: Claude Opus 5.5 (1M context) --- .../ProviderTurnStartService.ts | 8 ++- .../SelectionRestart.integration.test.ts | 57 ++++++++++++++++++- 2 files changed, 60 insertions(+), 5 deletions(-) diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 94d41e98e293..57aa0815c7aa 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -223,11 +223,12 @@ export const layer: Layer.Layer< if (Option.isNone(child) || !isProviderNativeSubagentThread(child.value.thread)) continue; linkedChildThreadIds.add(threadId); // A child thread has no runs, so the recovery read leaves out its - // streaming replies, provider turns and its pending requests' nodes - // and items; read them by thread instead. + // streaming replies, provider turns, subagents whose item already + // settled, and its pending requests' nodes and items; read them by + // thread instead. const records = yield* projectionStore.getThreadRecords( threadId, - ["messages", "nodes", "turnItems", "providerTurns"], + ["messages", "nodes", "turnItems", "providerTurns", "subagents"], { messageRoles: ["assistant"], turnItemTypes: ["approval_request", "user_input_request"], @@ -244,6 +245,7 @@ export const layer: Layer.Layer< threads.push({ ...child.value, providerTurns: records.providerTurns, + subagents: records.subagents, nodes: [ ...child.value.nodes, ...records.nodes.filter( diff --git a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts index bf747546d9f1..949ae368aaa2 100644 --- a/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts +++ b/apps/server/src/orchestration-v2/SelectionRestart.integration.test.ts @@ -128,6 +128,9 @@ function openTurnWork( const childApprovalNodeId = NodeId.make(`node:child-approval:${input.attemptId}`); const childRequestId = RuntimeRequestId.make(`request:child-approval:${input.attemptId}`); const nestedRunningSubagentId = NodeId.make(`node:nested-running-subagent:${input.attemptId}`); + const nestedSettledItemSubagentId = NodeId.make( + `node:nested-settled-item-subagent:${input.attemptId}`, + ); const requestId = RuntimeRequestId.make(`request:approval:${input.attemptId}`); return [ { @@ -421,6 +424,54 @@ function openTurnWork( result: null, }, }, + // A third reported its item done, but its row is still running. + { + type: "subagent.updated", + driver, + subagent: { + id: nestedSettledItemSubagentId, + threadId: childThreadId, + runId: null, + parentNodeId: subagentNodeId, + origin: "provider_native", + createdBy: "agent", + driver, + providerInstanceId: input.modelSelection.instanceId, + providerThreadId: null, + childThreadId: null, + nativeTaskRef: null, + prompt: "Check the tests", + title: null, + model: null, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + { + type: "turn_item.updated", + driver, + turnItem: { + ...base, + id: TurnItemId.make(`turn-item:nested-settled-item-subagent:${input.attemptId}`), + threadId: childThreadId, + runId: null, + nodeId: null, + ordinal: 5, + status: "completed", + completedAt: now, + type: "subagent", + subagentId: nestedSettledItemSubagentId, + origin: "provider_native", + driver, + providerInstanceId: input.modelSelection.instanceId, + childThreadId: null, + prompt: "Check the tests", + result: "done", + }, + }, { type: "turn_item.updated", driver, @@ -1164,12 +1215,14 @@ it.live("settles the work a restarted run inherited when its replacement never o ["approval_request", "failed"], ["subagent", "completed"], ["subagent", "failed"], + ["subagent", "completed"], ["command_execution", "failed"], ]); assert.isFalse(nativeChildStreaming); assert.deepEqual(nativeChildRequests, ["cancelled"]); - // A subagent the provider-native subagent launched ends with it. - assert.deepEqual(nativeChildSubagents, ["failed"]); + // Subagents the provider-native subagent launched end with it, including + // one whose item already reported done. + assert.deepEqual(nativeChildSubagents, ["failed", "failed"]); // So does the subagent's own provider turn, which no run attempt owns. assert.deepEqual(nativeChildProviderTurns, ["failed"]); assert.deepEqual(nativeGrandchildItems, [["command_execution", "failed"]]); From 9933fa4a328e837760dc4f26babca742539f6cbc Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Wed, 7 Oct 2026 17:08:02 -0700 Subject: [PATCH 09/11] fix(server): a failed start's cascade keeps a child that completed meanwhile A selection restart detaches the earlier attempt's session only after interrupting its turn, and the earlier attempt's ingest fiber keeps writing a provider-native child's subagent, node and item updates without a run gate. When the restarted attempt then fails to start, the failure reads the run's open subagents and commits a cascade that settles them `failed`. A child that finished between that read and the commit had its completed status and result overwritten. The failure write now sets a new `guardSettledWork` option on `writeIfRunCurrent`. Inside the write transaction it rereads each subagent, node and turn item the batch updates and drops the update when the row has already ended (completed, failed, cancelled or interrupted). Rows that are still open settle as before. Co-Authored-By: Claude Opus 5.5 (1M context) --- apps/server/src/orchestration-v2/EventSink.ts | 68 ++++++++++- .../ProviderTurnStartService.test.ts | 114 ++++++++++++++++++ .../ProviderTurnStartService.ts | 3 + 3 files changed, 180 insertions(+), 5 deletions(-) diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index ea513924fad7..a36711999968 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -90,6 +90,12 @@ export interface EventSinkV2Shape { * provider's own approval cancellation stays authoritative. */ readonly guardPendingRequestCancellations?: boolean; + /** + * Drops subagent, node and turn-item updates for rows that ended after an + * earlier projection read showed them open. For a writer that settles + * work it no longer receives events for. + */ + readonly guardSettledWork?: boolean; readonly commandId?: CommandId; readonly threadId: ThreadId; readonly runId: RunId; @@ -363,16 +369,68 @@ const layerBase: Layer.Layer< } }); }); + // A cleanup settles subagents, nodes and items it read as open. One that + // ended on its own before the commit keeps its own outcome and result. + const guardSettledWork = (events: ReadonlyArray) => + Effect.gen(function* () { + const ended = (status: string) => + status === "completed" || + status === "failed" || + status === "cancelled" || + status === "interrupted"; + const threads = new Map< + ThreadId, + ProjectionStore.ProjectionRecords<"nodes" | "subagents"> | undefined + >(); + const records = (threadId: ThreadId) => + Effect.gen(function* () { + if (threads.has(threadId)) return threads.get(threadId); + const read = yield* projectionStore + .getThreadRecords(threadId, ["nodes", "subagents"]) + .pipe( + Effect.catchTags({ + ProjectionStoreThreadNotFoundError: () => Effect.succeed(undefined), + }), + ); + threads.set(threadId, read); + return read; + }); + const kept: Array = []; + for (const event of events) { + const current = + event.type === "subagent.updated" + ? (yield* records(event.payload.threadId))?.subagents.find( + (row) => row.id === event.payload.id, + ) + : event.type === "node.updated" + ? (yield* records(event.payload.threadId))?.nodes.find( + (row) => row.id === event.payload.id, + ) + : event.type === "turn-item.updated" + ? yield* projectionStore.getTurnItem({ + threadId: event.payload.threadId, + itemId: event.payload.id, + }) + : undefined; + if (current == null || !ended(current.status)) kept.push(event); + } + return kept; + }); const guardCancellations = (input: { readonly guardPendingUserInputCancellations?: boolean; readonly guardPendingRequestCancellations?: boolean; + readonly guardSettledWork?: boolean; readonly events: ReadonlyArray; }) => - input.guardPendingRequestCancellations === true - ? guardRequestCancellations(input.events, true) - : input.guardPendingUserInputCancellations === true - ? guardRequestCancellations(input.events, false) - : Effect.succeed(input.events); + Effect.gen(function* () { + const events: ReadonlyArray = + input.guardPendingRequestCancellations === true + ? yield* guardRequestCancellations(input.events, true) + : input.guardPendingUserInputCancellations === true + ? yield* guardRequestCancellations(input.events, false) + : input.events; + return input.guardSettledWork === true ? yield* guardSettledWork(events) : events; + }); const normalizeEvents = (events: ReadonlyArray) => { const runOrdinals = new Map( diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index ebbc85283bfb..84695b5619c4 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -21,6 +21,7 @@ import { type OrchestrationV2ExecutionNode, type OrchestrationV2Run, type OrchestrationV2RuntimeRequest, + type OrchestrationV2Subagent, type OrchestrationV2ThreadProjection, type OrchestrationV2TurnItem, OrchestrationV2DomainEvent, @@ -1082,6 +1083,8 @@ function makePersistedStartFailureHarness() { // `beforeFailureWrite` runs once the start has read its projection and // given up on the provider, right before the failure is committed. const run_ = (options: { + /** Rows the run inherited besides the approval. */ + readonly seed?: ReadonlyArray; readonly beforeFailureWrite?: Effect.Effect< void, EventSink.EventSinkV2Error, @@ -1108,6 +1111,7 @@ function makePersistedStartFailureHarness() { }), ), }); + if (options.seed !== undefined) yield* eventSink.write({ events: options.seed }); const gatedSink = EventSink.EventSinkV2.of({ ...eventSink, writeIfRunCurrent: (input) => @@ -1161,8 +1165,12 @@ function makePersistedStartFailureHarness() { }).pipe(Effect.provide(Layer.fresh(layerPersistence))); return { now, + driver, + instanceId, threadId, runId, + rootNodeId, + providerThreadId, approvalRequest, approvalNode, approvalItem, @@ -1308,3 +1316,109 @@ effectIt.effect("fails the run when the work it inherited cannot be read", () => ]); }), ); + +effectIt.effect("keeps a native subagent that completed while a failed start is written", () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const subagentNode: OrchestrationV2ExecutionNode = { + ...harness.approvalNode, + id: NodeId.make("node-subagent-start-failure"), + kind: "subagent", + status: "running", + runtimeRequestId: null, + }; + const subagent: OrchestrationV2Subagent = { + id: subagentNode.id, + threadId: harness.threadId, + runId: harness.runId, + parentNodeId: harness.rootNodeId, + origin: "provider_native", + createdBy: "agent", + driver: harness.driver, + providerInstanceId: harness.instanceId, + providerThreadId: harness.providerThreadId, + childThreadId: null, + nativeTaskRef: null, + prompt: "Explore the repo", + title: "Explorer", + model: null, + status: "running", + result: null, + startedAt: harness.now, + completedAt: null, + updatedAt: harness.now, + }; + const subagentItem: OrchestrationV2TurnItem = { + id: TurnItemId.make("item-subagent-start-failure"), + threadId: harness.threadId, + runId: harness.runId, + nodeId: subagentNode.id, + providerThreadId: harness.providerThreadId, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 2, + status: "running", + title: "Explorer", + startedAt: harness.now, + completedAt: null, + updatedAt: harness.now, + type: "subagent", + subagentId: subagent.id, + origin: "provider_native", + driver: harness.driver, + providerInstanceId: harness.instanceId, + childThreadId: null, + prompt: subagent.prompt, + result: null, + }; + const base = { threadId: harness.threadId, runId: harness.runId, occurredAt: harness.now }; + const subagentEvents = ( + label: string, + update: { readonly status: "running" | "completed"; readonly result: string | null }, + ): ReadonlyArray => { + const completedAt = update.status === "completed" ? harness.now : null; + return [ + { + ...base, + id: EventId.make(`event-subagent-node-${label}`), + type: "node.updated", + nodeId: subagentNode.id, + payload: { ...subagentNode, status: update.status, completedAt }, + }, + { + ...base, + id: EventId.make(`event-subagent-${label}`), + type: "subagent.updated", + nodeId: subagent.id, + payload: { ...subagent, ...update, completedAt }, + }, + { + ...base, + id: EventId.make(`event-subagent-item-${label}`), + type: "turn-item.updated", + nodeId: subagentNode.id, + payload: { ...subagentItem, ...update, completedAt }, + }, + ]; + }; + const projection = yield* harness.run({ + seed: subagentEvents("running", { status: "running", result: null }), + // The earlier attempt's session still reports on its native subagent; + // it finishes after the failure read it as running. + beforeFailureWrite: Effect.gen(function* () { + yield* (yield* EventSink.EventSinkV2).write({ + events: subagentEvents("completed", { status: "completed", result: "Found it" }), + }); + }), + }); + + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.subagents).toMatchObject([{ status: "completed", result: "Found it" }]); + expect(projection.nodes.find((node) => node.id === subagentNode.id)?.status).toBe("completed"); + expect(projection.turnItems.find((item) => item.id === subagentItem.id)).toMatchObject({ + status: "completed", + result: "Found it", + }); + }), +); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 57aa0815c7aa..e33d6313d901 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -466,6 +466,9 @@ export const layer: Layer.Layer< // An answer or approval that lands while the failure is written wins // over cancelling the request it resolved. guardPendingRequestCancellations: true, + // A provider-native child the earlier attempt still reports on can + // finish between the inherited-work read and this commit. + guardSettledWork: true, events: status === "failed" ? [ From ba3fa152db07abf25bd04f59138111cbf224b652 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Sat, 10 Oct 2026 20:44:07 -0700 Subject: [PATCH 10/11] fix(server): a Codex approval answered right away completes its item Codex emitted a request's runtime-request event before its turn item, and each event commits on its own. An answer that landed between the two found no item to settle, so the item stayed `waiting` after the request and node resolved. Emit the node and item first and the request last, so a request is only answerable once everything it settles exists. Co-Authored-By: Claude Opus 5.5 --- .../Adapters/CodexAdapterV2.ts | 96 +++++++++++-------- 1 file changed, 56 insertions(+), 40 deletions(-) diff --git a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts index 398337b22e2c..fa1a28685db4 100644 --- a/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts +++ b/apps/server/src/orchestration-v2/Adapters/CodexAdapterV2.ts @@ -4975,16 +4975,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); yield* Ref.update(finalAnswerItemIdsByTurn, (current) => { const ids = current.get(payload.turnId); @@ -5082,16 +5084,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5142,16 +5146,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5203,16 +5209,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5288,16 +5296,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5347,16 +5357,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5407,16 +5419,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(decision).pipe( @@ -5465,16 +5479,18 @@ export const makeCodexAdapterV2 = Effect.fn("makeCodexAdapterV2")(function* ( driver: CODEX_PROVIDER, node: artifacts.node, }); + // The request is answerable once it is recorded, so its node and item + // land first and an answer always finds them to settle. yield* emitProviderEvent({ - type: "runtime_request.updated", + type: "turn_item.updated", driver: CODEX_PROVIDER, - threadId: artifacts.node.threadId, - runtimeRequest: artifacts.request, + turnItem: artifacts.turnItem, }); yield* emitProviderEvent({ - type: "turn_item.updated", + type: "runtime_request.updated", driver: CODEX_PROVIDER, - turnItem: artifacts.turnItem, + threadId: artifacts.node.threadId, + runtimeRequest: artifacts.request, }); const resolved = yield* Deferred.await(answers).pipe( From 366a6415e5bccabf71777290f9382c1595653a27 Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Sat, 10 Oct 2026 21:44:37 -0700 Subject: [PATCH 11/11] fix(server): a failed restart keeps a child's final reply and still fails when its recheck cannot read When a restarted run's replacement attempt never starts, the failure write settles the work the run inherited. The event sink rechecked that settlement in the write transaction, but two gaps remained: - A recheck read error (getThreadRecords, getTurnItem, getRuntimeRequest) aborted the whole write, leaving the run and attempt `starting`. - Streaming messages, streaming turn items and provider turns were written from the earlier read's snapshot. A provider-native child that finished or streamed further before the commit had its final reply truncated, and a provider turn that completed was rewritten `failed`. writeIfRunCurrent now takes the inherited work as a separate `settlement`. The sink rechecks it in the transaction. Rows that ended since are dropped, and rows still open are settled from their current state, so text that grew survives. If the recheck cannot read, the settlement is dropped with a warning and the run, attempt and root still commit. Co-Authored-By: Claude Opus 5.5 --- apps/server/src/orchestration-v2/EventSink.ts | 160 +++++++++++++----- .../ProviderTurnStartService.test.ts | 154 ++++++++++++++++- .../ProviderTurnStartService.ts | 46 +++-- 3 files changed, 289 insertions(+), 71 deletions(-) diff --git a/apps/server/src/orchestration-v2/EventSink.ts b/apps/server/src/orchestration-v2/EventSink.ts index a36711999968..2473dedb49fb 100644 --- a/apps/server/src/orchestration-v2/EventSink.ts +++ b/apps/server/src/orchestration-v2/EventSink.ts @@ -7,10 +7,12 @@ import { RunAttemptId, RunId, RuntimeRequestId, + MessageId, NodeId, type ProjectId, ThreadId, } from "@t3tools/contracts"; +import * as Cause from "effect/Cause"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -85,17 +87,14 @@ export interface EventSinkV2Shape { readonly writeIfRunCurrent: (input: { readonly guardPendingUserInputCancellations?: boolean; /** - * Guards approvals and every other request kind too. For a writer that - * cancels requests an earlier projection read showed pending; a - * provider's own approval cancellation stays authoritative. + * Best-effort settlement of work an earlier projection read showed open, + * committed ahead of `events`. For a writer that settles work it no longer + * receives events for. It is rechecked in the transaction: a request + * answered or a row that ended since keeps its own outcome, and a row still + * open is settled from its current state. If the recheck cannot read, the + * settlement is dropped and `events` still commit. */ - readonly guardPendingRequestCancellations?: boolean; - /** - * Drops subagent, node and turn-item updates for rows that ended after an - * earlier projection read showed them open. For a writer that settles - * work it no longer receives events for. - */ - readonly guardSettledWork?: boolean; + readonly settlement?: ReadonlyArray; readonly commandId?: CommandId; readonly threadId: ThreadId; readonly runId: RunId; @@ -369,24 +368,36 @@ const layerBase: Layer.Layer< } }); }); - // A cleanup settles subagents, nodes and items it read as open. One that - // ended on its own before the commit keeps its own outcome and result. - const guardSettledWork = (events: ReadonlyArray) => + // A cleanup settles work it read as open. A row that ended on its own + // before the commit keeps its own outcome and result. One still open is + // settled from its current row, so a reply that streamed further since the + // read keeps its text. + const refreshSettlement = (events: ReadonlyArray) => Effect.gen(function* () { const ended = (status: string) => status === "completed" || status === "failed" || status === "cancelled" || status === "interrupted"; + const messageIds = new Map>(); + for (const event of events) { + if (event.type !== "message.updated") continue; + const ids = messageIds.get(event.payload.threadId) ?? []; + ids.push(event.payload.id); + messageIds.set(event.payload.threadId, ids); + } const threads = new Map< ThreadId, - ProjectionStore.ProjectionRecords<"nodes" | "subagents"> | undefined + | ProjectionStore.ProjectionRecords<"nodes" | "subagents" | "providerTurns" | "messages"> + | undefined >(); const records = (threadId: ThreadId) => Effect.gen(function* () { if (threads.has(threadId)) return threads.get(threadId); const read = yield* projectionStore - .getThreadRecords(threadId, ["nodes", "subagents"]) + .getThreadRecords(threadId, ["nodes", "subagents", "providerTurns", "messages"], { + messageIds: messageIds.get(threadId) ?? [], + }) .pipe( Effect.catchTags({ ProjectionStoreThreadNotFoundError: () => Effect.succeed(undefined), @@ -397,40 +408,100 @@ const layerBase: Layer.Layer< }); const kept: Array = []; for (const event of events) { - const current = - event.type === "subagent.updated" - ? (yield* records(event.payload.threadId))?.subagents.find( - (row) => row.id === event.payload.id, - ) - : event.type === "node.updated" - ? (yield* records(event.payload.threadId))?.nodes.find( - (row) => row.id === event.payload.id, - ) - : event.type === "turn-item.updated" - ? yield* projectionStore.getTurnItem({ - threadId: event.payload.threadId, - itemId: event.payload.id, - }) - : undefined; - if (current == null || !ended(current.status)) kept.push(event); + switch (event.type) { + case "subagent.updated": { + const current = (yield* records(event.payload.threadId))?.subagents.find( + (row) => row.id === event.payload.id, + ); + if (current === undefined || !ended(current.status)) kept.push(event); + break; + } + case "node.updated": { + const current = (yield* records(event.payload.threadId))?.nodes.find( + (row) => row.id === event.payload.id, + ); + if (current === undefined || !ended(current.status)) kept.push(event); + break; + } + case "turn-item.updated": { + const current = yield* projectionStore.getTurnItem({ + threadId: event.payload.threadId, + itemId: event.payload.id, + }); + if (current === null) kept.push(event); + else if (!ended(current.status)) { + const settled = { + status: event.payload.status, + completedAt: event.payload.completedAt, + updatedAt: event.payload.updatedAt, + }; + kept.push({ + ...event, + payload: + "streaming" in current + ? { ...current, ...settled, streaming: false } + : { ...current, ...settled }, + }); + } + break; + } + case "message.updated": { + const current = (yield* records(event.payload.threadId))?.messages.find( + (row) => row.id === event.payload.id, + ); + if (current === undefined) kept.push(event); + else if (current.streaming) { + kept.push({ + ...event, + payload: { ...current, streaming: false, updatedAt: event.payload.updatedAt }, + }); + } + break; + } + case "provider-turn.updated": { + const current = (yield* records(event.threadId))?.providerTurns.find( + (row) => row.id === event.payload.id, + ); + if (current === undefined) kept.push(event); + else if (!ended(current.status)) { + kept.push({ + ...event, + payload: { + ...current, + status: event.payload.status, + completedAt: event.payload.completedAt, + }, + }); + } + break; + } + default: + kept.push(event); + } } return kept; }); + // The settlement is optional; the events it rides with are not. A recheck + // that cannot read drops the whole settlement rather than the write. + const guardSettlement = (events: ReadonlyArray) => + guardRequestCancellations(events, true).pipe( + Effect.flatMap(refreshSettlement), + Effect.catchCauseIf( + (cause) => !Cause.hasInterruptsOnly(cause), + (cause) => + Effect.logWarning("event sink dropped a settlement it could not recheck", { + eventCount: events.length, + cause: Cause.pretty(cause), + }).pipe(Effect.as([])), + ), + ); const guardCancellations = (input: { readonly guardPendingUserInputCancellations?: boolean; - readonly guardPendingRequestCancellations?: boolean; - readonly guardSettledWork?: boolean; readonly events: ReadonlyArray; }) => - Effect.gen(function* () { - const events: ReadonlyArray = - input.guardPendingRequestCancellations === true - ? yield* guardRequestCancellations(input.events, true) - : input.guardPendingUserInputCancellations === true - ? yield* guardRequestCancellations(input.events, false) - : input.events; - return input.guardSettledWork === true ? yield* guardSettledWork(events) : events; - }); + input.guardPendingUserInputCancellations === true + ? guardRequestCancellations(input.events, false) + : Effect.succeed(input.events); const normalizeEvents = (events: ReadonlyArray) => { const runOrdinals = new Map( @@ -550,7 +621,10 @@ const layerBase: Layer.Layer< }; } - const normalized = yield* normalizeEvents(yield* guardCancellations(input)); + const normalized = yield* normalizeEvents([ + ...(input.settlement === undefined ? [] : yield* guardSettlement(input.settlement)), + ...(yield* guardCancellations(input)), + ]); const storedEvents = yield* eventStore.append({ ...(input.commandId === undefined ? {} : { commandId: input.commandId }), events: normalized, diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts index 2f3678f944f7..709c7bf350d8 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.test.ts @@ -19,6 +19,7 @@ import { ProjectId, type OrchestrationV2AppThread, type OrchestrationV2ExecutionNode, + type OrchestrationV2ProviderTurn, type OrchestrationV2Run, type OrchestrationV2RuntimeRequest, type OrchestrationV2Subagent, @@ -1078,8 +1079,29 @@ function makePersistedStartFailureHarness() { const layerStores = Layer.mergeAll(EventStore.layer, ProjectionStore.layer).pipe( Layer.provideMerge(layerDatabase), ); - const layerSink = EventSink.layer.pipe(Layer.provide(layerStores)); - const layerPersistence = Layer.mergeAll(layerStores, layerSink, IdAllocator.layer); + // The sink can read through its own view of the projection store. + const layerPersistence = ( + sinkStore?: ( + store: ProjectionStore.ProjectionStoreV2Shape, + ) => ProjectionStore.ProjectionStoreV2Shape, + ) => + Layer.mergeAll( + layerStores, + EventSink.layer.pipe( + Layer.provide( + sinkStore === undefined + ? Layer.empty + : Layer.effect( + ProjectionStore.ProjectionStoreV2, + Effect.gen(function* () { + return sinkStore(yield* ProjectionStore.ProjectionStoreV2); + }), + ), + ), + Layer.provide(layerStores), + ), + IdAllocator.layer, + ); // `beforeFailureWrite` runs once the start has read its projection and // given up on the provider, right before the failure is committed. const run_ = (options: { @@ -1094,6 +1116,10 @@ function makePersistedStartFailureHarness() { store: ProjectionStore.ProjectionStoreV2Shape, eventSink: EventSink.EventSinkV2Shape, ) => ProjectionStore.ProjectionStoreV2Shape; + /** What the event sink reads inside its write transaction. */ + readonly sinkProjectionStore?: ( + store: ProjectionStore.ProjectionStoreV2Shape, + ) => ProjectionStore.ProjectionStoreV2Shape; }) => Effect.gen(function* () { const eventSink = yield* EventSink.EventSinkV2; @@ -1162,13 +1188,14 @@ function makePersistedStartFailureHarness() { ); yield* service.start({ threadId, runId }); return yield* store.getThreadProjection(threadId); - }).pipe(Effect.provide(Layer.fresh(layerPersistence))); + }).pipe(Effect.provide(Layer.fresh(layerPersistence(options.sinkProjectionStore)))); return { now, driver, instanceId, threadId, runId, + attemptId, rootNodeId, providerThreadId, approvalRequest, @@ -1422,3 +1449,124 @@ effectIt.effect("keeps a native subagent that completed while a failed start is }); }), ); + +// The guarded settlement is optional; failing the run is not. +const unavailable = (threadId: ThreadId) => + Effect.fail( + new ProjectionStore.ProjectionStoreReadError({ threadId, cause: "database unavailable" }), + ); +const failingRechecks: ReadonlyArray< + readonly [ + string, + (store: ProjectionStore.ProjectionStoreV2Shape) => ProjectionStore.ProjectionStoreV2Shape, + ] +> = [ + ["request", (store) => ({ ...store, getRuntimeRequest: unavailable })], + ["records", (store) => ({ ...store, getThreadRecords: unavailable })], + ["turn item", (store) => ({ ...store, getTurnItem: ({ threadId }) => unavailable(threadId) })], +]; +for (const [read, sinkProjectionStore] of failingRechecks) { + effectIt.effect(`fails the run when the write cannot recheck an inherited ${read}`, () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const projection = yield* harness.run({ sinkProjectionStore }); + + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect(projection.attempts.map((attempt) => attempt.status)).toEqual(["failed"]); + expect(projection.nodes.find((node) => node.id === harness.rootNodeId)?.status).toBe( + "failed", + ); + // The settlement it could not recheck is left as it was. + expect(projection.runtimeRequests).toMatchObject([{ status: "pending" }]); + expect(projection.turnItems).toMatchObject([ + { type: "approval_request", status: "waiting" }, + { type: "error", title: "Provider session failed to open" }, + ]); + }), + ); +} + +effectIt.effect( + "keeps replies and provider turns that moved on while a failed start is written", + () => + Effect.gen(function* () { + const harness = makePersistedStartFailureHarness(); + const base = { threadId: harness.threadId, runId: harness.runId, occurredAt: harness.now }; + const reply = ( + id: string, + update: { readonly text: string; readonly streaming: boolean }, + ): OrchestrationV2DomainEvent => ({ + ...base, + id: EventId.make(`event-${id}-${update.text.length}`), + type: "message.updated", + nodeId: harness.rootNodeId, + payload: { + id: MessageId.make(id), + threadId: harness.threadId, + runId: harness.runId, + nodeId: harness.rootNodeId, + role: "assistant", + createdBy: "agent", + creationSource: "provider", + attachments: [], + createdAt: harness.now, + updatedAt: harness.now, + ...update, + }, + }); + const providerTurn: OrchestrationV2ProviderTurn = { + id: ProviderTurnId.make("provider-turn-earlier-attempt"), + providerThreadId: harness.providerThreadId, + nodeId: harness.rootNodeId, + runAttemptId: harness.attemptId, + nativeTurnRef: null, + ordinal: 1, + status: "running", + startedAt: harness.now, + completedAt: null, + }; + const turn = (status: "running" | "completed"): OrchestrationV2DomainEvent => ({ + ...base, + id: EventId.make(`event-provider-turn-${status}`), + type: "provider-turn.updated", + nodeId: harness.rootNodeId, + payload: { + ...providerTurn, + status, + completedAt: status === "completed" ? harness.now : null, + }, + }); + const projection = yield* harness.run({ + seed: [ + reply("message-finished", { text: "Partial", streaming: true }), + reply("message-growing", { text: "Still", streaming: true }), + turn("running"), + ], + // The earlier attempt's session still reports after the failure read + // these as streaming and running: one reply finishes, the other grows, + // and the provider turn completes. + beforeFailureWrite: Effect.gen(function* () { + yield* (yield* EventSink.EventSinkV2).write({ + events: [ + reply("message-finished", { text: "Partial reply, now final", streaming: false }), + reply("message-growing", { text: "Still writing more", streaming: true }), + turn("completed"), + ], + }); + }), + }); + + expect(projection.runs.map((run) => run.status)).toEqual(["failed"]); + expect( + projection.messages + .filter((message) => message.role === "assistant") + .map(({ id, text, streaming }) => ({ id, text, streaming })), + ).toEqual([ + { id: "message-finished", text: "Partial reply, now final", streaming: false }, + { id: "message-growing", text: "Still writing more", streaming: false }, + ]); + expect(projection.providerTurns.find((row) => row.id === providerTurn.id)?.status).toBe( + "completed", + ); + }), +); diff --git a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts index 791556d6a5fb..d9982b0c4260 100644 --- a/apps/server/src/orchestration-v2/ProviderTurnStartService.ts +++ b/apps/server/src/orchestration-v2/ProviderTurnStartService.ts @@ -464,31 +464,27 @@ export const layer: Layer.Layer< runId, activeAttemptId: attempt.id, expectedStatus: "starting", - // An answer or approval that lands while the failure is written wins - // over cancelling the request it resolved. - guardPendingRequestCancellations: true, - // A provider-native child the earlier attempt still reports on can - // finish between the inherited-work read and this commit. - guardSettledWork: true, - events: - status === "failed" - ? [ - // Failing the run matters more than settling what it - // inherited: a read that fails here must not leave it - // `starting`. - ...(yield* inheritedWorkSettlement({ run, now }).pipe( - Effect.catchCause((cause) => - Cause.hasInterruptsOnly(cause) - ? Effect.failCause(cause) - : Effect.logWarning( - "provider turn start could not settle the work a failed run inherited", - { threadId: projection.thread.id, runId, cause: Cause.pretty(cause) }, - ).pipe(Effect.as([])), - ), - )), - ...events, - ] - : events, + // The sink rechecks this in the transaction: an answer, or a + // provider-native child the earlier attempt still reports on, can + // land between the inherited-work read and this commit. + ...(status === "failed" + ? { + // Failing the run matters more than settling what it + // inherited: a read that fails here must not leave it + // `starting`. + settlement: yield* inheritedWorkSettlement({ run, now }).pipe( + Effect.catchCause((cause) => + Cause.hasInterruptsOnly(cause) + ? Effect.failCause(cause) + : Effect.logWarning( + "provider turn start could not settle the work a failed run inherited", + { threadId: projection.thread.id, runId, cause: Cause.pretty(cause) }, + ).pipe(Effect.as([])), + ), + ), + } + : {}), + events, }); }, );