Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
226 changes: 226 additions & 0 deletions apps/server/src/orchestration-v2/FoundationPersistence.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
23 changes: 22 additions & 1 deletion apps/server/src/orchestration-v2/ProjectionStore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3363,6 +3363,16 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
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
Expand Down Expand Up @@ -3730,7 +3740,8 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
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')
Expand Down Expand Up @@ -3867,6 +3878,16 @@ export const layer: Layer.Layer<ProjectionStoreV2, never, SqlClient.SqlClient> =
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
`,
Expand Down
55 changes: 55 additions & 0 deletions apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -439,6 +439,61 @@ export const make = Effect.gen(function* () {
});
}
}
// A provider-native subagent thread has no runs: its work is a runless

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The new runless-root-turn recovery path needs a focused test. Existing recovery tests cover runless background items and run-backed nodes, but not discovery and cancellation of a runless child root_turn with a streaming reasoning item. Could you add a test using the real projection store and test layers for external services that verifies the child is selected, the node and item are cancelled, and streaming becomes false?

Posted via Macroscope — Effect Service Conventions

// 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 },
Comment thread
macroscopeapp[bot] marked this conversation as resolved.
});
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.
Expand Down
Loading