diff --git a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts index 1e0952b4aea0..92ce52358197 100644 --- a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts +++ b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts @@ -2802,6 +2802,232 @@ it.layer(TestLayer)("orchestration V2 foundation persistence", (it) => { }), ); + it.effect("settles a native subagent's child thread when its provider process is gone", () => + Effect.gen(function* () { + const eventSink = yield* EventSinkV2; + const projectionStore = yield* ProjectionStoreV2; + const now = yield* DateTime.now; + const parentId = ThreadId.make("thread:foundation-native-subagent-parent"); + const childId = ThreadId.make("thread:foundation-native-subagent-child"); + const runId = RunId.make("run:foundation-native-subagent"); + const subagentId = NodeId.make("node:foundation-native-subagent"); + const childRootId = NodeId.make("node:foundation-native-subagent-child-root"); + const parent = makeThread(parentId, now); + const child: OrchestrationV2AppThread = { + ...makeThread(childId, now), + createdBy: "agent", + creationSource: "provider", + lineage: { + parentThreadId: parentId, + relationshipToParent: "subagent", + rootThreadId: parentId, + }, + forkedFrom: { type: "node", nodeId: subagentId }, + }; + const node = (input: { + readonly id: NodeId; + readonly threadId: ThreadId; + readonly runId: RunId | null; + readonly kind: "root_turn" | "subagent"; + }) => ({ + ...input, + parentNodeId: null, + rootNodeId: input.id, + status: "running" as const, + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }); + // The parent run settled while its background subagent kept working, + // then the server died. Recovery already cancels the parent's subagent + // item, entity, and node; the child's runless root turn lives on another + // thread and must be settled too. + yield* eventSink.commitCommand({ + commandId: CommandId.make("command:foundation-native-subagent"), + threadId: parentId, + commandType: "foundation.native-subagent", + acceptedAt: now, + events: [ + threadCreatedEvent({ + id: "event:foundation-native-subagent:parent", + thread: parent, + now, + }), + threadCreatedEvent({ id: "event:foundation-native-subagent:child", thread: child, now }), + { + id: EventId.make("event:foundation-native-subagent:run"), + type: "run.created", + threadId: parentId, + runId, + providerInstanceId, + occurredAt: now, + payload: { + id: runId, + threadId: parentId, + ordinal: 1, + providerInstanceId, + modelSelection, + providerThreadId: null, + userMessageId: MessageId.make("message:foundation-native-subagent"), + rootNodeId: null, + activeAttemptId: null, + status: "completed", + queuePosition: null, + requestedAt: now, + startedAt: now, + completedAt: now, + checkpointId: null, + contextHandoffId: null, + }, + }, + { + id: EventId.make("event:foundation-native-subagent:subagent-node"), + type: "node.updated", + threadId: parentId, + runId, + nodeId: subagentId, + occurredAt: now, + payload: node({ id: subagentId, threadId: parentId, runId, kind: "subagent" }), + }, + { + id: EventId.make("event:foundation-native-subagent:child-root"), + type: "node.updated", + threadId: childId, + nodeId: childRootId, + occurredAt: now, + payload: node({ id: childRootId, threadId: childId, runId: null, kind: "root_turn" }), + }, + { + // Claude's live "Subagent progress" item in the child, as the + // adapter writes it while task_progress frames arrive. + id: EventId.make("event:foundation-native-subagent:child-progress"), + type: "turn-item.updated", + threadId: childId, + nodeId: childRootId, + occurredAt: now, + payload: { + id: TurnItemId.make("item:foundation-native-subagent:progress"), + threadId: childId, + runId: null, + nodeId: childRootId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 101, + type: "reasoning", + status: "running", + title: "Subagent progress", + startedAt: now, + completedAt: null, + updatedAt: now, + text: "Running git diff --stat", + streaming: true, + }, + }, + { + id: EventId.make("event:foundation-native-subagent:subagent"), + type: "subagent.updated", + threadId: parentId, + runId, + nodeId: subagentId, + driver: providerDriver, + providerInstanceId, + occurredAt: now, + payload: { + id: subagentId, + threadId: parentId, + runId, + parentNodeId: subagentId, + origin: "provider_native", + createdBy: "agent", + driver: providerDriver, + providerInstanceId, + providerThreadId: null, + childThreadId: childId, + nativeTaskRef: null, + prompt: "Audit the adapters", + title: null, + model: null, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + { + id: EventId.make("event:foundation-native-subagent:item"), + type: "turn-item.updated", + threadId: parentId, + runId, + nodeId: subagentId, + occurredAt: now, + payload: { + id: TurnItemId.make("item:foundation-native-subagent"), + threadId: parentId, + runId, + nodeId: subagentId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + type: "subagent", + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + subagentId, + origin: "provider_native", + driver: providerDriver, + providerInstanceId, + childThreadId: childId, + prompt: "Audit the adapters", + result: null, + }, + }, + ], + effects: [], + }); + + const recovery = yield* ProviderRuntimeRecovery.make.pipe( + Effect.provide(ServerSettings.layerTest()), + Effect.provideService( + OrchestrationEffectWorkerV2, + OrchestrationEffectWorkerV2.of({ + awaitWork: Effect.void, + runRecoveryOnce: Effect.succeed(false), + runOnce: Effect.succeed(false), + nextClaimableAt: Effect.succeed(Option.none()), + drain: () => Effect.succeed(0), + }), + ), + ); + assert.include(yield* projectionStore.getRecoveryThreadIds("runtime"), childId); + yield* recovery.recover; + + const parentProjection = yield* projectionStore.getThreadProjection(parentId); + assert.equal(parentProjection.subagents[0]?.status, "cancelled"); + const childProjection = yield* projectionStore.getThreadProjection(childId); + const childRoot = childProjection.nodes.find((candidate) => candidate.id === childRootId); + assert.equal(childRoot?.status, "cancelled"); + assert.isNotNull(childRoot?.completedAt ?? null); + // Nothing inside the child keeps reading as live work either. + const progress = childProjection.turnItems.find((item) => item.type === "reasoning"); + assert.equal(progress?.status, "cancelled"); + assert.isFalse(progress?.type === "reasoning" && progress.streaming); + assert.isNotNull(progress?.completedAt ?? null); + assert.notInclude(yield* projectionStore.getRecoveryThreadIds("runtime"), childId); + }), + ); + it.effect("allocates collision-free positions beyond 100 items and rebuilds equivalently", () => Effect.gen(function* () { const eventSink = yield* EventSinkV2; diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index de41d1660158..a25179eab8fe 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -3363,6 +3363,16 @@ export const layer: Layer.Layer = CROSS JOIN orchestration_v2_projection_subagents AS subagents ON subagents.provider_thread_id = pending_provider_threads.provider_thread_id UNION + SELECT subagents.child_thread_id FROM orchestration_v2_projection_subagents AS subagents + WHERE subagents.child_thread_id IS NOT NULL + AND EXISTS ( + SELECT 1 FROM orchestration_v2_projection_nodes AS node + WHERE node.thread_id = subagents.child_thread_id + AND node.run_id IS NULL + AND node.kind = 'root_turn' + AND node.status IN ('pending', 'running', 'waiting') + ) + UNION SELECT item.thread_id FROM orchestration_v2_projection_turn_items AS item WHERE NOT EXISTS ( SELECT 1 FROM orchestration_v2_projection_runs AS run @@ -3730,7 +3740,8 @@ export const layer: Layer.Layer = WHERE node.thread_id = ${threadId} AND node.status IN ('pending', 'starting', 'running', 'waiting') AND ( - node.run_id IN ( + (node.run_id IS NULL AND node.kind = 'root_turn') + OR node.run_id IN ( SELECT run_id FROM orchestration_v2_projection_runs WHERE thread_id = ${threadId} AND status IN ('queued', 'preparing', 'starting', 'running', 'waiting') @@ -3867,6 +3878,16 @@ export const layer: Layer.Layer = AND status IN ('queued', 'preparing', 'starting', 'running', 'waiting') ) OR item.type IN ('command_execution', 'dynamic_tool', 'subagent') + OR ( + item.run_id IS NULL + AND item.node_id IN ( + SELECT node_id FROM orchestration_v2_projection_nodes + WHERE thread_id = ${threadId} + AND run_id IS NULL + AND kind = 'root_turn' + AND status IN ('pending', 'running', 'waiting') + ) + ) ) ORDER BY item.ordinal ASC, item.turn_item_id ASC `, diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index a069ccb0cfdc..e64bc03be76c 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -439,6 +439,61 @@ export const make = Effect.gen(function* () { }); } } + // A provider-native subagent thread has no runs: its work is a runless + // root turn, plus items under it (Claude's live progress item), that + // only the dead provider process could settle. Left running, the child + // would show as working forever. + const cancelledStaleItemIds = new Set( + events.flatMap((event) => (event.type === "turn-item.updated" ? [event.payload.id] : [])), + ); + for (const node of projection.nodes) { + if ( + node.kind !== "root_turn" || + node.runId !== null || + !isNonterminalNodeStatus(node.status) || + cancelledStaleNodeIds.has(node.id) + ) { + continue; + } + cancelledStaleNodeIds.add(node.id); + events.push({ + id: yield* allocateEventId(), + type: "node.updated", + threadId: projection.thread.id, + nodeId: node.id, + providerInstanceId: projection.thread.providerInstanceId, + occurredAt: now, + payload: { ...node, status: "cancelled", completedAt: now }, + }); + for (const item of projection.turnItems) { + if ( + item.nodeId !== node.id || + item.runId !== null || + !isNonterminalTurnItemStatus(item.status) || + cancelledStaleItemIds.has(item.id) + ) { + continue; + } + cancelledStaleItemIds.add(item.id); + events.push({ + id: yield* allocateEventId(), + type: "turn-item.updated", + threadId: projection.thread.id, + nodeId: node.id, + providerInstanceId: projection.thread.providerInstanceId, + occurredAt: now, + payload: { + ...item, + status: "cancelled", + completedAt: now, + updatedAt: now, + ...(item.type === "reasoning" || item.type === "assistant_message" + ? { streaming: false } + : {}), + }, + }); + } + } // All provider processes are gone on startup/shutdown: clear any // persisted Waiting roster (including idle threads from settled roots) // and idle active threads without resurrecting active status.