From a2ef2d8b2b6c597b1e3a09db17b04528e2fa4577 Mon Sep 17 00:00:00 2001 From: Jake Leventhal Date: Mon, 5 Oct 2026 21:24:58 -0400 Subject: [PATCH 1/5] fix(orchestration): reconcile held delegated work without stale running state --- ...rchestratorMcpService.queued-tasks.test.ts | 247 +++++++ apps/server/src/mcp/OrchestratorMcpService.ts | 17 +- .../src/mcp/toolkits/orchestrator/tools.ts | 2 +- .../DelegatedCompletionDelivery.test.ts | 621 +++++++++++++++++- .../FoundationPersistence.test.ts | 1 + .../src/orchestration-v2/Orchestrator.ts | 187 ++++-- .../ProjectionRecovery.test.ts | 125 ++-- .../ProjectionSettlement.test.ts | 31 + .../src/orchestration-v2/ProjectionStore.ts | 26 +- .../ProviderRuntimeRecoveryService.test.ts | 58 +- .../ProviderRuntimeRecoveryService.ts | 22 + packages/contracts/src/orchestratorMcp.ts | 2 + .../src/server/subagentProjection.ts | 14 +- ...chestrationV2PendingBackgroundWork.test.ts | 67 ++ .../orchestrationV2PendingBackgroundWork.ts | 17 +- 15 files changed, 1271 insertions(+), 166 deletions(-) create mode 100644 apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts diff --git a/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts new file mode 100644 index 000000000000..dec6ca9d1fde --- /dev/null +++ b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts @@ -0,0 +1,247 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { assert, it } from "@effect/vitest"; +import { + EnvironmentId, + MessageId, + NodeId, + type OrchestrationV2ThreadProjection, + type OrchestrationV2ServerCommand, + ProjectId, + ProviderInstanceId, + RunId, + ThreadId, +} from "@t3tools/contracts"; +import * as DateTime from "effect/DateTime"; +import * as Effect from "effect/Effect"; +import * as Layer from "effect/Layer"; + +import * as ProviderAdapterRegistry from "../orchestration-v2/ProviderAdapterRegistry.ts"; +import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; +import * as ProjectService from "../project/ProjectService.ts"; +import * as SecretRequests from "../secrets/SecretRequests.ts"; +import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; +import * as ScheduledTaskService from "../scheduledTasks/ScheduledTaskService.ts"; +import type { McpInvocationScope } from "./McpInvocationContext.ts"; +import * as OrchestratorMcpService from "./OrchestratorMcpService.ts"; + +const parentThreadId = ThreadId.make("thread:queued-task-parent"); +const childThreadId = ThreadId.make("thread:queued-task-child"); +const taskId = NodeId.make("node:queued-task"); +const now = DateTime.makeUnsafe("2026-10-03T11:00:00Z"); +const providerInstanceId = ProviderInstanceId.make("claude_two"); +const scope: McpInvocationScope = { + requestNamespace: "test", + client: undefined, + environmentId: EnvironmentId.make("environment:queued-task"), + thread: { + threadId: parentThreadId, + providerSessionId: "provider-session:queued-task", + providerInstanceId, + }, + capabilities: new Set(["orchestration"]), + issuedAt: 1, +}; + +function fixture( + options: { + published?: boolean; + executing?: boolean; + failCancellation?: boolean; + secondQueued?: boolean; + deliveryPending?: boolean; + } = {}, +) { + const projectId = ProjectId.make("project:queued-task"); + const thread = { + id: childThreadId, + projectId, + title: "Child", + modelSelection: { instanceId: providerInstanceId, model: "claude-opus-5-5" }, + lineage: { parentThreadId, relationshipToParent: "subagent", rootThreadId: parentThreadId }, + createdAt: now, + updatedAt: now, + deletedAt: null, + archivedAt: null, + settledOverride: null, + settledAt: null, + runtimeMode: "full-access", + interactionMode: "default", + }; + const run = (ordinal: number, status: string) => ({ + id: RunId.make(`run:queued-task:${ordinal}`), + threadId: childThreadId, + ordinal, + status, + userMessageId: MessageId.make(`message:queued-task:${ordinal}`), + requestedAt: now, + startedAt: status === "queued" ? null : now, + completedAt: status === "failed" || status === "completed" ? now : null, + queueHeld: status === "queued", + modelSelection: thread.modelSelection, + providerInstanceId, + }); + const child = { + thread, + runs: [ + run(1, "failed"), + run(2, "queued"), + run(3, "completed"), + run(4, "completed"), + ...(options.executing ? [run(5, "running")] : []), + ...(options.secondQueued ? [run(6, "queued")] : []), + ], + messages: [ + { + id: MessageId.make("result:queued-task:4"), + runId: RunId.make("run:queued-task:4"), + role: "assistant", + text: "Latest completed follow-up result", + updatedAt: now, + }, + ], + turnItems: [], + providerThreads: [{ status: "idle", pendingBackgroundTasks: [] }], + runtimeRequests: [], + contextTransfers: [], + subagents: [], + } as unknown as OrchestrationV2ThreadProjection; + const parent = { + thread: { + ...thread, + id: parentThreadId, + lineage: { parentThreadId: null, relationshipToParent: null }, + }, + runs: [], + contextTransfers: [], + subagents: [ + { + id: taskId, + threadId: parentThreadId, + childThreadId, + origin: "app_owned", + status: options.published ? "completed" : "running", + providerInstanceId, + model: "claude-opus-5-5", + result: options.published ? "Published original result" : null, + completionDelivery: { state: options.deliveryPending ? "pending" : "disposed" }, + }, + ], + } as unknown as OrchestrationV2ThreadProjection; + const commands: OrchestrationV2ServerCommand[] = []; + const layer = OrchestratorMcpService.layer.pipe( + Layer.provide( + Layer.mergeAll( + NodeServices.layer, + Layer.mock(ThreadManagementService.ThreadManagementService)({ + getThreadShell: () => Effect.succeed(child.thread as never), + getThreadRecords: (id) => Effect.succeed(id === parentThreadId ? parent : child), + getProjectThreadRecords: () => Effect.succeed(child), + stopDelegatedTasks: () => Effect.void, + getTimelinePage: () => Effect.succeed({ items: [], totalItems: 0, hasMore: false }), + dispatch: (command) => + Effect.suspend(() => { + commands.push(command); + return options.failCancellation && command.type === "thread.stop" + ? Effect.fail(new Error("Cancellation failed") as never) + : Effect.succeed({} as never); + }), + }), + Layer.mock(ProviderRegistry.ProviderRegistry)({ getProviders: Effect.succeed([]) }), + Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({ + list: () => Effect.succeed([]), + }), + Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), + Layer.mock(ProjectService.ProjectService)({}), + ), + ), + ); + return { child, parent, commands, layer }; +} + +it.effect("reports an old held continuation as queued and keeps the latest result readable", () => { + const { commands, layer } = fixture(); + return Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const task = yield* service.taskStatus(scope, taskId); + assert.equal(task.status, "queued"); + assert.equal(task.workState, "working"); + assert.equal(task.hasPendingChildRuns, true); + assert.equal(task.latestTerminalRunId, RunId.make("run:queued-task:4")); + assert.equal(task.latestTerminalSummary, "Latest completed follow-up result"); + assert.isNull(task.summary); + const read = yield* service.readThread(scope, { threadId: childThreadId, runLimit: 1 }); + assert.equal(read.thread.status, "completed"); + assert.isNull(read.thread.activeRunId); + assert.equal(read.thread.pendingRequestCount, 0); + assert.equal(read.thread.queuedRunCount, 1); + assert.equal(read.thread.heldQueuedRunCount, 1); + assert.equal(read.recentRuns.length, 1); + assert.deepEqual(commands, []); + }).pipe(Effect.provide(layer)); +}); + +it.effect("stops queued-only unpublished work while preserving held input", () => { + const { commands, child, layer } = fixture({ secondQueued: true }); + return Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const result = yield* service.cancelTask(scope, { taskId, clientRequestId: "cancel-held" }); + assert.equal(result.status, "cancel_requested"); + assert.deepEqual( + commands.map((command) => command.type), + ["thread.stop"], + ); + const command = commands[0]!; + if (command.type !== "thread.stop") return yield* Effect.die("Expected thread stop"); + assert.equal(command.threadId, childThreadId); + assert.deepEqual( + child.runs.filter((run) => run.status === "queued").map((run) => run.queueHeld), + [true, true], + ); + }).pipe(Effect.provide(layer)); +}); + +it.effect("keeps truly executing child work running ahead of held queued work", () => { + const { layer } = fixture({ executing: true }); + return Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const result = yield* service.taskStatus(scope, taskId); + assert.equal(result.status, "running"); + assert.equal(result.hasPendingChildRuns, true); + }).pipe(Effect.provide(layer)); +}); + +it.effect("keeps a published result terminal while stopping later backing-thread work", () => { + const { commands, layer } = fixture({ published: true, executing: true }); + return Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const status = yield* service.taskStatus(scope, taskId); + assert.equal(status.status, "completed"); + assert.equal(status.summary, "Published original result"); + assert.equal(status.hasPendingChildRuns, true); + const cancel = yield* service.cancelTask(scope, { + taskId, + clientRequestId: "cancel-published", + }); + assert.equal(cancel.status, "completed"); + assert.deepEqual( + commands.map((command) => command.type), + ["thread.stop"], + ); + }).pipe(Effect.provide(layer)); +}); + +it.effect("does not dispose completion delivery when stopping the child fails", () => { + const { commands, layer } = fixture({ failCancellation: true, deliveryPending: true }); + return Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const error = yield* service + .cancelTask(scope, { taskId, clientRequestId: "cancel-fails" }) + .pipe(Effect.flip); + assert.equal(error.code, "task_not_cancellable"); + assert.deepEqual( + commands.map((command) => command.type), + ["thread.stop"], + ); + }).pipe(Effect.provide(layer)); +}); diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index b218a76d1ab5..3998aadfb78c 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -698,6 +698,10 @@ function threadDetail( pendingRequestCount: projection.runtimeRequests.filter( (request) => request.status === "pending", ).length, + queuedRunCount: projection.runs.filter((run) => run.status === "queued").length, + heldQueuedRunCount: projection.runs.filter( + (run) => run.status === "queued" && run.queueHeld === true, + ).length, archived: projection.thread.archivedAt !== null, ...threadSettlement(projection.thread), // From the shell, like the list, so read and list agree on snooze state. @@ -1281,9 +1285,11 @@ const make = Effect.gen(function* () { ) : workState === "result_available" ? taskStatusForRun(progress.resultRun ?? childRun) - : taskStatusForRun(childRun) === "queued" - ? "queued" - : "running"; + : progress.pendingRun !== undefined + ? taskStatusForRun(progress.pendingRun) + : taskStatusForRun(childRun) === "queued" + ? "queued" + : "running"; const derivedResult = task.result !== null ? task.result @@ -2069,11 +2075,12 @@ const make = Effect.gen(function* () { const child = yield* loadProjection(current.childThreadId); if ( current.workState !== "waiting_for_children" && - ThreadManagementService.latestActiveRun(child) === undefined + ThreadManagementService.latestActiveRun(child) === undefined && + !child.runs.some((run) => run.status === "queued") ) { return yield* failure( "task_not_cancellable", - `Delegated task ${input.taskId} has no interruptible child run.`, + `Delegated task ${input.taskId} has no interruptible or queued child run.`, ); } yield* stopWithinLimit; diff --git a/apps/server/src/mcp/toolkits/orchestrator/tools.ts b/apps/server/src/mcp/toolkits/orchestrator/tools.ts index ac0ebb63b7dd..68802c92c0c4 100644 --- a/apps/server/src/mcp/toolkits/orchestrator/tools.ts +++ b/apps/server/src/mcp/toolkits/orchestrator/tools.ts @@ -192,7 +192,7 @@ const ThreadListTool = Tool.make("t3_thread_list", { const ThreadReadTool = Tool.make("t3_thread_read", { description: - "Read durable state and a paginated timeline from any T3 thread in this environment. The default messages view returns user messages, assistant messages, and proposed plans; activity returns all summarized timeline items. Reading an untruncated terminal assistant result from this parent thread's direct app-owned child acknowledges that child's automatic completion delivery. Continue with afterPosition=nextPosition. Recover long item text with itemId and textOffset=nextTextOffset until nextTextOffset is null; offsets count UTF-16 code units. The thread also reports its snooze state. To link a thread for the user, write `[title](t3-thread://v1/)` with the threadId exactly as returned, not URL-encoded; T3 Code shows the thread's current title.", + "Read durable state and a paginated timeline from any T3 thread in this environment. The default messages view returns user messages, assistant messages, and proposed plans; activity returns all summarized timeline items. pendingRequestCount counts provider approval/input requests; queuedRunCount and heldQueuedRunCount report unstarted message runs independently of the latest run status. Reading an untruncated terminal assistant result from this parent thread's direct app-owned child acknowledges that child's automatic completion delivery. Continue with afterPosition=nextPosition. Recover long item text with itemId and textOffset=nextTextOffset until nextTextOffset is null; offsets count UTF-16 code units. The thread also reports its snooze state. To link a thread for the user, write `[title](t3-thread://v1/)` with the threadId exactly as returned, not URL-encoded; T3 Code shows the thread's current title.", parameters: OrchestratorMcpThreadReadInput, success: OrchestratorMcpThreadReadResult, failure: OrchestratorMcpFailure, diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 5a41e5fd08e2..719d2bdf1b4d 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -7,6 +7,8 @@ import { EventId, MessageId, type ModelSelection, + type OrchestrationV2AppThread, + TurnItemId, NodeId, type OrchestrationV2Run, ProjectId, @@ -14,6 +16,7 @@ import { ProviderInstanceId, ProviderThreadId, RunId, + RunAttemptId, ThreadId, TurnItemId, } from "@t3tools/contracts"; @@ -37,6 +40,8 @@ import * as WorkspacePaths from "../workspace/WorkspacePaths.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as EventSink from "./EventSink.ts"; import * as Orchestrator from "./Orchestrator.ts"; +import * as ProviderRuntimeRecoveryService from "./ProviderRuntimeRecoveryService.ts"; +import type { ProviderAdapterV2Shape } from "@t3tools/provider-core/server/ProviderAdapter"; import { continueRestartedRun } from "./RestartContinuation.ts"; import * as RuntimeLayer from "./runtimeLayer.ts"; import * as ProviderTurnStartServiceTestkit from "./ProviderTurnStartService.testkit.ts"; @@ -103,7 +108,7 @@ const layerTestProviderInstanceRegistry = Layer.succeed( }, ); -const layerTest = Layer.mergeAll(RuntimeLayer.layer, RuntimeLayer.layerEventSink).pipe( +const layerOrchestrationTest = Layer.mergeAll(RuntimeLayer.layer, RuntimeLayer.layerEventSink).pipe( Layer.provideMerge(RuntimeLayer.layerProjectService), Layer.provide( Layer.mock(WorkspacePaths.WorkspacePaths)({ @@ -140,6 +145,8 @@ const layerTest = Layer.mergeAll(RuntimeLayer.layer, RuntimeLayer.layerEventSink Layer.provide(layerPlatformTest), ); +const layerTest = layerOrchestrationTest.pipe(Layer.provide(SqlitePersistence.layerMemory)); + const seedParentWithTerminalTask = (input: { readonly threadId: ThreadId; readonly projectId: ProjectId; @@ -292,6 +299,356 @@ const seedParentWithTerminalTask = (input: { }); it.layer(layerTest)("delegated completion delivery repairs", (it) => { + it.effect( + "keeps an older held continuation pending after later terminal runs and tracks resumed work", + () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const sink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:held-child-parent"); + const childThreadId = ThreadId.make("thread:held-child"); + const taskId = NodeId.make("node:held-child-task"); + const parentRunId = RunId.make("run:held-child-parent"); + const rootNodeId = NodeId.make("node:held-child-parent-root"); + const projectId = ProjectId.make("project:held-child"); + yield* seedParentWithTerminalTask({ + threadId, + projectId, + runId: parentRunId, + rootNodeId, + taskId, + deliveryState: "disposed", + now, + }); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("command:held-child-create"), + threadId: childThreadId, + projectId, + title: "Held continuation", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "agent", + creationSource: "mcp", + }); + const parent = yield* orchestrator.getThreadProjection(threadId); + const child = yield* orchestrator.getThreadProjection(childThreadId); + const task = { + ...parent.subagents[0]!, + childThreadId, + status: "running" as const, + result: null, + completedAt: null, + completionDelivery: undefined, + }; + const makeRun = ( + ordinal: number, + status: OrchestrationV2Run["status"], + ): OrchestrationV2Run => ({ + ...parent.runs[0]!, + id: RunId.make(`run:held-child:${ordinal}`), + threadId: childThreadId, + ordinal, + providerThreadId: null, + rootNodeId: null, + activeAttemptId: null, + userMessageId: MessageId.make(`message:held-child:${ordinal}`), + status, + startedAt: status === "queued" ? null : now, + completedAt: status === "queued" ? null : now, + queueHeld: status === "queued", + delegatedCompletion: undefined, + }); + const runs = [ + makeRun(1, "failed"), + makeRun(2, "queued"), + makeRun(3, "completed"), + makeRun(4, "completed"), + ]; + const afterSequence = yield* sink.latestSequence(); + yield* sink.write({ + events: [ + { + id: EventId.make("event:held-child-link"), + type: "thread.metadata-updated", + threadId: childThreadId, + occurredAt: now, + payload: { + ...child.thread, + lineage: { + parentThreadId: threadId, + relationshipToParent: "subagent", + rootThreadId: threadId, + }, + forkedFrom: { type: "node", nodeId: taskId }, + }, + }, + { + id: EventId.make("event:held-child-task"), + type: "subagent.updated", + threadId, + occurredAt: now, + payload: task, + }, + { + id: EventId.make("event:held-child-node"), + type: "node.updated", + threadId, + occurredAt: now, + payload: { + id: taskId, + threadId, + runId: parentRunId, + parentNodeId: rootNodeId, + rootNodeId, + kind: "subagent", + status: "running", + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + id: EventId.make("event:held-child-item"), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: { + id: TurnItemId.make("item:held-child"), + threadId, + runId: parentRunId, + nodeId: taskId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + type: "subagent", + subagentId: taskId, + prompt: task.prompt, + origin: "app_owned", + driver, + providerInstanceId: modelSelection.instanceId, + childThreadId, + result: null, + }, + }, + ...runs.map((run) => ({ + id: EventId.make(`event:${run.id}`), + type: "run.updated" as const, + threadId: childThreadId, + runId: run.id, + occurredAt: now, + payload: run, + })), + { + id: EventId.make("event:held-child-final-result"), + type: "message.updated", + threadId: childThreadId, + occurredAt: now, + payload: { + id: MessageId.make("message:held-child-result"), + threadId: childThreadId, + runId: runs[3]!.id, + nodeId: null, + role: "assistant", + createdBy: "agent", + creationSource: "provider", + text: "Latest completed result", + attachments: [], + streaming: false, + createdAt: now, + updatedAt: now, + }, + }, + ], + }); + const awaitTaskStatus = (sequence: number, status: "pending" | "running" | "completed") => + sink.stream({ afterSequence: sequence, eventType: "subagent.updated" }).pipe( + Stream.filter( + (stored) => + stored.event.type === "subagent.updated" && + stored.event.payload.id === taskId && + stored.event.payload.status === status, + ), + Stream.take(1), + Stream.runDrain, + ); + yield* awaitTaskStatus(afterSequence, "pending"); + const pending = yield* orchestrator.getThreadProjection(threadId); + assert.equal(pending.subagents[0]?.status, "pending"); + assert.isNull(pending.subagents[0]?.result); + assert.equal(pending.nodes.find((node) => node.id === taskId)?.status, "pending"); + assert.equal(pending.turnItems.find((item) => item.type === "subagent")?.status, "pending"); + assert.isFalse( + pending.contextTransfers.some((transfer) => transfer.type === "subagent_result"), + ); + const held = yield* orchestrator.getThreadProjection(childThreadId); + assert.equal(held.runs.find((run) => run.id === runs[1]!.id)?.status, "queued"); + assert.isTrue(held.runs.find((run) => run.id === runs[1]!.id)?.queueHeld); + assert.lengthOf(held.runs, 4); + assert.equal( + held.messages.find((message) => message.role === "assistant")?.text, + "Latest completed result", + ); + + // A real resumed run must restore Running even when later ordinals + // already hold completed history. Recovery never performs this resume. + const beforeResume = yield* sink.latestSequence(); + yield* sink.write({ + events: [ + { + id: EventId.make("event:held-child-resumed"), + type: "run.updated", + threadId: childThreadId, + runId: runs[1]!.id, + occurredAt: now, + payload: { + ...runs[1]!, + status: "running", + queueHeld: false, + startedAt: now, + completedAt: null, + }, + }, + ], + }); + yield* awaitTaskStatus(beforeResume, "running"); + const running = yield* orchestrator.getThreadProjection(threadId); + assert.equal(running.subagents[0]?.status, "running"); + assert.equal(running.nodes.find((node) => node.id === taskId)?.status, "running"); + assert.equal(running.turnItems.find((item) => item.type === "subagent")?.status, "running"); + assert.isNull(running.subagents[0]?.result); + const beforeComplete = yield* sink.latestSequence(); + yield* sink.write({ + events: [ + { + id: EventId.make("event:held-child-resume-completed"), + type: "run.updated", + threadId: childThreadId, + runId: runs[1]!.id, + occurredAt: now, + payload: { + ...runs[1]!, + status: "completed", + queueHeld: false, + startedAt: now, + completedAt: now, + }, + }, + ], + }); + yield* awaitTaskStatus(beforeComplete, "completed"); + const completed = yield* orchestrator.getThreadProjection(threadId); + assert.equal(completed.subagents[0]?.result, "Latest completed result"); + assert.lengthOf( + completed.contextTransfers.filter((transfer) => transfer.type === "subagent_result"), + 1, + ); + + const siblingThreadId = ThreadId.make("thread:held-child-live-sibling"); + const siblingTaskId = NodeId.make("node:held-child-live-sibling"); + const beforeFollowup = yield* sink.latestSequence(); + const followup = makeRun(5, "queued"); + yield* sink.write({ + events: [ + { + id: EventId.make("event:held-child-new-followup"), + type: "run.created", + threadId: childThreadId, + runId: followup.id, + occurredAt: now, + payload: followup, + }, + { + id: EventId.make("event:held-child-live-sibling"), + type: "thread.created", + threadId: siblingThreadId, + occurredAt: now, + payload: { + ...child.thread, + id: siblingThreadId, + lineage: { + parentThreadId: threadId, + relationshipToParent: "subagent", + rootThreadId: threadId, + }, + forkedFrom: { type: "node", nodeId: siblingTaskId }, + }, + }, + { + id: EventId.make("event:held-child-live-sibling-task"), + type: "subagent.updated", + threadId, + occurredAt: now, + payload: { + ...task, + id: siblingTaskId, + childThreadId: siblingThreadId, + status: "pending", + }, + }, + { + id: EventId.make("event:held-child-live-sibling-run"), + type: "run.created", + threadId: siblingThreadId, + occurredAt: now, + payload: { + ...makeRun(1, "running"), + id: RunId.make("run:held-child-live-sibling"), + threadId: siblingThreadId, + completedAt: null, + }, + }, + ], + }); + // The same created-run stream processes the follow-up before this live + // sibling's receipt. The published result must remain terminal. + yield* sink.stream({ afterSequence: beforeFollowup, eventType: "subagent.updated" }).pipe( + Stream.filter( + (stored) => + stored.event.type === "subagent.updated" && + stored.event.payload.id === siblingTaskId && + stored.event.payload.status === "running", + ), + Stream.take(1), + Stream.runDrain, + ); + const afterFollowup = yield* orchestrator.getThreadProjection(threadId); + const published = afterFollowup.subagents.find((candidate) => candidate.id === taskId); + assert.equal(published?.status, "completed"); + assert.equal(published?.result, "Latest completed result"); + assert.equal( + afterFollowup.subagents.find((candidate) => candidate.id === siblingTaskId)?.status, + "running", + ); + assert.lengthOf( + afterFollowup.contextTransfers.filter((transfer) => transfer.type === "subagent_result"), + 1, + ); + const childAfterFollowup = yield* orchestrator.getThreadProjection(childThreadId); + assert.equal( + childAfterFollowup.runs.find((run) => run.id === followup.id)?.status, + "queued", + ); + assert.isTrue(childAfterFollowup.runs.find((run) => run.id === followup.id)?.queueHeld); + }), + ); + it.effect("acceptance batches pending siblings without acknowledging their results", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; @@ -1466,3 +1823,265 @@ it.layer(layerTest)("delegated tasks across a server restart", (it) => { }), ); }); + +it.effect.each([ + { + heldQueue: true, + name: "repairs a persisted Running task with held older input during startup without replaying it", + }, + { + heldQueue: false, + name: "publishes the recovered terminal result of an active child without queued work", + }, +])("$name", ({ heldQueue }) => { + const parentId = ThreadId.make("thread:startup-queue-parent"); + const childId = ThreadId.make("thread:startup-queue-child"); + const taskId = NodeId.make("node:startup-queue-task"); + const parentRunId = RunId.make("run:startup-queue-parent"); + const rootNodeId = NodeId.make("node:startup-queue-parent"); + const seed = Effect.gen(function* () { + const sink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const thread = (id: ThreadId): OrchestrationV2AppThread => ({ + id, + projectId: ProjectId.make("project:startup-queue"), + title: "Persisted held queue", + createdBy: "agent", + creationSource: "mcp", + providerInstanceId: modelSelection.instanceId, + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + activeProviderThreadId: null, + lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: parentId }, + forkedFrom: null, + createdAt: now, + updatedAt: now, + archivedAt: null, + settledOverride: null, + settledAt: null, + lastVisitedAt: null, + deletedAt: null, + }); + const run = (ordinal: number, status: OrchestrationV2Run["status"]): OrchestrationV2Run => ({ + id: RunId.make(`run:startup-queue:${ordinal}`), + threadId: childId, + ordinal, + providerInstanceId: modelSelection.instanceId, + modelSelection, + providerThreadId: null, + userMessageId: MessageId.make(`message:startup-queue:${ordinal}`), + rootNodeId: null, + activeAttemptId: null, + status, + requestedAt: now, + startedAt: status === "queued" ? null : now, + completedAt: status === "queued" || status === "running" ? null : now, + queueHeld: status === "queued", + checkpointId: null, + contextHandoffId: null, + }); + yield* sink.write({ + events: [ + { + id: EventId.make("event:startup-queue-parent"), + type: "thread.created", + threadId: parentId, + occurredAt: now, + payload: thread(parentId), + }, + { + id: EventId.make("event:startup-queue-child"), + type: "thread.created", + threadId: childId, + occurredAt: now, + payload: { + ...thread(childId), + lineage: { + parentThreadId: parentId, + relationshipToParent: "subagent", + rootThreadId: parentId, + }, + forkedFrom: { type: "node", nodeId: taskId }, + }, + }, + { + id: EventId.make("event:startup-queue-parent-run"), + type: "run.updated", + threadId: parentId, + runId: parentRunId, + occurredAt: now, + payload: { ...run(1, "completed"), id: parentRunId, threadId: parentId, rootNodeId }, + }, + { + id: EventId.make("event:startup-queue-node"), + type: "node.updated", + threadId: parentId, + occurredAt: now, + payload: { + id: taskId, + threadId: parentId, + runId: parentRunId, + parentNodeId: rootNodeId, + rootNodeId, + kind: "subagent", + status: "running", + countsForRun: false, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + runtimeRequestId: null, + checkpointScopeId: null, + startedAt: now, + completedAt: null, + }, + }, + { + id: EventId.make("event:startup-queue-item"), + type: "turn-item.updated", + threadId: parentId, + occurredAt: now, + payload: { + id: TurnItemId.make("item:startup-queue-task"), + threadId: parentId, + runId: parentRunId, + nodeId: taskId, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + type: "subagent", + subagentId: taskId, + origin: "app_owned", + driver, + providerInstanceId: modelSelection.instanceId, + childThreadId: childId, + prompt: "Held continuation", + result: null, + }, + }, + { + id: EventId.make("event:startup-queue-task"), + type: "subagent.updated", + threadId: parentId, + occurredAt: now, + payload: { + id: taskId, + threadId: parentId, + runId: parentRunId, + parentNodeId: NodeId.make("node:startup-queue-parent"), + origin: "app_owned", + createdBy: "agent", + driver, + providerInstanceId: modelSelection.instanceId, + providerThreadId: null, + childThreadId: childId, + nativeTaskRef: null, + prompt: "Held continuation", + title: null, + model: null, + status: "running", + result: null, + startedAt: now, + completedAt: null, + updatedAt: now, + }, + }, + ...(heldQueue + ? [run(1, "failed"), run(2, "queued"), run(3, "completed"), run(4, "completed")] + : [run(1, "failed"), run(2, "running")] + ).map((payload) => ({ + id: EventId.make(`event:${payload.id}`), + type: "run.updated" as const, + threadId: childId, + runId: payload.id, + occurredAt: now, + payload, + })), + { + id: EventId.make("event:startup-queue-message"), + type: "message.updated", + threadId: childId, + occurredAt: now, + payload: { + id: MessageId.make("message:startup-queue:2"), + threadId: childId, + runId: RunId.make("run:startup-queue:2"), + nodeId: null, + role: "user", + createdBy: "agent", + creationSource: "mcp", + text: "Keep this genuine queued continuation", + attachments: [], + streaming: false, + createdAt: now, + updatedAt: now, + }, + }, + ], + }); + }); + const seededPersistence = Layer.effectDiscard(seed).pipe( + Layer.provideMerge( + RuntimeLayer.layerEventSink.pipe(Layer.provideMerge(SqlitePersistence.layerMemory)), + ), + ); + return Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const beforeRuntimeRecovery = yield* orchestrator.getThreadProjection(parentId); + assert.equal(beforeRuntimeRecovery.subagents[0]?.status, "running"); + const runtimeRecovery = yield* ProviderRuntimeRecoveryService.ProviderRuntimeRecoveryService; + yield* runtimeRecovery.recover; + yield* orchestrator.recoverDelegatedTasks; + const parent = yield* orchestrator.getThreadProjection(parentId); + const child = yield* orchestrator.getThreadProjection(childId); + if (heldQueue) { + assert.equal(parent.subagents[0]?.status, "pending"); + assert.equal(parent.nodes.find((node) => node.id === taskId)?.status, "pending"); + assert.equal(parent.turnItems.find((item) => item.type === "subagent")?.status, "pending"); + assert.isNull(parent.subagents[0]?.result); + assert.isFalse( + parent.contextTransfers.some((transfer) => transfer.type === "subagent_result"), + ); + const queued = child.runs.find((candidate) => candidate.ordinal === 2); + assert.equal(queued?.status, "queued"); + assert.isTrue(queued?.queueHeld); + assert.isNull(queued?.activeAttemptId); + assert.isNull(queued?.providerThreadId); + assert.equal( + child.messages.find((message) => message.runId === queued?.id)?.text, + "Keep this genuine queued continuation", + ); + assert.lengthOf(child.runs, 4); + assert.isFalse( + child.runs.some((candidate) => + ["preparing", "starting", "running"].includes(candidate.status), + ), + ); + } else { + assert.equal(parent.subagents[0]?.status, "cancelled"); + assert.equal(parent.nodes.find((node) => node.id === taskId)?.status, "cancelled"); + assert.equal(parent.turnItems.find((item) => item.type === "subagent")?.status, "cancelled"); + assert.equal(parent.subagents[0]?.result, "Child task ended with status cancelled."); + assert.lengthOf( + parent.contextTransfers.filter((transfer) => transfer.type === "subagent_result"), + 1, + ); + assert.equal(child.runs.find((candidate) => candidate.ordinal === 2)?.status, "cancelled"); + assert.lengthOf(child.runs, 2); + assert.isFalse( + child.runs.some((candidate) => + ["preparing", "starting", "running", "queued"].includes(candidate.status), + ), + ); + } + }).pipe(Effect.provide(layerOrchestrationTest.pipe(Layer.provide(seededPersistence)))); +}); diff --git a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts index 9a1e43a8b44b..fd0290afbdc9 100644 --- a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts +++ b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts @@ -2779,6 +2779,7 @@ it.layer(layerTest)("orchestration V2 foundation persistence", (it) => { { id: sessionId, driver: "codex", providerInstanceId, status: "running" }, ], providerTurns: [{ providerThreadId, runAttemptId: attemptId, status: "running" }], + subagents: [], } as unknown as OrchestrationV2ThreadProjection; const recovery = yield* ProviderRuntimeRecovery.make.pipe( Effect.provide( diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 977b225f5256..bef44d78b954 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -9572,8 +9572,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); /** - * Transfers a terminal child's result into its parent and offers the parent - * wake. Every mutation here targets the PARENT thread, so callers must hold + * Synchronizes unpublished child work, then transfers a settled child's + * result into its parent and offers the parent wake. Every mutation here targets the PARENT thread, so callers must hold * the parent thread's dispatch lock rather than the child's: the * delegated_task.wake-policy handler rewrites the same subagent row under * that lock with a full-row payload, and unserialized writers clobber each @@ -9600,38 +9600,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return; } const progress = delegatedTaskProgress(childControls); - if (progress.state !== "result_available") return; - const childRun = progress.resultRun; - if (childRun === undefined) return; - const terminalStatus = delegatedTaskTerminalStatus(childRun.status); - if (terminalStatus === null) { - return; - } - if ( - yield* childAwaitsRestartContinuation( - childControls.runs, - childRun, - options?.settledContinuationOf, - ) - ) { - return; - } - - const childResult = yield* projectionStore.getThreadRecords( - childThreadId, - ["messages", "turnItems"], - { - messageRoles: ["assistant"], - messageRunIds: [childRun.id], - turnItemRunId: childRun.id, - turnItemTypes: ["assistant_message", "error"], - }, - ); - const childProjection = { - ...childControls, - messages: childResult.messages, - turnItems: childResult.turnItems, - }; const parentThreadId = childControls.thread.lineage.parentThreadId; const parentProjection = yield* projectionStore.getThreadRecords( parentThreadId, @@ -9666,16 +9634,105 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return; } + // A held queued continuation still belongs to this task. Publish no + // result until it settles, but keep the parent's durable status honest. + if (task.result !== null) return; + const parentNode = parentProjection.nodes.find((candidate) => candidate.id === task.id); + const parentTurnItem = parentProjection.turnItems.find( + (candidate) => candidate.type === "subagent" && candidate.subagentId === task.id, + ); + if (progress.state !== "result_available") { + const status = + progress.state === "waiting_for_children" + ? ("waiting" as const) + : progress.pendingRun === undefined || progress.pendingRun.status === "queued" + ? ("pending" as const) + : ("running" as const); + if ( + task.status === status && + (parentNode === undefined || parentNode.status === status) && + (parentTurnItem === undefined || parentTurnItem.status === status) + ) { + return; + } + const now = yield* DateTime.now; + yield* writeSystemEvents([ + { + type: "subagent.updated", + threadId: parentThreadId, + ...(task.runId === null ? {} : { runId: task.runId }), + nodeId: task.id, + driver: task.driver, + occurredAt: now, + payload: { ...task, status, completedAt: null, updatedAt: now }, + }, + ...(parentNode === undefined + ? [] + : [ + { + type: "node.updated" as const, + threadId: parentThreadId, + ...(parentNode.runId === null ? {} : { runId: parentNode.runId }), + nodeId: parentNode.id, + driver: task.driver, + occurredAt: now, + payload: { ...parentNode, status, completedAt: null }, + }, + ]), + ...(parentTurnItem === undefined + ? [] + : [ + { + type: "turn-item.updated" as const, + threadId: parentThreadId, + ...(parentTurnItem.runId === null ? {} : { runId: parentTurnItem.runId }), + ...(parentTurnItem.nodeId === null ? {} : { nodeId: parentTurnItem.nodeId }), + driver: task.driver, + occurredAt: now, + payload: { ...parentTurnItem, status, completedAt: null, updatedAt: now }, + }, + ]), + ]); + return; + } + const childRun = progress.resultRun; + if (childRun === undefined) return; + const terminalStatus = delegatedTaskTerminalStatus(childRun.status); + if (terminalStatus === null) { + return; + } + + if ( + yield* childAwaitsRestartContinuation( + childControls.runs, + childRun, + options?.settledContinuationOf, + ) + ) { + return; + } + + const childResult = yield* projectionStore.getThreadRecords( + childThreadId, + ["messages", "turnItems"], + { + messageRoles: ["assistant"], + messageRunIds: [childRun.id], + turnItemRunId: childRun.id, + turnItemTypes: ["assistant_message", "error"], + }, + ); + const childProjection = { + ...childControls, + messages: childResult.messages, + turnItems: childResult.turnItems, + }; const now = yield* DateTime.now; const result = subagentResultForRun(childProjection, childRun); const parentRun = task.runId === null ? undefined : parentProjection.runs.find((candidate) => candidate.id === task.runId); - const parentNode = parentProjection.nodes.find((candidate) => candidate.id === task.id); - const parentTurnItem = parentProjection.turnItems.find( - (candidate) => candidate.type === "subagent" && candidate.subagentId === task.id, - ); const updatedTask: OrchestrationV2Subagent = { ...task, providerThreadId: childRun.providerThreadId, @@ -10629,7 +10686,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const dispatchWithReceipt = (command: OrchestrationV2ServerCommand) => threadDispatch.withLock(commandThreadId(command), dispatchWithReceiptEffect(command)); - const handleTerminalRun = (stored: OrchestrationV2StoredEvent) => + const handleRunUpdate = (stored: OrchestrationV2StoredEvent) => Effect.gen(function* () { const threadId = stored.event.threadId; // finalize writes the parent thread and startNextQueuedRun writes this @@ -10638,6 +10695,20 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio // while holding the parent lock, so nesting the parent lock inside the // child lock here would invert that order, and the keyed executor's // semaphores are neither reentrant nor deadlock-aware. + // Nonterminal updates synchronize the task without promoting queued work. + if ( + String(stored.commandId).startsWith("command:runtime-reconcile:") || + stored.event.type !== "run.updated" || + !["completed", "interrupted", "failed", "cancelled", "rolled_back"].includes( + stored.event.payload.status, + ) + ) { + const parentThreadId = yield* appOwnedSubagentParentThreadId(threadId); + if (parentThreadId !== undefined) { + yield* threadDispatch.withLock(parentThreadId, finalizeAppOwnedSubagent(threadId)); + } + return; + } if (stored.event.type === "run.updated") { yield* threadDispatch.withLock( threadId, @@ -10667,7 +10738,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } }).pipe( Effect.catchCause((cause) => - Effect.logWarning("Failed to react to terminal V2 run", { + Effect.logWarning("Failed to react to V2 run update", { threadId: stored.event.threadId, sequence: stored.sequence, cause, @@ -10675,28 +10746,24 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ); - // Historical terminal events are already represented by the projections - // below. Replaying the full event table on every server start delays live - // queue promotion in proportion to the lifetime size of the database. - const terminalEventsAfterSequence = yield* eventSink.latestSequence().pipe(Effect.orDie); + // Recover historical child state from projections below. Live run updates + // synchronize unpublished tasks; only terminal updates promote the queue. + // Replaying full history would delay this work as the database grows. + const runEventsAfterSequence = yield* eventSink.latestSequence().pipe(Effect.orDie); // Queue promotion can wait on a provider or a thread lock. Subscribe to run // updates before buffering so that wait never retains unrelated tool bodies. - yield* eventSink - .stream({ afterSequence: terminalEventsAfterSequence, eventType: "run.updated" }) - .pipe( - Stream.filter( - (stored) => - stored.event.type === "run.updated" && - !String(stored.commandId).startsWith("command:runtime-reconcile:") && - (stored.event.payload.status === "completed" || - stored.event.payload.status === "interrupted" || - stored.event.payload.status === "failed" || - stored.event.payload.status === "cancelled" || - stored.event.payload.status === "rolled_back"), - ), - Stream.runForEach(handleTerminalRun), - Effect.forkDetach, - ); + yield* Stream.merge( + eventSink.stream({ afterSequence: runEventsAfterSequence, eventType: "run.updated" }), + eventSink.stream({ afterSequence: runEventsAfterSequence, eventType: "run.created" }), + ).pipe( + Stream.filter( + (stored) => + !String(stored.commandId).startsWith("command:runtime-reconcile:") && + (stored.event.type === "run.updated" || stored.event.type === "run.created"), + ), + Stream.runForEach(handleRunUpdate), + Effect.forkDetach, + ); // Settles child results and completion deliveries whose runs ended without // the listener above: before this boot, or in runtime reconciliation, which diff --git a/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts b/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts index c7ea0530d351..5c7046e93e83 100644 --- a/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts @@ -224,71 +224,70 @@ it.effect("selects unfinished recovery work without reading settled thread histo }).pipe(Effect.provide(layerTest)), ); -it.effect("recovers terminal subagent results until their cross-thread transfer exists", () => - Effect.gen(function* () { - const projections = yield* ProjectionStore.ProjectionStoreV2; - const sql = yield* SqlClient.SqlClient; - const now = yield* DateTime.now; - const parent = yield* createThread("subagent-parent"); - const children: Array = []; - for (const name of ["terminal", "archived", "deleted", "running", "held", "queued"]) { - const child = yield* createThread(`subagent-${name}`, { - lineage: { parentThreadId: parent, relationshipToParent: "subagent", rootThreadId: parent }, - forkedFrom: { type: "node", nodeId: NodeId.make(`node:${name}`) }, - archivedAt: name === "archived" ? now : null, - deletedAt: name === "deleted" ? now : null, - }); - yield* createRun(child, name === "running" ? "running" : "completed"); - // A wake queued behind the result: held by Stop or a restart, or still deliverable. - if (name === "held" || name === "queued") { - yield* createRun(child, "queued", { - ordinal: 2, - startedAt: null, - ...(name === "held" ? { queueHeld: true } : {}), +it.effect( + "recovers unpublished subagent lifecycle states until their cross-thread transfer exists", + () => + Effect.gen(function* () { + const projections = yield* ProjectionStore.ProjectionStoreV2; + const sql = yield* SqlClient.SqlClient; + const now = yield* DateTime.now; + const parent = yield* createThread("subagent-parent"); + const children: Array = []; + for (const name of ["terminal", "archived", "deleted", "running", "queued"]) { + const child = yield* createThread(`subagent-${name}`, { + lineage: { + parentThreadId: parent, + relationshipToParent: "subagent", + rootThreadId: parent, + }, + forkedFrom: { type: "node", nodeId: NodeId.make(`node:${name}`) }, + archivedAt: name === "archived" ? now : null, + deletedAt: name === "deleted" ? now : null, }); + yield* createRun( + child, + name === "running" ? "running" : name === "queued" ? "queued" : "completed", + ); + children.push(child); } - children.push(child); - } - const terminalChildren = new Set([children[0]!, children[1]!, children[4]!]); - assert.deepEqual( - new Set(yield* projections.getRecoveryThreadIds("subagent-results")), - terminalChildren, - ); - - const completed = children[0]!; - const transferId = ContextTransferId.make("transfer:recovery:subagent-result"); - yield* projections.apply({ - id: EventId.make("event:recovery:subagent-result"), - type: "context-transfer.created", - threadId: parent, - occurredAt: now, - payload: { - id: transferId, - type: "subagent_result", - sourceThreadId: completed, - targetThreadId: parent, - sourcePoint: { threadId: completed }, - basePoint: null, - sourceProviderInstanceId: providerInstanceId, - targetProviderInstanceId: providerInstanceId, - targetRunId: null, - status: "pending", - resolution: null, - createdBy: "system", - error: null, - createdAt: now, - updatedAt: now, - consumedAt: null, - }, - }); - const archivedChild = children[1]; - assert.isDefined(archivedChild); - assert.deepEqual( - new Set(yield* projections.getRecoveryThreadIds("subagent-results")), - new Set([archivedChild, children[4]!]), - ); - assert.deepEqual(yield* projections.getUnreadableThreadIds(), []); - yield* sql` + assert.deepEqual( + new Set(yield* projections.getRecoveryThreadIds("subagent-results")), + new Set(children.filter((_, index) => index !== 2)), + ); + const completed = children[0]!; + const transferId = ContextTransferId.make("transfer:recovery:subagent-result"); + yield* projections.apply({ + id: EventId.make("event:recovery:subagent-result"), + type: "context-transfer.created", + threadId: parent, + occurredAt: now, + payload: { + id: transferId, + type: "subagent_result", + sourceThreadId: completed, + targetThreadId: parent, + sourcePoint: { threadId: completed }, + basePoint: null, + sourceProviderInstanceId: providerInstanceId, + targetProviderInstanceId: providerInstanceId, + targetRunId: null, + status: "pending", + resolution: null, + createdBy: "system", + error: null, + createdAt: now, + updatedAt: now, + consumedAt: null, + }, + }); + const archivedChild = children[1]; + assert.isDefined(archivedChild); + assert.deepEqual( + new Set(yield* projections.getRecoveryThreadIds("subagent-results")), + new Set([archivedChild, children[3], children[4]]), + ); + assert.deepEqual(yield* projections.getUnreadableThreadIds(), []); + yield* sql` UPDATE orchestration_v2_projection_context_transfers SET payload_json = '{}' WHERE context_transfer_id = ${transferId} `; diff --git a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts index f626d0c38583..aa6a8c4fadc1 100644 --- a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts @@ -143,6 +143,37 @@ const createItem = Effect.fn(function* ( it.effect.each([ ["sql", layerSql], ["memory", ProjectionStore.layerMemory], +] as const)( + "%s: excludes abandoned foreground tools from shells while retaining their history for recovery", + ([, testLayer]) => + Effect.gen(function* () { + const store = yield* ProjectionStore.ProjectionStoreV2; + for (const status of ["failed", "interrupted", "cancelled"] as const) { + const threadId = yield* createThread(`orphan-handoff-${status}`); + const failedRunId = yield* createRun(threadId, status); + yield* createItem(threadId, failedRunId, "running"); + yield* createRun(threadId, "completed", 2); + yield* createRun(threadId, "completed", 3); + const allShells = yield* store.getShellSnapshot(); + const scopedShell = yield* store.getThreadShell(threadId); + const candidates = yield* store.getSettlementCandidates(threadId); + assert.deepEqual( + allShells.threads.find((thread) => thread.id === threadId)?.pendingBackgroundTasks, + [], + ); + assert.deepEqual(scopedShell?.pendingBackgroundTasks, []); + assert.deepEqual(candidates[0]?.pendingBackgroundTasks, []); + const history = yield* store.getThreadRecords(threadId, ["runs", "turnItems"]); + assert.equal(history.turnItems[0]?.status, "running"); + assert.equal(history.runs.find((run) => run.id === failedRunId)?.status, status); + assert.include(yield* store.getRecoveryThreadIds("runtime"), threadId); + } + }).pipe(Effect.provide(testLayer)), +); + +it.effect.each([ + ["sql", SqlLayer], + ["memory", ProjectionStore.layerMemory], ] as const)( "%s: discovers settlement work with the same activity and background semantics as the shell", ([, testLayer]) => diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index 8c33faef808b..cab22deaa8a3 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -556,9 +556,7 @@ function needsRecovery( projection.thread.lineage.relationshipToParent === "subagent" && parentThreadId !== null && projection.thread.forkedFrom?.type === "node" && - ["completed", "interrupted", "failed", "cancelled", "rolled_back"].includes( - latestUnheldRun(projection.runs)?.status ?? "idle", - ) && + projection.runs.length > 0 && !projection.contextTransfers.some( (transfer) => transfer.type === "subagent_result" && @@ -3576,17 +3574,10 @@ export const layer: Layer.Layer = json_extract(child.payload_json, '$.lineage.relationshipToParent') = 'subagent' AND json_extract(child.payload_json, '$.lineage.parentThreadId') IS NOT NULL AND json_extract(child.payload_json, '$.forkedFrom.type') = 'node' - -- A held queue waits for the user, so the newest unheld run - -- decides, matching latestUnheldRun. - AND ( - SELECT status FROM orchestration_v2_projection_runs + AND EXISTS ( + SELECT 1 FROM orchestration_v2_projection_runs WHERE thread_id = child.thread_id - AND NOT ( - status = 'queued' - AND json_extract(payload_json, '$.queueHeld') IS 1 - ) - ORDER BY ordinal DESC LIMIT 1 - ) IN ('completed', 'interrupted', 'failed', 'cancelled', 'rolled_back') + ) AND NOT EXISTS ( SELECT 1 FROM orchestration_v2_projection_context_transfers WHERE source_thread_id = child.thread_id @@ -5353,10 +5344,11 @@ export const layer: Layer.Layer = ON r.run_id = i.run_id WHERE i.type IN ('command_execution', 'dynamic_tool', 'subagent') AND i.status NOT IN ('completed', 'interrupted', 'failed', 'cancelled') - -- A rolled-back run's items are abandoned, not pending. Without - -- this the shell reports Waiting for work no one will finish, - -- matching the item_count query's exclusion above. + -- Foreground tools on failed/stopped runs cannot finish. Commands + -- and native children can outlive their root; keep those visible. AND (i.run_id IS NULL OR r.status <> 'rolled_back') + AND (i.type <> 'dynamic_tool' OR r.status IS NULL + OR r.status NOT IN ('failed', 'interrupted', 'cancelled')) ` : sql` SELECT i.thread_id, i.payload_json @@ -5366,6 +5358,8 @@ export const layer: Layer.Layer = WHERE i.type IN ('command_execution', 'dynamic_tool', 'subagent') AND i.status NOT IN ('completed', 'interrupted', 'failed', 'cancelled') AND (i.run_id IS NULL OR r.status <> 'rolled_back') + AND (i.type <> 'dynamic_tool' OR r.status IS NULL + OR r.status NOT IN ('failed', 'interrupted', 'cancelled')) AND i.thread_id IN ${sql.in(threadIds)} `; diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts index da764d41d730..c251126b6347 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.test.ts @@ -147,6 +147,7 @@ it.effect("expires orphaned runtime requests before command readiness", () => { providerThreads: [], runs: [], nodes: [], + subagents: [], } as unknown as OrchestrationV2ThreadProjection; const layer = ProviderRuntimeRecovery.layer.pipe( Layer.provide(ServerSettings.layerTest()), @@ -198,6 +199,7 @@ it.effect("preserves async questions across startup and shutdown", () => { runs: [], nodes: [], turnItems: [], + subagents: [], } as unknown as OrchestrationV2ThreadProjection; const commitCommand = vi.fn(() => Effect.die("an async question needs no process-loss write")); const layer = ProviderRuntimeRecovery.layer.pipe( @@ -247,6 +249,7 @@ it.effect("uses the same reconciliation path to cancel runtime requests during s providerThreads: [], runs: [], nodes: [], + subagents: [], } as unknown as OrchestrationV2ThreadProjection; const layer = ProviderRuntimeRecovery.layer.pipe( Layer.provide(ServerSettings.layerTest()), @@ -1064,9 +1067,9 @@ it.effect( }, ); -it.effect( - "terminalizes a leftover nonpersistent dynamic_tool on a settled run after process loss", - () => { +it.effect.each(["completed", "failed", "interrupted", "cancelled"] as const)( + "terminalizes a leftover nonpersistent dynamic_tool on a %s run after process loss", + (runStatus) => { const threadId = ThreadId.make("thread_recovery_orphan_wait"); const settledRunId = RunId.make("run_recovery_orphan_wait_settled"); const providerThreadId = ProviderThreadId.make("provider_thread_recovery_orphan_wait"); @@ -1095,7 +1098,7 @@ it.effect( }, ], providerTurns: [], - runs: [{ id: settledRunId, status: "completed", providerInstanceId: codexInstanceId }], + runs: [{ id: settledRunId, status: runStatus, providerInstanceId: codexInstanceId }], attempts: [], nodes: [ { id: orphanWaitNodeId, runId: settledRunId, status: "running", kind: "tool_call" }, @@ -1189,9 +1192,9 @@ it.effect( }, ); -it.effect( - "terminalizes the linked subagent and node for a stale subagent item on a settled run", - () => { +it.effect.each(["completed", "running", "waiting"] as const)( + "cancels native tasks on a %s parent run while preserving app-owned task ownership", + (runStatus) => { const threadId = ThreadId.make("thread_recovery_subagent"); const settledRunId = RunId.make("run_recovery_subagent_settled"); const providerThreadId = ProviderThreadId.make("provider_thread_recovery_subagent"); @@ -1200,6 +1203,18 @@ it.effect( const staleItemId = TurnItemId.make("turn_item_recovery_subagent_stale"); const doneItemId = TurnItemId.make("turn_item_recovery_subagent_done"); const claudeInstanceId = ProviderInstanceId.make("claude"); + const appOwnedTasks = (["pending", "running", "waiting", "completed"] as const).map( + (status) => ({ + id: NodeId.make(`node_recovery_app_owned_${status}`), + runId: settledRunId, + driver: ProviderDriverKind.make("claude"), + providerInstanceId: claudeInstanceId, + origin: "app_owned" as const, + childThreadId: ThreadId.make(`thread_recovery_app_owned_${status}`), + status, + result: status === "completed" ? "Published child result" : null, + }), + ); let committedInput: Parameters[0] | null = null; const projection = { @@ -1216,12 +1231,12 @@ it.effect( }, ], providerTurns: [], - // Settled run: the stale-item loop owns it, not the nonterminal loop. - runs: [{ id: settledRunId, status: "completed", providerInstanceId: claudeInstanceId }], + runs: [{ id: settledRunId, status: runStatus, providerInstanceId: claudeInstanceId }], attempts: [], nodes: [ { id: staleSubagentNodeId, runId: settledRunId, status: "running" }, { id: doneSubagentNodeId, runId: settledRunId, status: "completed" }, + ...appOwnedTasks.map((task) => ({ id: task.id, runId: settledRunId, status: task.status })), ], subagents: [ { @@ -1240,6 +1255,7 @@ it.effect( status: "completed", result: "done", }, + ...appOwnedTasks, ], messages: [], turnItems: [ @@ -1253,6 +1269,16 @@ it.effect( subagentId: staleSubagentNodeId, providerInstanceId: claudeInstanceId, }, + ...appOwnedTasks.map((task) => ({ + id: TurnItemId.make(`turn_item_${task.id}`), + runId: settledRunId, + nodeId: task.id, + providerThreadId, + type: "subagent", + status: task.status === "completed" ? "running" : task.status, + subagentId: task.id, + providerInstanceId: claudeInstanceId, + })), { id: doneItemId, runId: settledRunId, @@ -1295,6 +1321,20 @@ it.effect( yield* (yield* ProviderRuntimeRecovery.ProviderRuntimeRecoveryService).reconcile("startup"); const events = committedInput?.events ?? []; + // App-owned children recover from their own runs, including a held queue. + // Neither parent process loss nor a stale item may rewrite their results. + for (const task of appOwnedTasks) { + assert.isFalse( + events.some( + (event) => + (event.type === "subagent.updated" && event.payload.id === task.id) || + (event.type === "node.updated" && event.payload.id === task.id) || + (event.type === "turn-item.updated" && event.payload.nodeId === task.id), + ), + ); + } + assert.equal(appOwnedTasks[3]?.result, "Published child result"); + // Only the nonterminal subagent item is cancelled. const turnItemCancels = events.filter( (event) => event.type === "turn-item.updated" && event.payload.status === "cancelled", diff --git a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts index 7a063bec065b..2ff86b19e53c 100644 --- a/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts +++ b/apps/server/src/orchestration-v2/ProviderRuntimeRecoveryService.ts @@ -3,6 +3,7 @@ import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import { CommandId, type OrchestrationV2DomainEvent, + type NodeId, type ProviderThreadId, type OrchestrationV2RestartCancelledBackgroundWork, type OrchestrationV2Subagent, @@ -113,6 +114,19 @@ function isAppOwnedDelegationItem( return item.type === "subagent" && isAppOwnedDelegation(item); } +function appOwnedSubagentIds(projection: ProjectionStore.ProjectionRuntimeRecoveryState) { + return new Set( + projection.subagents.filter((task) => task.origin === "app_owned").map((task) => task.id), + ); +} + +function isAppOwnedSubagentItem( + item: OrchestrationV2ThreadProjection["turnItems"][number], + taskIds: ReadonlySet, +): boolean { + return item.type === "subagent" && taskIds.has(item.subagentId); +} + function providerThreadHasPendingBackgroundTasks( providerThread: OrchestrationV2ThreadProjection["providerThreads"][number], ): boolean { @@ -192,6 +206,9 @@ export const make = Effect.gen(function* () { continueAfterRestart: boolean, ) { const now = yield* DateTime.now; + // App-owned children have their own durable runs and recovery. Losing + // the parent provider session does not cancel the delegated task. + const taskIds = appOwnedSubagentIds(projection); const runs = [] as Array; for (const run of nonterminalRuns(projection)) { if (run.status === "waiting") { @@ -272,6 +289,7 @@ export const make = Effect.gen(function* () { const recordCancelledBackgroundItem = ( item: OrchestrationV2ThreadProjection["turnItems"][number], ) => { + if (isAppOwnedSubagentItem(item, taskIds)) return; if (!isBackgroundCapableTurnItemType(item.type)) return; const work = cancelledTurnItemWork(item); if (work === undefined) return; @@ -345,6 +363,7 @@ export const make = Effect.gen(function* () { for (const node of projection.nodes.filter( (candidate) => candidate.runId === run.id && + !taskIds.has(candidate.id) && !messageRequestNodeIds.has(candidate.id) && !delegatedTaskNodeIds.has(candidate.id) && (candidate.status === "pending" || @@ -418,6 +437,7 @@ export const make = Effect.gen(function* () { for (const item of projection.turnItems.filter( (candidate) => candidate.runId === run.id && + !isAppOwnedSubagentItem(candidate, taskIds) && (candidate.nodeId === null || !messageRequestNodeIds.has(candidate.nodeId)) && !isAppOwnedDelegationItem(candidate) && (candidate.status === "pending" || @@ -446,6 +466,7 @@ export const make = Effect.gen(function* () { const recoveredNonterminalRunIds = new Set(runs.map((run) => run.id)); const cancelledStaleNodeIds = new Set(); for (const item of projection.turnItems ?? []) { + if (isAppOwnedSubagentItem(item, taskIds)) continue; if (item.runId !== null && recoveredNonterminalRunIds.has(item.runId)) { continue; } @@ -557,6 +578,7 @@ export const make = Effect.gen(function* () { for (const item of projection.turnItems) { if ( item.nodeId !== node.id || + isAppOwnedSubagentItem(item, taskIds) || item.runId !== null || !isNonterminalTurnItemStatus(item.status) || cancelledStaleItemIds.has(item.id) diff --git a/packages/contracts/src/orchestratorMcp.ts b/packages/contracts/src/orchestratorMcp.ts index 0ab3e0b9b86c..69e2ab3831c7 100644 --- a/packages/contracts/src/orchestratorMcp.ts +++ b/packages/contracts/src/orchestratorMcp.ts @@ -371,6 +371,8 @@ export const OrchestratorMcpThreadDetail = Schema.Struct({ runCount: NonNegativeInt, itemCount: NonNegativeInt, pendingRequestCount: NonNegativeInt, + queuedRunCount: Schema.optionalKey(NonNegativeInt), + heldQueuedRunCount: Schema.optionalKey(NonNegativeInt), archived: Schema.Boolean, settled: Schema.Boolean, settledAt: Schema.NullOr(IsoDateTime), diff --git a/packages/provider-core/src/server/subagentProjection.ts b/packages/provider-core/src/server/subagentProjection.ts index 7ee8d7c1aa8d..ff2a81290d6c 100644 --- a/packages/provider-core/src/server/subagentProjection.ts +++ b/packages/provider-core/src/server/subagentProjection.ts @@ -229,11 +229,12 @@ export function delegatedTaskProgress(projection: { const workRuns = projection.runs.filter( (run) => !monitorRuns.has(run.id) && run.status !== "rolled_back", ); - // A held queue waits for the user to resume it (after Stop, a restart, or a - // provider failure), so its runs are not work the task still owes. - const active = workRuns.some( - (run) => !terminal(run.status) && !(run.status === "queued" && run.queueHeld === true), - ); + // A held queue entry can predate many finished follow-ups. It is still + // pending intent, but it must not make an idle provider look running. + const pendingRuns = workRuns + .filter((run) => !terminal(run.status) && !(run.status === "queued" && run.queueHeld === true)) + .toSorted((a, b) => b.ordinal - a.ordinal); + const pendingRun = pendingRuns.find((run) => run.status !== "queued") ?? pendingRuns[0]; const children = projection.subagents.some( (task) => @@ -249,11 +250,12 @@ export function delegatedTaskProgress(projection: { .toSorted((a, b) => (runRanAfter(a, b) ? -1 : runRanAfter(b, a) ? 1 : 0))[0]; return { state: - active || resultRun === undefined + pendingRun !== undefined || resultRun === undefined ? ("working" as const) : children ? ("waiting_for_children" as const) : ("result_available" as const), resultRun, + pendingRun, }; } diff --git a/packages/shared/src/orchestrationV2PendingBackgroundWork.test.ts b/packages/shared/src/orchestrationV2PendingBackgroundWork.test.ts index 013a49f0c125..5b9668364aa4 100644 --- a/packages/shared/src/orchestrationV2PendingBackgroundWork.test.ts +++ b/packages/shared/src/orchestrationV2PendingBackgroundWork.test.ts @@ -3,6 +3,7 @@ import type { OrchestrationV2PendingBackgroundTask } from "@t3tools/contracts"; import { backgroundWorkHoldsCompletion, derivePendingBackgroundWork, + pendingBackgroundTurnItems, turnItemUpdateCanEndBackgroundWork, } from "./orchestrationV2PendingBackgroundWork.ts"; @@ -41,6 +42,72 @@ describe("backgroundWorkHoldsCompletion", () => { }); describe("derivePendingBackgroundWork", () => { + it.each(["failed", "interrupted", "cancelled"] as const)( + "excludes a %s run's orphan handoff while preserving live background work and queued history", + (status) => { + const runs = [ + { id: "run-1" as never, ordinal: 1, status }, + { id: "run-2" as never, ordinal: 2, status: "queued" as const, queueHeld: true }, + { id: "run-3" as never, ordinal: 3, status: "completed" as const }, + { id: "run-4" as never, ordinal: 4, status: "completed" as const }, + ]; + const turnItems = [ + { + id: "handoff", + type: "dynamic_tool" as const, + status: "running" as const, + runId: "run-1", + title: "mcp__t3-code__t3_worktree_handoff", + }, + { + id: "completed-run-background", + type: "dynamic_tool" as const, + status: "running" as const, + runId: "run-3", + title: "Await background result", + }, + { + id: "child", + type: "subagent" as const, + status: "running" as const, + runId: "run-1", + title: "Independent child", + }, + { + id: "command", + type: "command_execution" as const, + status: "running" as const, + runId: "run-1", + title: "Dev server", + }, + ]; + const tasks = derivePendingBackgroundWork({ + latestRun: runs[3], + runs, + turnItems, + providerThreads: [ + { + id: "provider" as never, + pendingBackgroundTasks: [{ taskId: "live-roster", kind: "background_task" as const }], + }, + ], + }); + expect(tasks.map((task) => task.taskId)).toEqual([ + "live-roster", + "completed-run-background", + "child", + "command", + ]); + expect(pendingBackgroundTurnItems({ turnItems, runs }).map((item) => item.id)).toEqual([ + "completed-run-background", + "child", + "command", + ]); + expect(runs[1]).toMatchObject({ status: "queued", queueHeld: true }); + expect(turnItems[0]?.status).toBe("running"); + }, + ); + it("returns empty while the latest run is not settled", () => { const tasks = derivePendingBackgroundWork({ latestRun: { id: "run-1" as never, ordinal: 1, status: "running" }, diff --git a/packages/shared/src/orchestrationV2PendingBackgroundWork.ts b/packages/shared/src/orchestrationV2PendingBackgroundWork.ts index 5a17bad321ed..4234a4a82446 100644 --- a/packages/shared/src/orchestrationV2PendingBackgroundWork.ts +++ b/packages/shared/src/orchestrationV2PendingBackgroundWork.ts @@ -104,7 +104,7 @@ type PendingBackgroundWorkTurnItem = { readonly type: OrchestrationV2TurnItem["type"]; readonly status: OrchestrationV2TurnItem["status"]; readonly title: string | null; - /** When present and the run is rolled_back, the item is abandoned, not pending. */ + /** Run ownership distinguishes abandoned foreground tools from live background work. */ readonly runId?: OrchestrationV2Run["id"] | string | null; readonly nativeItemRef?: { readonly nativeId: string | null; @@ -190,15 +190,22 @@ export function pendingBackgroundTurnItems run.status === "rolled_back").map((run) => String(run.id)), ); + const abandonedToolRunIds = new Set( + (input.runs ?? []) + .filter((run) => ["failed", "interrupted", "cancelled"].includes(run.status)) + .map((run) => String(run.id)), + ); return input.turnItems.filter( (item) => BACKGROUND_TURN_ITEM_TYPES.has(item.type) && isOrchestrationV2WorkActive(item.status) && !(item.type === "dynamic_tool" && isPersistentDynamicToolInput(item.input)) && - // Null/absent run id stays eligible; only known rolled_back runs drop. + // A failed foreground MCP call cannot finish after its run has ended. + // Commands and native children can outlive the root; their state remains authoritative. (item.runId === undefined || item.runId === null || - !rolledBackRunIds.has(String(item.runId))), + (!rolledBackRunIds.has(String(item.runId)) && + !(item.type === "dynamic_tool" && abandonedToolRunIds.has(String(item.runId))))), ); } @@ -216,7 +223,7 @@ export function pendingBackgroundTurnItems Date: Tue, 6 Oct 2026 14:01:48 -0400 Subject: [PATCH 2/5] fix(orchestration): adapt recovery tests to renamed upstream services --- apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts | 2 +- apps/server/src/orchestration-v2/ProjectionSettlement.test.ts | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts index dec6ca9d1fde..28a85fecb4cf 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts @@ -19,7 +19,7 @@ import * as ProviderAdapterRegistry from "../orchestration-v2/ProviderAdapterReg import * as ThreadManagementService from "../orchestration-v2/ThreadManagementService.ts"; import * as ProjectService from "../project/ProjectService.ts"; import * as SecretRequests from "../secrets/SecretRequests.ts"; -import * as ProviderRegistry from "../provider/Services/ProviderRegistry.ts"; +import * as ProviderRegistry from "../provider/ProviderRegistry.ts"; import * as ScheduledTaskService from "../scheduledTasks/ScheduledTaskService.ts"; import type { McpInvocationScope } from "./McpInvocationContext.ts"; import * as OrchestratorMcpService from "./OrchestratorMcpService.ts"; diff --git a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts index aa6a8c4fadc1..61d378bac486 100644 --- a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts @@ -172,7 +172,7 @@ it.effect.each([ ); it.effect.each([ - ["sql", SqlLayer], + ["sql", layerSql], ["memory", ProjectionStore.layerMemory], ] as const)( "%s: discovers settlement work with the same activity and background semantics as the shell", From 7b1ce91834a62565872302aee9f45eb956daff42 Mon Sep 17 00:00:00 2001 From: Jake Leventhal Date: Thu, 8 Oct 2026 15:42:54 -0400 Subject: [PATCH 3/5] fix(orchestration): preserve held queue completion and current test fixtures --- .../src/mcp/OrchestratorMcpService.queued-tasks.test.ts | 1 + .../src/orchestration-v2/DelegatedCompletionDelivery.test.ts | 1 - packages/provider-core/src/server/subagentProjection.ts | 4 ++-- 3 files changed, 3 insertions(+), 3 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts index 28a85fecb4cf..ff20ca9366f8 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts @@ -137,6 +137,7 @@ function fixture( getThreadRecords: (id) => Effect.succeed(id === parentThreadId ? parent : child), getProjectThreadRecords: () => Effect.succeed(child), stopDelegatedTasks: () => Effect.void, + delegatedTaskResultPending: () => Effect.succeed(true), getTimelinePage: () => Effect.succeed({ items: [], totalItems: 0, hasMore: false }), dispatch: (command) => Effect.suspend(() => { diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 719d2bdf1b4d..5fd824d2b7cc 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -8,7 +8,6 @@ import { MessageId, type ModelSelection, type OrchestrationV2AppThread, - TurnItemId, NodeId, type OrchestrationV2Run, ProjectId, diff --git a/packages/provider-core/src/server/subagentProjection.ts b/packages/provider-core/src/server/subagentProjection.ts index ff2a81290d6c..e552f147b23e 100644 --- a/packages/provider-core/src/server/subagentProjection.ts +++ b/packages/provider-core/src/server/subagentProjection.ts @@ -232,7 +232,7 @@ export function delegatedTaskProgress(projection: { // A held queue entry can predate many finished follow-ups. It is still // pending intent, but it must not make an idle provider look running. const pendingRuns = workRuns - .filter((run) => !terminal(run.status) && !(run.status === "queued" && run.queueHeld === true)) + .filter((run) => !terminal(run.status)) .toSorted((a, b) => b.ordinal - a.ordinal); const pendingRun = pendingRuns.find((run) => run.status !== "queued") ?? pendingRuns[0]; const children = @@ -250,7 +250,7 @@ export function delegatedTaskProgress(projection: { .toSorted((a, b) => (runRanAfter(a, b) ? -1 : runRanAfter(b, a) ? 1 : 0))[0]; return { state: - pendingRun !== undefined || resultRun === undefined + pendingRuns.some((run) => !(run.status === "queued" && run.queueHeld === true)) || resultRun === undefined ? ("working" as const) : children ? ("waiting_for_children" as const) From a376f76905d1de06aab604be99995d67b9ff1898 Mon Sep 17 00:00:00 2001 From: Jake Leventhal Date: Thu, 8 Oct 2026 15:43:59 -0400 Subject: [PATCH 4/5] fix(test): retain upstream held queue activity semantics --- .../server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts index ff20ca9366f8..229740d12611 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.queued-tasks.test.ts @@ -167,7 +167,8 @@ it.effect("reports an old held continuation as queued and keeps the latest resul const task = yield* service.taskStatus(scope, taskId); assert.equal(task.status, "queued"); assert.equal(task.workState, "working"); - assert.equal(task.hasPendingChildRuns, true); + // Held input remains visible in the queue counts, but does not count as active child work. + assert.equal(task.hasPendingChildRuns, false); assert.equal(task.latestTerminalRunId, RunId.make("run:queued-task:4")); assert.equal(task.latestTerminalSummary, "Latest completed follow-up result"); assert.isNull(task.summary); From 77c5d36ed4fc284e100dcdc75699f3e2f0ba5ab1 Mon Sep 17 00:00:00 2001 From: Jake Leventhal Date: Thu, 8 Oct 2026 15:45:21 -0400 Subject: [PATCH 5/5] fix(orchestration): retain upstream held queue completion repairs --- apps/server/src/mcp/OrchestratorMcpService.ts | 8 +- .../DelegatedCompletionDelivery.test.ts | 619 +----------------- .../FoundationPersistence.test.ts | 1 - .../src/orchestration-v2/Orchestrator.ts | 187 ++---- .../ProjectionRecovery.test.ts | 125 ++-- .../ProjectionSettlement.test.ts | 31 - .../src/orchestration-v2/ProjectionStore.ts | 26 +- .../src/server/subagentProjection.ts | 14 +- 8 files changed, 152 insertions(+), 859 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index 3998aadfb78c..b2c7aecfc517 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -1245,6 +1245,10 @@ const make = Effect.gen(function* () { const childRun = delegatedTaskRun(childControls, task); const terminalRun = latestTerminalResultRun(childControls, childRun); const progress = delegatedTaskProgress(childControls); + const pendingRuns = childControls.runs + .filter((run) => !ThreadManagementService.isTerminalRunStatus(run.status)) + .toSorted((a, b) => b.ordinal - a.ordinal); + const pendingRun = pendingRuns.find((run) => run.status !== "queued") ?? pendingRuns[0]; const resultRunIds = [ ...new Set( [progress.resultRun?.id, terminalRun?.id].filter((id): id is RunId => id !== undefined), @@ -1285,8 +1289,8 @@ const make = Effect.gen(function* () { ) : workState === "result_available" ? taskStatusForRun(progress.resultRun ?? childRun) - : progress.pendingRun !== undefined - ? taskStatusForRun(progress.pendingRun) + : pendingRun !== undefined + ? taskStatusForRun(pendingRun) : taskStatusForRun(childRun) === "queued" ? "queued" : "running"; diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 5fd824d2b7cc..0c5e38b8a6d8 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -7,7 +7,6 @@ import { EventId, MessageId, type ModelSelection, - type OrchestrationV2AppThread, NodeId, type OrchestrationV2Run, ProjectId, @@ -15,7 +14,6 @@ import { ProviderInstanceId, ProviderThreadId, RunId, - RunAttemptId, ThreadId, TurnItemId, } from "@t3tools/contracts"; @@ -40,7 +38,6 @@ import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; import * as EventSink from "./EventSink.ts"; import * as Orchestrator from "./Orchestrator.ts"; import * as ProviderRuntimeRecoveryService from "./ProviderRuntimeRecoveryService.ts"; -import type { ProviderAdapterV2Shape } from "@t3tools/provider-core/server/ProviderAdapter"; import { continueRestartedRun } from "./RestartContinuation.ts"; import * as RuntimeLayer from "./runtimeLayer.ts"; import * as ProviderTurnStartServiceTestkit from "./ProviderTurnStartService.testkit.ts"; @@ -107,7 +104,7 @@ const layerTestProviderInstanceRegistry = Layer.succeed( }, ); -const layerOrchestrationTest = Layer.mergeAll(RuntimeLayer.layer, RuntimeLayer.layerEventSink).pipe( +const layerTest = Layer.mergeAll(RuntimeLayer.layer, RuntimeLayer.layerEventSink).pipe( Layer.provideMerge(RuntimeLayer.layerProjectService), Layer.provide( Layer.mock(WorkspacePaths.WorkspacePaths)({ @@ -144,8 +141,6 @@ const layerOrchestrationTest = Layer.mergeAll(RuntimeLayer.layer, RuntimeLayer.l Layer.provide(layerPlatformTest), ); -const layerTest = layerOrchestrationTest.pipe(Layer.provide(SqlitePersistence.layerMemory)); - const seedParentWithTerminalTask = (input: { readonly threadId: ThreadId; readonly projectId: ProjectId; @@ -298,356 +293,6 @@ const seedParentWithTerminalTask = (input: { }); it.layer(layerTest)("delegated completion delivery repairs", (it) => { - it.effect( - "keeps an older held continuation pending after later terminal runs and tracks resumed work", - () => - Effect.gen(function* () { - const orchestrator = yield* Orchestrator.OrchestratorV2; - const sink = yield* EventSink.EventSinkV2; - const now = yield* DateTime.now; - const threadId = ThreadId.make("thread:held-child-parent"); - const childThreadId = ThreadId.make("thread:held-child"); - const taskId = NodeId.make("node:held-child-task"); - const parentRunId = RunId.make("run:held-child-parent"); - const rootNodeId = NodeId.make("node:held-child-parent-root"); - const projectId = ProjectId.make("project:held-child"); - yield* seedParentWithTerminalTask({ - threadId, - projectId, - runId: parentRunId, - rootNodeId, - taskId, - deliveryState: "disposed", - now, - }); - yield* orchestrator.dispatch({ - type: "thread.create", - commandId: CommandId.make("command:held-child-create"), - threadId: childThreadId, - projectId, - title: "Held continuation", - modelSelection, - runtimeMode: "full-access", - interactionMode: "default", - branch: null, - worktreePath: null, - createdBy: "agent", - creationSource: "mcp", - }); - const parent = yield* orchestrator.getThreadProjection(threadId); - const child = yield* orchestrator.getThreadProjection(childThreadId); - const task = { - ...parent.subagents[0]!, - childThreadId, - status: "running" as const, - result: null, - completedAt: null, - completionDelivery: undefined, - }; - const makeRun = ( - ordinal: number, - status: OrchestrationV2Run["status"], - ): OrchestrationV2Run => ({ - ...parent.runs[0]!, - id: RunId.make(`run:held-child:${ordinal}`), - threadId: childThreadId, - ordinal, - providerThreadId: null, - rootNodeId: null, - activeAttemptId: null, - userMessageId: MessageId.make(`message:held-child:${ordinal}`), - status, - startedAt: status === "queued" ? null : now, - completedAt: status === "queued" ? null : now, - queueHeld: status === "queued", - delegatedCompletion: undefined, - }); - const runs = [ - makeRun(1, "failed"), - makeRun(2, "queued"), - makeRun(3, "completed"), - makeRun(4, "completed"), - ]; - const afterSequence = yield* sink.latestSequence(); - yield* sink.write({ - events: [ - { - id: EventId.make("event:held-child-link"), - type: "thread.metadata-updated", - threadId: childThreadId, - occurredAt: now, - payload: { - ...child.thread, - lineage: { - parentThreadId: threadId, - relationshipToParent: "subagent", - rootThreadId: threadId, - }, - forkedFrom: { type: "node", nodeId: taskId }, - }, - }, - { - id: EventId.make("event:held-child-task"), - type: "subagent.updated", - threadId, - occurredAt: now, - payload: task, - }, - { - id: EventId.make("event:held-child-node"), - type: "node.updated", - threadId, - occurredAt: now, - payload: { - id: taskId, - threadId, - runId: parentRunId, - parentNodeId: rootNodeId, - rootNodeId, - kind: "subagent", - status: "running", - countsForRun: false, - providerThreadId: null, - providerTurnId: null, - nativeItemRef: null, - runtimeRequestId: null, - checkpointScopeId: null, - startedAt: now, - completedAt: null, - }, - }, - { - id: EventId.make("event:held-child-item"), - type: "turn-item.updated", - threadId, - occurredAt: now, - payload: { - id: TurnItemId.make("item:held-child"), - threadId, - runId: parentRunId, - nodeId: taskId, - providerThreadId: null, - providerTurnId: null, - nativeItemRef: null, - parentItemId: null, - ordinal: 1, - status: "running", - title: null, - startedAt: now, - completedAt: null, - updatedAt: now, - type: "subagent", - subagentId: taskId, - prompt: task.prompt, - origin: "app_owned", - driver, - providerInstanceId: modelSelection.instanceId, - childThreadId, - result: null, - }, - }, - ...runs.map((run) => ({ - id: EventId.make(`event:${run.id}`), - type: "run.updated" as const, - threadId: childThreadId, - runId: run.id, - occurredAt: now, - payload: run, - })), - { - id: EventId.make("event:held-child-final-result"), - type: "message.updated", - threadId: childThreadId, - occurredAt: now, - payload: { - id: MessageId.make("message:held-child-result"), - threadId: childThreadId, - runId: runs[3]!.id, - nodeId: null, - role: "assistant", - createdBy: "agent", - creationSource: "provider", - text: "Latest completed result", - attachments: [], - streaming: false, - createdAt: now, - updatedAt: now, - }, - }, - ], - }); - const awaitTaskStatus = (sequence: number, status: "pending" | "running" | "completed") => - sink.stream({ afterSequence: sequence, eventType: "subagent.updated" }).pipe( - Stream.filter( - (stored) => - stored.event.type === "subagent.updated" && - stored.event.payload.id === taskId && - stored.event.payload.status === status, - ), - Stream.take(1), - Stream.runDrain, - ); - yield* awaitTaskStatus(afterSequence, "pending"); - const pending = yield* orchestrator.getThreadProjection(threadId); - assert.equal(pending.subagents[0]?.status, "pending"); - assert.isNull(pending.subagents[0]?.result); - assert.equal(pending.nodes.find((node) => node.id === taskId)?.status, "pending"); - assert.equal(pending.turnItems.find((item) => item.type === "subagent")?.status, "pending"); - assert.isFalse( - pending.contextTransfers.some((transfer) => transfer.type === "subagent_result"), - ); - const held = yield* orchestrator.getThreadProjection(childThreadId); - assert.equal(held.runs.find((run) => run.id === runs[1]!.id)?.status, "queued"); - assert.isTrue(held.runs.find((run) => run.id === runs[1]!.id)?.queueHeld); - assert.lengthOf(held.runs, 4); - assert.equal( - held.messages.find((message) => message.role === "assistant")?.text, - "Latest completed result", - ); - - // A real resumed run must restore Running even when later ordinals - // already hold completed history. Recovery never performs this resume. - const beforeResume = yield* sink.latestSequence(); - yield* sink.write({ - events: [ - { - id: EventId.make("event:held-child-resumed"), - type: "run.updated", - threadId: childThreadId, - runId: runs[1]!.id, - occurredAt: now, - payload: { - ...runs[1]!, - status: "running", - queueHeld: false, - startedAt: now, - completedAt: null, - }, - }, - ], - }); - yield* awaitTaskStatus(beforeResume, "running"); - const running = yield* orchestrator.getThreadProjection(threadId); - assert.equal(running.subagents[0]?.status, "running"); - assert.equal(running.nodes.find((node) => node.id === taskId)?.status, "running"); - assert.equal(running.turnItems.find((item) => item.type === "subagent")?.status, "running"); - assert.isNull(running.subagents[0]?.result); - const beforeComplete = yield* sink.latestSequence(); - yield* sink.write({ - events: [ - { - id: EventId.make("event:held-child-resume-completed"), - type: "run.updated", - threadId: childThreadId, - runId: runs[1]!.id, - occurredAt: now, - payload: { - ...runs[1]!, - status: "completed", - queueHeld: false, - startedAt: now, - completedAt: now, - }, - }, - ], - }); - yield* awaitTaskStatus(beforeComplete, "completed"); - const completed = yield* orchestrator.getThreadProjection(threadId); - assert.equal(completed.subagents[0]?.result, "Latest completed result"); - assert.lengthOf( - completed.contextTransfers.filter((transfer) => transfer.type === "subagent_result"), - 1, - ); - - const siblingThreadId = ThreadId.make("thread:held-child-live-sibling"); - const siblingTaskId = NodeId.make("node:held-child-live-sibling"); - const beforeFollowup = yield* sink.latestSequence(); - const followup = makeRun(5, "queued"); - yield* sink.write({ - events: [ - { - id: EventId.make("event:held-child-new-followup"), - type: "run.created", - threadId: childThreadId, - runId: followup.id, - occurredAt: now, - payload: followup, - }, - { - id: EventId.make("event:held-child-live-sibling"), - type: "thread.created", - threadId: siblingThreadId, - occurredAt: now, - payload: { - ...child.thread, - id: siblingThreadId, - lineage: { - parentThreadId: threadId, - relationshipToParent: "subagent", - rootThreadId: threadId, - }, - forkedFrom: { type: "node", nodeId: siblingTaskId }, - }, - }, - { - id: EventId.make("event:held-child-live-sibling-task"), - type: "subagent.updated", - threadId, - occurredAt: now, - payload: { - ...task, - id: siblingTaskId, - childThreadId: siblingThreadId, - status: "pending", - }, - }, - { - id: EventId.make("event:held-child-live-sibling-run"), - type: "run.created", - threadId: siblingThreadId, - occurredAt: now, - payload: { - ...makeRun(1, "running"), - id: RunId.make("run:held-child-live-sibling"), - threadId: siblingThreadId, - completedAt: null, - }, - }, - ], - }); - // The same created-run stream processes the follow-up before this live - // sibling's receipt. The published result must remain terminal. - yield* sink.stream({ afterSequence: beforeFollowup, eventType: "subagent.updated" }).pipe( - Stream.filter( - (stored) => - stored.event.type === "subagent.updated" && - stored.event.payload.id === siblingTaskId && - stored.event.payload.status === "running", - ), - Stream.take(1), - Stream.runDrain, - ); - const afterFollowup = yield* orchestrator.getThreadProjection(threadId); - const published = afterFollowup.subagents.find((candidate) => candidate.id === taskId); - assert.equal(published?.status, "completed"); - assert.equal(published?.result, "Latest completed result"); - assert.equal( - afterFollowup.subagents.find((candidate) => candidate.id === siblingTaskId)?.status, - "running", - ); - assert.lengthOf( - afterFollowup.contextTransfers.filter((transfer) => transfer.type === "subagent_result"), - 1, - ); - const childAfterFollowup = yield* orchestrator.getThreadProjection(childThreadId); - assert.equal( - childAfterFollowup.runs.find((run) => run.id === followup.id)?.status, - "queued", - ); - assert.isTrue(childAfterFollowup.runs.find((run) => run.id === followup.id)?.queueHeld); - }), - ); - it.effect("acceptance batches pending siblings without acknowledging their results", () => Effect.gen(function* () { const orchestrator = yield* Orchestrator.OrchestratorV2; @@ -1822,265 +1467,3 @@ it.layer(layerTest)("delegated tasks across a server restart", (it) => { }), ); }); - -it.effect.each([ - { - heldQueue: true, - name: "repairs a persisted Running task with held older input during startup without replaying it", - }, - { - heldQueue: false, - name: "publishes the recovered terminal result of an active child without queued work", - }, -])("$name", ({ heldQueue }) => { - const parentId = ThreadId.make("thread:startup-queue-parent"); - const childId = ThreadId.make("thread:startup-queue-child"); - const taskId = NodeId.make("node:startup-queue-task"); - const parentRunId = RunId.make("run:startup-queue-parent"); - const rootNodeId = NodeId.make("node:startup-queue-parent"); - const seed = Effect.gen(function* () { - const sink = yield* EventSink.EventSinkV2; - const now = yield* DateTime.now; - const thread = (id: ThreadId): OrchestrationV2AppThread => ({ - id, - projectId: ProjectId.make("project:startup-queue"), - title: "Persisted held queue", - createdBy: "agent", - creationSource: "mcp", - providerInstanceId: modelSelection.instanceId, - modelSelection, - runtimeMode: "full-access", - interactionMode: "default", - branch: null, - worktreePath: null, - activeProviderThreadId: null, - lineage: { parentThreadId: null, relationshipToParent: null, rootThreadId: parentId }, - forkedFrom: null, - createdAt: now, - updatedAt: now, - archivedAt: null, - settledOverride: null, - settledAt: null, - lastVisitedAt: null, - deletedAt: null, - }); - const run = (ordinal: number, status: OrchestrationV2Run["status"]): OrchestrationV2Run => ({ - id: RunId.make(`run:startup-queue:${ordinal}`), - threadId: childId, - ordinal, - providerInstanceId: modelSelection.instanceId, - modelSelection, - providerThreadId: null, - userMessageId: MessageId.make(`message:startup-queue:${ordinal}`), - rootNodeId: null, - activeAttemptId: null, - status, - requestedAt: now, - startedAt: status === "queued" ? null : now, - completedAt: status === "queued" || status === "running" ? null : now, - queueHeld: status === "queued", - checkpointId: null, - contextHandoffId: null, - }); - yield* sink.write({ - events: [ - { - id: EventId.make("event:startup-queue-parent"), - type: "thread.created", - threadId: parentId, - occurredAt: now, - payload: thread(parentId), - }, - { - id: EventId.make("event:startup-queue-child"), - type: "thread.created", - threadId: childId, - occurredAt: now, - payload: { - ...thread(childId), - lineage: { - parentThreadId: parentId, - relationshipToParent: "subagent", - rootThreadId: parentId, - }, - forkedFrom: { type: "node", nodeId: taskId }, - }, - }, - { - id: EventId.make("event:startup-queue-parent-run"), - type: "run.updated", - threadId: parentId, - runId: parentRunId, - occurredAt: now, - payload: { ...run(1, "completed"), id: parentRunId, threadId: parentId, rootNodeId }, - }, - { - id: EventId.make("event:startup-queue-node"), - type: "node.updated", - threadId: parentId, - occurredAt: now, - payload: { - id: taskId, - threadId: parentId, - runId: parentRunId, - parentNodeId: rootNodeId, - rootNodeId, - kind: "subagent", - status: "running", - countsForRun: false, - providerThreadId: null, - providerTurnId: null, - nativeItemRef: null, - runtimeRequestId: null, - checkpointScopeId: null, - startedAt: now, - completedAt: null, - }, - }, - { - id: EventId.make("event:startup-queue-item"), - type: "turn-item.updated", - threadId: parentId, - occurredAt: now, - payload: { - id: TurnItemId.make("item:startup-queue-task"), - threadId: parentId, - runId: parentRunId, - nodeId: taskId, - providerThreadId: null, - providerTurnId: null, - nativeItemRef: null, - parentItemId: null, - ordinal: 1, - status: "running", - title: null, - startedAt: now, - completedAt: null, - updatedAt: now, - type: "subagent", - subagentId: taskId, - origin: "app_owned", - driver, - providerInstanceId: modelSelection.instanceId, - childThreadId: childId, - prompt: "Held continuation", - result: null, - }, - }, - { - id: EventId.make("event:startup-queue-task"), - type: "subagent.updated", - threadId: parentId, - occurredAt: now, - payload: { - id: taskId, - threadId: parentId, - runId: parentRunId, - parentNodeId: NodeId.make("node:startup-queue-parent"), - origin: "app_owned", - createdBy: "agent", - driver, - providerInstanceId: modelSelection.instanceId, - providerThreadId: null, - childThreadId: childId, - nativeTaskRef: null, - prompt: "Held continuation", - title: null, - model: null, - status: "running", - result: null, - startedAt: now, - completedAt: null, - updatedAt: now, - }, - }, - ...(heldQueue - ? [run(1, "failed"), run(2, "queued"), run(3, "completed"), run(4, "completed")] - : [run(1, "failed"), run(2, "running")] - ).map((payload) => ({ - id: EventId.make(`event:${payload.id}`), - type: "run.updated" as const, - threadId: childId, - runId: payload.id, - occurredAt: now, - payload, - })), - { - id: EventId.make("event:startup-queue-message"), - type: "message.updated", - threadId: childId, - occurredAt: now, - payload: { - id: MessageId.make("message:startup-queue:2"), - threadId: childId, - runId: RunId.make("run:startup-queue:2"), - nodeId: null, - role: "user", - createdBy: "agent", - creationSource: "mcp", - text: "Keep this genuine queued continuation", - attachments: [], - streaming: false, - createdAt: now, - updatedAt: now, - }, - }, - ], - }); - }); - const seededPersistence = Layer.effectDiscard(seed).pipe( - Layer.provideMerge( - RuntimeLayer.layerEventSink.pipe(Layer.provideMerge(SqlitePersistence.layerMemory)), - ), - ); - return Effect.gen(function* () { - const orchestrator = yield* Orchestrator.OrchestratorV2; - const beforeRuntimeRecovery = yield* orchestrator.getThreadProjection(parentId); - assert.equal(beforeRuntimeRecovery.subagents[0]?.status, "running"); - const runtimeRecovery = yield* ProviderRuntimeRecoveryService.ProviderRuntimeRecoveryService; - yield* runtimeRecovery.recover; - yield* orchestrator.recoverDelegatedTasks; - const parent = yield* orchestrator.getThreadProjection(parentId); - const child = yield* orchestrator.getThreadProjection(childId); - if (heldQueue) { - assert.equal(parent.subagents[0]?.status, "pending"); - assert.equal(parent.nodes.find((node) => node.id === taskId)?.status, "pending"); - assert.equal(parent.turnItems.find((item) => item.type === "subagent")?.status, "pending"); - assert.isNull(parent.subagents[0]?.result); - assert.isFalse( - parent.contextTransfers.some((transfer) => transfer.type === "subagent_result"), - ); - const queued = child.runs.find((candidate) => candidate.ordinal === 2); - assert.equal(queued?.status, "queued"); - assert.isTrue(queued?.queueHeld); - assert.isNull(queued?.activeAttemptId); - assert.isNull(queued?.providerThreadId); - assert.equal( - child.messages.find((message) => message.runId === queued?.id)?.text, - "Keep this genuine queued continuation", - ); - assert.lengthOf(child.runs, 4); - assert.isFalse( - child.runs.some((candidate) => - ["preparing", "starting", "running"].includes(candidate.status), - ), - ); - } else { - assert.equal(parent.subagents[0]?.status, "cancelled"); - assert.equal(parent.nodes.find((node) => node.id === taskId)?.status, "cancelled"); - assert.equal(parent.turnItems.find((item) => item.type === "subagent")?.status, "cancelled"); - assert.equal(parent.subagents[0]?.result, "Child task ended with status cancelled."); - assert.lengthOf( - parent.contextTransfers.filter((transfer) => transfer.type === "subagent_result"), - 1, - ); - assert.equal(child.runs.find((candidate) => candidate.ordinal === 2)?.status, "cancelled"); - assert.lengthOf(child.runs, 2); - assert.isFalse( - child.runs.some((candidate) => - ["preparing", "starting", "running", "queued"].includes(candidate.status), - ), - ); - } - }).pipe(Effect.provide(layerOrchestrationTest.pipe(Layer.provide(seededPersistence)))); -}); diff --git a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts index fd0290afbdc9..9a1e43a8b44b 100644 --- a/apps/server/src/orchestration-v2/FoundationPersistence.test.ts +++ b/apps/server/src/orchestration-v2/FoundationPersistence.test.ts @@ -2779,7 +2779,6 @@ it.layer(layerTest)("orchestration V2 foundation persistence", (it) => { { id: sessionId, driver: "codex", providerInstanceId, status: "running" }, ], providerTurns: [{ providerThreadId, runAttemptId: attemptId, status: "running" }], - subagents: [], } as unknown as OrchestrationV2ThreadProjection; const recovery = yield* ProviderRuntimeRecovery.make.pipe( Effect.provide( diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index bef44d78b954..977b225f5256 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -9572,8 +9572,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); /** - * Synchronizes unpublished child work, then transfers a settled child's - * result into its parent and offers the parent wake. Every mutation here targets the PARENT thread, so callers must hold + * Transfers a terminal child's result into its parent and offers the parent + * wake. Every mutation here targets the PARENT thread, so callers must hold * the parent thread's dispatch lock rather than the child's: the * delegated_task.wake-policy handler rewrites the same subagent row under * that lock with a full-row payload, and unserialized writers clobber each @@ -9600,6 +9600,38 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return; } const progress = delegatedTaskProgress(childControls); + if (progress.state !== "result_available") return; + const childRun = progress.resultRun; + if (childRun === undefined) return; + const terminalStatus = delegatedTaskTerminalStatus(childRun.status); + if (terminalStatus === null) { + return; + } + if ( + yield* childAwaitsRestartContinuation( + childControls.runs, + childRun, + options?.settledContinuationOf, + ) + ) { + return; + } + + const childResult = yield* projectionStore.getThreadRecords( + childThreadId, + ["messages", "turnItems"], + { + messageRoles: ["assistant"], + messageRunIds: [childRun.id], + turnItemRunId: childRun.id, + turnItemTypes: ["assistant_message", "error"], + }, + ); + const childProjection = { + ...childControls, + messages: childResult.messages, + turnItems: childResult.turnItems, + }; const parentThreadId = childControls.thread.lineage.parentThreadId; const parentProjection = yield* projectionStore.getThreadRecords( parentThreadId, @@ -9634,105 +9666,16 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio return; } - // A held queued continuation still belongs to this task. Publish no - // result until it settles, but keep the parent's durable status honest. - if (task.result !== null) return; - const parentNode = parentProjection.nodes.find((candidate) => candidate.id === task.id); - const parentTurnItem = parentProjection.turnItems.find( - (candidate) => candidate.type === "subagent" && candidate.subagentId === task.id, - ); - if (progress.state !== "result_available") { - const status = - progress.state === "waiting_for_children" - ? ("waiting" as const) - : progress.pendingRun === undefined || progress.pendingRun.status === "queued" - ? ("pending" as const) - : ("running" as const); - if ( - task.status === status && - (parentNode === undefined || parentNode.status === status) && - (parentTurnItem === undefined || parentTurnItem.status === status) - ) { - return; - } - const now = yield* DateTime.now; - yield* writeSystemEvents([ - { - type: "subagent.updated", - threadId: parentThreadId, - ...(task.runId === null ? {} : { runId: task.runId }), - nodeId: task.id, - driver: task.driver, - occurredAt: now, - payload: { ...task, status, completedAt: null, updatedAt: now }, - }, - ...(parentNode === undefined - ? [] - : [ - { - type: "node.updated" as const, - threadId: parentThreadId, - ...(parentNode.runId === null ? {} : { runId: parentNode.runId }), - nodeId: parentNode.id, - driver: task.driver, - occurredAt: now, - payload: { ...parentNode, status, completedAt: null }, - }, - ]), - ...(parentTurnItem === undefined - ? [] - : [ - { - type: "turn-item.updated" as const, - threadId: parentThreadId, - ...(parentTurnItem.runId === null ? {} : { runId: parentTurnItem.runId }), - ...(parentTurnItem.nodeId === null ? {} : { nodeId: parentTurnItem.nodeId }), - driver: task.driver, - occurredAt: now, - payload: { ...parentTurnItem, status, completedAt: null, updatedAt: now }, - }, - ]), - ]); - return; - } - const childRun = progress.resultRun; - if (childRun === undefined) return; - const terminalStatus = delegatedTaskTerminalStatus(childRun.status); - if (terminalStatus === null) { - return; - } - - if ( - yield* childAwaitsRestartContinuation( - childControls.runs, - childRun, - options?.settledContinuationOf, - ) - ) { - return; - } - - const childResult = yield* projectionStore.getThreadRecords( - childThreadId, - ["messages", "turnItems"], - { - messageRoles: ["assistant"], - messageRunIds: [childRun.id], - turnItemRunId: childRun.id, - turnItemTypes: ["assistant_message", "error"], - }, - ); - const childProjection = { - ...childControls, - messages: childResult.messages, - turnItems: childResult.turnItems, - }; const now = yield* DateTime.now; const result = subagentResultForRun(childProjection, childRun); const parentRun = task.runId === null ? undefined : parentProjection.runs.find((candidate) => candidate.id === task.runId); + const parentNode = parentProjection.nodes.find((candidate) => candidate.id === task.id); + const parentTurnItem = parentProjection.turnItems.find( + (candidate) => candidate.type === "subagent" && candidate.subagentId === task.id, + ); const updatedTask: OrchestrationV2Subagent = { ...task, providerThreadId: childRun.providerThreadId, @@ -10686,7 +10629,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const dispatchWithReceipt = (command: OrchestrationV2ServerCommand) => threadDispatch.withLock(commandThreadId(command), dispatchWithReceiptEffect(command)); - const handleRunUpdate = (stored: OrchestrationV2StoredEvent) => + const handleTerminalRun = (stored: OrchestrationV2StoredEvent) => Effect.gen(function* () { const threadId = stored.event.threadId; // finalize writes the parent thread and startNextQueuedRun writes this @@ -10695,20 +10638,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio // while holding the parent lock, so nesting the parent lock inside the // child lock here would invert that order, and the keyed executor's // semaphores are neither reentrant nor deadlock-aware. - // Nonterminal updates synchronize the task without promoting queued work. - if ( - String(stored.commandId).startsWith("command:runtime-reconcile:") || - stored.event.type !== "run.updated" || - !["completed", "interrupted", "failed", "cancelled", "rolled_back"].includes( - stored.event.payload.status, - ) - ) { - const parentThreadId = yield* appOwnedSubagentParentThreadId(threadId); - if (parentThreadId !== undefined) { - yield* threadDispatch.withLock(parentThreadId, finalizeAppOwnedSubagent(threadId)); - } - return; - } if (stored.event.type === "run.updated") { yield* threadDispatch.withLock( threadId, @@ -10738,7 +10667,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } }).pipe( Effect.catchCause((cause) => - Effect.logWarning("Failed to react to V2 run update", { + Effect.logWarning("Failed to react to terminal V2 run", { threadId: stored.event.threadId, sequence: stored.sequence, cause, @@ -10746,24 +10675,28 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ), ); - // Recover historical child state from projections below. Live run updates - // synchronize unpublished tasks; only terminal updates promote the queue. - // Replaying full history would delay this work as the database grows. - const runEventsAfterSequence = yield* eventSink.latestSequence().pipe(Effect.orDie); + // Historical terminal events are already represented by the projections + // below. Replaying the full event table on every server start delays live + // queue promotion in proportion to the lifetime size of the database. + const terminalEventsAfterSequence = yield* eventSink.latestSequence().pipe(Effect.orDie); // Queue promotion can wait on a provider or a thread lock. Subscribe to run // updates before buffering so that wait never retains unrelated tool bodies. - yield* Stream.merge( - eventSink.stream({ afterSequence: runEventsAfterSequence, eventType: "run.updated" }), - eventSink.stream({ afterSequence: runEventsAfterSequence, eventType: "run.created" }), - ).pipe( - Stream.filter( - (stored) => - !String(stored.commandId).startsWith("command:runtime-reconcile:") && - (stored.event.type === "run.updated" || stored.event.type === "run.created"), - ), - Stream.runForEach(handleRunUpdate), - Effect.forkDetach, - ); + yield* eventSink + .stream({ afterSequence: terminalEventsAfterSequence, eventType: "run.updated" }) + .pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && + !String(stored.commandId).startsWith("command:runtime-reconcile:") && + (stored.event.payload.status === "completed" || + stored.event.payload.status === "interrupted" || + stored.event.payload.status === "failed" || + stored.event.payload.status === "cancelled" || + stored.event.payload.status === "rolled_back"), + ), + Stream.runForEach(handleTerminalRun), + Effect.forkDetach, + ); // Settles child results and completion deliveries whose runs ended without // the listener above: before this boot, or in runtime reconciliation, which diff --git a/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts b/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts index 5c7046e93e83..c7ea0530d351 100644 --- a/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionRecovery.test.ts @@ -224,70 +224,71 @@ it.effect("selects unfinished recovery work without reading settled thread histo }).pipe(Effect.provide(layerTest)), ); -it.effect( - "recovers unpublished subagent lifecycle states until their cross-thread transfer exists", - () => - Effect.gen(function* () { - const projections = yield* ProjectionStore.ProjectionStoreV2; - const sql = yield* SqlClient.SqlClient; - const now = yield* DateTime.now; - const parent = yield* createThread("subagent-parent"); - const children: Array = []; - for (const name of ["terminal", "archived", "deleted", "running", "queued"]) { - const child = yield* createThread(`subagent-${name}`, { - lineage: { - parentThreadId: parent, - relationshipToParent: "subagent", - rootThreadId: parent, - }, - forkedFrom: { type: "node", nodeId: NodeId.make(`node:${name}`) }, - archivedAt: name === "archived" ? now : null, - deletedAt: name === "deleted" ? now : null, +it.effect("recovers terminal subagent results until their cross-thread transfer exists", () => + Effect.gen(function* () { + const projections = yield* ProjectionStore.ProjectionStoreV2; + const sql = yield* SqlClient.SqlClient; + const now = yield* DateTime.now; + const parent = yield* createThread("subagent-parent"); + const children: Array = []; + for (const name of ["terminal", "archived", "deleted", "running", "held", "queued"]) { + const child = yield* createThread(`subagent-${name}`, { + lineage: { parentThreadId: parent, relationshipToParent: "subagent", rootThreadId: parent }, + forkedFrom: { type: "node", nodeId: NodeId.make(`node:${name}`) }, + archivedAt: name === "archived" ? now : null, + deletedAt: name === "deleted" ? now : null, + }); + yield* createRun(child, name === "running" ? "running" : "completed"); + // A wake queued behind the result: held by Stop or a restart, or still deliverable. + if (name === "held" || name === "queued") { + yield* createRun(child, "queued", { + ordinal: 2, + startedAt: null, + ...(name === "held" ? { queueHeld: true } : {}), }); - yield* createRun( - child, - name === "running" ? "running" : name === "queued" ? "queued" : "completed", - ); - children.push(child); } - assert.deepEqual( - new Set(yield* projections.getRecoveryThreadIds("subagent-results")), - new Set(children.filter((_, index) => index !== 2)), - ); - const completed = children[0]!; - const transferId = ContextTransferId.make("transfer:recovery:subagent-result"); - yield* projections.apply({ - id: EventId.make("event:recovery:subagent-result"), - type: "context-transfer.created", - threadId: parent, - occurredAt: now, - payload: { - id: transferId, - type: "subagent_result", - sourceThreadId: completed, - targetThreadId: parent, - sourcePoint: { threadId: completed }, - basePoint: null, - sourceProviderInstanceId: providerInstanceId, - targetProviderInstanceId: providerInstanceId, - targetRunId: null, - status: "pending", - resolution: null, - createdBy: "system", - error: null, - createdAt: now, - updatedAt: now, - consumedAt: null, - }, - }); - const archivedChild = children[1]; - assert.isDefined(archivedChild); - assert.deepEqual( - new Set(yield* projections.getRecoveryThreadIds("subagent-results")), - new Set([archivedChild, children[3], children[4]]), - ); - assert.deepEqual(yield* projections.getUnreadableThreadIds(), []); - yield* sql` + children.push(child); + } + const terminalChildren = new Set([children[0]!, children[1]!, children[4]!]); + assert.deepEqual( + new Set(yield* projections.getRecoveryThreadIds("subagent-results")), + terminalChildren, + ); + + const completed = children[0]!; + const transferId = ContextTransferId.make("transfer:recovery:subagent-result"); + yield* projections.apply({ + id: EventId.make("event:recovery:subagent-result"), + type: "context-transfer.created", + threadId: parent, + occurredAt: now, + payload: { + id: transferId, + type: "subagent_result", + sourceThreadId: completed, + targetThreadId: parent, + sourcePoint: { threadId: completed }, + basePoint: null, + sourceProviderInstanceId: providerInstanceId, + targetProviderInstanceId: providerInstanceId, + targetRunId: null, + status: "pending", + resolution: null, + createdBy: "system", + error: null, + createdAt: now, + updatedAt: now, + consumedAt: null, + }, + }); + const archivedChild = children[1]; + assert.isDefined(archivedChild); + assert.deepEqual( + new Set(yield* projections.getRecoveryThreadIds("subagent-results")), + new Set([archivedChild, children[4]!]), + ); + assert.deepEqual(yield* projections.getUnreadableThreadIds(), []); + yield* sql` UPDATE orchestration_v2_projection_context_transfers SET payload_json = '{}' WHERE context_transfer_id = ${transferId} `; diff --git a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts index 61d378bac486..f626d0c38583 100644 --- a/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts +++ b/apps/server/src/orchestration-v2/ProjectionSettlement.test.ts @@ -140,37 +140,6 @@ const createItem = Effect.fn(function* ( }); }); -it.effect.each([ - ["sql", layerSql], - ["memory", ProjectionStore.layerMemory], -] as const)( - "%s: excludes abandoned foreground tools from shells while retaining their history for recovery", - ([, testLayer]) => - Effect.gen(function* () { - const store = yield* ProjectionStore.ProjectionStoreV2; - for (const status of ["failed", "interrupted", "cancelled"] as const) { - const threadId = yield* createThread(`orphan-handoff-${status}`); - const failedRunId = yield* createRun(threadId, status); - yield* createItem(threadId, failedRunId, "running"); - yield* createRun(threadId, "completed", 2); - yield* createRun(threadId, "completed", 3); - const allShells = yield* store.getShellSnapshot(); - const scopedShell = yield* store.getThreadShell(threadId); - const candidates = yield* store.getSettlementCandidates(threadId); - assert.deepEqual( - allShells.threads.find((thread) => thread.id === threadId)?.pendingBackgroundTasks, - [], - ); - assert.deepEqual(scopedShell?.pendingBackgroundTasks, []); - assert.deepEqual(candidates[0]?.pendingBackgroundTasks, []); - const history = yield* store.getThreadRecords(threadId, ["runs", "turnItems"]); - assert.equal(history.turnItems[0]?.status, "running"); - assert.equal(history.runs.find((run) => run.id === failedRunId)?.status, status); - assert.include(yield* store.getRecoveryThreadIds("runtime"), threadId); - } - }).pipe(Effect.provide(testLayer)), -); - it.effect.each([ ["sql", layerSql], ["memory", ProjectionStore.layerMemory], diff --git a/apps/server/src/orchestration-v2/ProjectionStore.ts b/apps/server/src/orchestration-v2/ProjectionStore.ts index cab22deaa8a3..8c33faef808b 100644 --- a/apps/server/src/orchestration-v2/ProjectionStore.ts +++ b/apps/server/src/orchestration-v2/ProjectionStore.ts @@ -556,7 +556,9 @@ function needsRecovery( projection.thread.lineage.relationshipToParent === "subagent" && parentThreadId !== null && projection.thread.forkedFrom?.type === "node" && - projection.runs.length > 0 && + ["completed", "interrupted", "failed", "cancelled", "rolled_back"].includes( + latestUnheldRun(projection.runs)?.status ?? "idle", + ) && !projection.contextTransfers.some( (transfer) => transfer.type === "subagent_result" && @@ -3574,10 +3576,17 @@ export const layer: Layer.Layer = json_extract(child.payload_json, '$.lineage.relationshipToParent') = 'subagent' AND json_extract(child.payload_json, '$.lineage.parentThreadId') IS NOT NULL AND json_extract(child.payload_json, '$.forkedFrom.type') = 'node' - AND EXISTS ( - SELECT 1 FROM orchestration_v2_projection_runs + -- A held queue waits for the user, so the newest unheld run + -- decides, matching latestUnheldRun. + AND ( + SELECT status FROM orchestration_v2_projection_runs WHERE thread_id = child.thread_id - ) + AND NOT ( + status = 'queued' + AND json_extract(payload_json, '$.queueHeld') IS 1 + ) + ORDER BY ordinal DESC LIMIT 1 + ) IN ('completed', 'interrupted', 'failed', 'cancelled', 'rolled_back') AND NOT EXISTS ( SELECT 1 FROM orchestration_v2_projection_context_transfers WHERE source_thread_id = child.thread_id @@ -5344,11 +5353,10 @@ export const layer: Layer.Layer = ON r.run_id = i.run_id WHERE i.type IN ('command_execution', 'dynamic_tool', 'subagent') AND i.status NOT IN ('completed', 'interrupted', 'failed', 'cancelled') - -- Foreground tools on failed/stopped runs cannot finish. Commands - -- and native children can outlive their root; keep those visible. + -- A rolled-back run's items are abandoned, not pending. Without + -- this the shell reports Waiting for work no one will finish, + -- matching the item_count query's exclusion above. AND (i.run_id IS NULL OR r.status <> 'rolled_back') - AND (i.type <> 'dynamic_tool' OR r.status IS NULL - OR r.status NOT IN ('failed', 'interrupted', 'cancelled')) ` : sql` SELECT i.thread_id, i.payload_json @@ -5358,8 +5366,6 @@ export const layer: Layer.Layer = WHERE i.type IN ('command_execution', 'dynamic_tool', 'subagent') AND i.status NOT IN ('completed', 'interrupted', 'failed', 'cancelled') AND (i.run_id IS NULL OR r.status <> 'rolled_back') - AND (i.type <> 'dynamic_tool' OR r.status IS NULL - OR r.status NOT IN ('failed', 'interrupted', 'cancelled')) AND i.thread_id IN ${sql.in(threadIds)} `; diff --git a/packages/provider-core/src/server/subagentProjection.ts b/packages/provider-core/src/server/subagentProjection.ts index e552f147b23e..7ee8d7c1aa8d 100644 --- a/packages/provider-core/src/server/subagentProjection.ts +++ b/packages/provider-core/src/server/subagentProjection.ts @@ -229,12 +229,11 @@ export function delegatedTaskProgress(projection: { const workRuns = projection.runs.filter( (run) => !monitorRuns.has(run.id) && run.status !== "rolled_back", ); - // A held queue entry can predate many finished follow-ups. It is still - // pending intent, but it must not make an idle provider look running. - const pendingRuns = workRuns - .filter((run) => !terminal(run.status)) - .toSorted((a, b) => b.ordinal - a.ordinal); - const pendingRun = pendingRuns.find((run) => run.status !== "queued") ?? pendingRuns[0]; + // A held queue waits for the user to resume it (after Stop, a restart, or a + // provider failure), so its runs are not work the task still owes. + const active = workRuns.some( + (run) => !terminal(run.status) && !(run.status === "queued" && run.queueHeld === true), + ); const children = projection.subagents.some( (task) => @@ -250,12 +249,11 @@ export function delegatedTaskProgress(projection: { .toSorted((a, b) => (runRanAfter(a, b) ? -1 : runRanAfter(b, a) ? 1 : 0))[0]; return { state: - pendingRuns.some((run) => !(run.status === "queued" && run.queueHeld === true)) || resultRun === undefined + active || resultRun === undefined ? ("working" as const) : children ? ("waiting_for_children" as const) : ("result_available" as const), resultRun, - pendingRun, }; }