From f6be005ea4e96aa8657ed1593c93ccd2a311b053 Mon Sep 17 00:00:00 2001 From: ashx-j <281260375+ashx-j@users.noreply.github.com> Date: Wed, 7 Oct 2026 11:29:20 +0200 Subject: [PATCH 1/4] fix(server): keep delegated tasks waiting for background commands --- .../OrchestratorMcpService.activity.test.ts | 1 + .../src/mcp/OrchestratorMcpService.test.ts | 284 ++++++++++-------- apps/server/src/mcp/OrchestratorMcpService.ts | 8 +- .../DelegatedCompletionDelivery.test.ts | 121 ++++++++ .../src/orchestration-v2/Orchestrator.ts | 64 +++- .../SubagentProjection.test.ts | 53 ++++ .../orchestration-v2/SubagentProjection.ts | 7 +- 7 files changed, 405 insertions(+), 133 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts b/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts index 7c6f9d28c950..c35de8a5a0e6 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.activity.test.ts @@ -310,6 +310,7 @@ it("taskStatus returns task.providerInstanceId rather than the driver kind", asy ], subagents: [], providerThreads: [], + turnItems: [], updatedAt: now, } as unknown as OrchestrationV2ThreadProjection; diff --git a/apps/server/src/mcp/OrchestratorMcpService.test.ts b/apps/server/src/mcp/OrchestratorMcpService.test.ts index 18b444b85177..2c73db802c12 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.test.ts @@ -12,6 +12,7 @@ import { ProviderInstanceId, RunId, ThreadId, + TurnItemId, type OrchestrationV2ThreadProjection, type ServerProvider, } from "@t3tools/contracts"; @@ -41,135 +42,164 @@ import { idleThreadProjection, liveThreadShell } from "./McpToolAccess.testkit.t import * as OrchestratorMcpService from "./OrchestratorMcpService.ts"; describe("OrchestratorMcpService", () => { - it.effect("retries terminal acknowledgement with a fresh command id", () => - Effect.gen(function* () { - const parentThreadId = ThreadId.make("thread:mcp-ack-parent"); - const childThreadId = ThreadId.make("thread:mcp-ack-child"); - const childRunId = RunId.make("run:mcp-ack-child"); - const taskId = NodeId.make("node:mcp-ack-task"); - const acknowledgementCommandIds = yield* Ref.make>([]); - const acknowledgementAttempts = yield* Ref.make(0); - const parentProjection = { - thread: { id: parentThreadId }, - runs: [], - contextTransfers: [], - subagents: [ - { - id: taskId, + it.effect.each(["nested task", "background command"] as const)( + "waits for a %s before acknowledging the result and retries failed acknowledgement", + (backgroundWork) => + Effect.gen(function* () { + const parentThreadId = ThreadId.make("thread:mcp-ack-parent"); + const childThreadId = ThreadId.make("thread:mcp-ack-child"); + const childRunId = RunId.make("run:mcp-ack-child"); + const taskId = NodeId.make("node:mcp-ack-task"); + const acknowledgementCommandIds = yield* Ref.make>([]); + const acknowledgementAttempts = yield* Ref.make(0); + const parentProjection = { + thread: { id: parentThreadId }, + runs: [], + contextTransfers: [], + subagents: [ + { + id: taskId, + threadId: parentThreadId, + origin: "app_owned", + childThreadId, + driver: "codex", + model: "gpt-5.6-terra", + result: "terminal result", + completionDelivery: { state: "pending" }, + }, + ], + } as unknown as OrchestrationV2ThreadProjection; + const childProjection = { + thread: { id: childThreadId }, + runs: [ + { + id: childRunId, + ordinal: 1, + status: "completed", + startedAt: DateTime.makeUnsafe("2026-10-03T10:00:00Z"), + completedAt: DateTime.makeUnsafe("2026-10-03T10:10:00Z"), + }, + { + id: RunId.make("run:mcp-ack-continuation"), + ordinal: 2, + status: "failed", + startedAt: DateTime.makeUnsafe("2026-10-03T10:02:00Z"), + completedAt: DateTime.makeUnsafe("2026-10-03T10:05:00Z"), + }, + ], + contextTransfers: [], + messages: [], + subagents: [], + providerThreads: [], + turnItems: [], + } as unknown as OrchestrationV2ThreadProjection; + let hasNestedWork = true; + const layerDependencies = Layer.mergeAll( + NodeServices.layer, + Layer.mock(ThreadManagementService.ThreadManagementService)({ + getThreadRecords: (threadId, records) => + Effect.succeed( + threadId === parentThreadId + ? hasNestedWork + ? { + ...parentProjection, + subagents: parentProjection.subagents.map((task) => ({ + ...task, + result: null, + status: "running" as const, + })), + } + : parentProjection + : hasNestedWork + ? { + ...childProjection, + subagents: + backgroundWork === "nested task" + ? [{ ...parentProjection.subagents[0]!, status: "running" as const }] + : [], + turnItems: + backgroundWork === "background command" && records.includes("turnItems") + ? [ + { + id: TurnItemId.make("background-command"), + threadId: childThreadId, + runId: childRunId, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + type: "command_execution" as const, + status: "running" as const, + title: null, + startedAt: null, + completedAt: null, + updatedAt: DateTime.makeUnsafe("2026-10-03T10:10:00Z"), + input: "test --watch", + }, + ] + : [], + } + : childProjection, + ), + dispatch: (command) => + Ref.update(acknowledgementCommandIds, (commandIds) => [ + ...commandIds, + String(command.commandId), + ]).pipe( + Effect.andThen(Ref.updateAndGet(acknowledgementAttempts, (count) => count + 1)), + Effect.flatMap((attempt) => + attempt === 1 + ? Effect.fail(new Error("simulated acknowledgement failure") as never) + : Effect.succeed({} as never), + ), + ), + }), + Layer.mock(ProviderRegistry.ProviderRegistry)({ getProviders: Effect.succeed([]) }), + Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({ + list: () => Effect.succeed([]), + }), + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(SecretRequests.SecretRequests)({}), + Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), + ); + const scope: McpInvocationScope = { + environmentId: EnvironmentId.make("environment:mcp-ack"), + requestNamespace: "provider-session:mcp-ack", + thread: { threadId: parentThreadId, - origin: "app_owned", - childThreadId, - driver: "codex", - model: "gpt-5.6-terra", - result: "terminal result", - completionDelivery: { state: "pending" }, + providerSessionId: "provider-session:mcp-ack", + providerInstanceId: ProviderInstanceId.make("codex"), }, - ], - } as unknown as OrchestrationV2ThreadProjection; - const childProjection = { - thread: { id: childThreadId }, - runs: [ - { - id: childRunId, - ordinal: 1, - status: "completed", - startedAt: DateTime.makeUnsafe("2026-10-03T10:00:00Z"), - completedAt: DateTime.makeUnsafe("2026-10-03T10:10:00Z"), - }, - { - id: RunId.make("run:mcp-ack-continuation"), - ordinal: 2, - status: "failed", - startedAt: DateTime.makeUnsafe("2026-10-03T10:02:00Z"), - completedAt: DateTime.makeUnsafe("2026-10-03T10:05:00Z"), - }, - ], - contextTransfers: [], - messages: [], - subagents: [], - providerThreads: [], - } as unknown as OrchestrationV2ThreadProjection; - let hasNestedWork = true; - const layerDependencies = Layer.mergeAll( - NodeServices.layer, - Layer.mock(ThreadManagementService.ThreadManagementService)({ - getThreadRecords: (threadId) => - Effect.succeed( - threadId === parentThreadId - ? hasNestedWork - ? { - ...parentProjection, - subagents: parentProjection.subagents.map((task) => ({ - ...task, - result: null, - status: "running" as const, - })), - } - : parentProjection - : hasNestedWork - ? { - ...childProjection, - subagents: [ - { ...parentProjection.subagents[0]!, status: "running" as const }, - ], - } - : childProjection, - ), - dispatch: (command) => - Ref.update(acknowledgementCommandIds, (commandIds) => [ - ...commandIds, - String(command.commandId), - ]).pipe( - Effect.andThen(Ref.updateAndGet(acknowledgementAttempts, (count) => count + 1)), - Effect.flatMap((attempt) => - attempt === 1 - ? Effect.fail(new Error("simulated acknowledgement failure") as never) - : Effect.succeed({} as never), - ), - ), - }), - Layer.mock(ProviderRegistry.ProviderRegistry)({ getProviders: Effect.succeed([]) }), - Layer.mock(ProviderAdapterRegistry.ProviderAdapterRegistryV2)({ - list: () => Effect.succeed([]), - }), - Layer.mock(ProjectService.ProjectService)({}), - Layer.mock(SecretRequests.SecretRequests)({}), - Layer.mock(ScheduledTaskService.ScheduledTaskService)({}), - ); - const scope: McpInvocationScope = { - environmentId: EnvironmentId.make("environment:mcp-ack"), - requestNamespace: "provider-session:mcp-ack", - thread: { - threadId: parentThreadId, - providerSessionId: "provider-session:mcp-ack", - providerInstanceId: ProviderInstanceId.make("codex"), - }, - client: undefined, - capabilities: new Set(["orchestration"]), - issuedAt: 1, - }; + client: undefined, + capabilities: new Set(["orchestration"]), + issuedAt: 1, + }; - yield* Effect.gen(function* () { - const service = yield* OrchestratorMcpService.OrchestratorMcpService; - const pending = yield* service.taskStatus(scope, taskId); - assert.equal(pending.status, "running"); - assert.equal(pending.workState, "waiting_for_children"); - assert.isNull(pending.summary); - assert.equal(yield* Ref.get(acknowledgementAttempts), 0); - hasNestedWork = false; - const error = yield* service.taskStatus(scope, taskId).pipe(Effect.flip); - assert.equal(error.code, "orchestration_error"); + yield* Effect.gen(function* () { + const service = yield* OrchestratorMcpService.OrchestratorMcpService; + const pending = yield* service.taskStatus(scope, taskId); + assert.equal(pending.status, "running"); + assert.equal(pending.workState, "waiting_for_children"); + assert.isNull(pending.summary); + assert.equal(yield* Ref.get(acknowledgementAttempts), 0); + hasNestedWork = false; + const error = yield* service.taskStatus(scope, taskId).pipe(Effect.flip); + assert.equal(error.code, "orchestration_error"); - const result = yield* service.taskStatus(scope, taskId); - assert.equal(result.status, "completed"); - assert.equal(result.summary, "terminal result"); - assert.equal(result.latestTerminalRunId, childRunId); - assert.equal(result.latestTerminalStatus, "completed"); - const commandIds = yield* Ref.get(acknowledgementCommandIds); - assert.equal(commandIds.length, 2); - assert.notEqual(commandIds[0], commandIds[1]); - }).pipe(Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(layerDependencies)))); - }), + const result = yield* service.taskStatus(scope, taskId); + assert.equal(result.status, "completed"); + assert.equal(result.summary, "terminal result"); + assert.equal(result.latestTerminalRunId, childRunId); + assert.equal(result.latestTerminalStatus, "completed"); + const commandIds = yield* Ref.get(acknowledgementCommandIds); + assert.equal(commandIds.length, 2); + assert.notEqual(commandIds[0], commandIds[1]); + }).pipe( + Effect.provide(OrchestratorMcpService.layer.pipe(Layer.provide(layerDependencies))), + ); + }), ); it.effect("reports a restart-cut child as working until its continuation settles", () => @@ -289,6 +319,7 @@ describe("OrchestratorMcpService", () => { messages: [], subagents: [], providerThreads: [], + turnItems: [], } as unknown as OrchestrationV2ThreadProjection; const layerDependencies = Layer.mergeAll( NodeServices.layer, @@ -365,6 +396,7 @@ describe("OrchestratorMcpService", () => { messages: [], subagents: [], providerThreads: [], + turnItems: [], } as unknown as OrchestrationV2ThreadProjection; const layerDependencies = Layer.mergeAll( NodeServices.layer, @@ -444,6 +476,7 @@ describe("OrchestratorMcpService", () => { messages: [], subagents: [], providerThreads: [], + turnItems: [], } as unknown as OrchestrationV2ThreadProjection; const layerDependencies = Layer.mergeAll( NodeServices.layer, @@ -530,6 +563,7 @@ describe("OrchestratorMcpService", () => { messages: [], subagents: [], providerThreads: [], + turnItems: [], } as unknown as OrchestrationV2ThreadProjection; const layerDependencies = Layer.mergeAll( NodeServices.layer, diff --git a/apps/server/src/mcp/OrchestratorMcpService.ts b/apps/server/src/mcp/OrchestratorMcpService.ts index 3d4059d70cef..df4c372223e9 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.ts @@ -1210,8 +1210,12 @@ const make = Effect.gen(function* () { const childControls = yield* threadManagement .getThreadRecords( task.childThreadId, - ["runs", "messages", "contextTransfers", "subagents", "providerThreads"], - { messageRoles: ["user"] }, + ["runs", "messages", "contextTransfers", "subagents", "providerThreads", "turnItems"], + { + messageRoles: ["user"], + turnItemTypes: ["command_execution", "dynamic_tool", "subagent"], + turnItemStatuses: ["pending", "running", "waiting"], + }, ) .pipe(Effect.mapError(threadManagementFailure)); const childRun = delegatedTaskRun(childControls, task); diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 6f1819a672b7..087105dbe966 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -14,6 +14,8 @@ import { ProviderThreadId, RunId, ThreadId, + TurnItemId, + type OrchestrationV2TurnItem, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; import * as Effect from "effect/Effect"; @@ -901,6 +903,125 @@ const seedRestartCancelledChild = (input: { return { taskId, childThreadId, childRunId }; }); +it.layer(layerTest)("delegated background work", (it) => { + it.effect.each(["completed", "interrupted"] as const)( + "delivers the child result only after its background command is %s", + (status) => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const eventSink = yield* EventSink.EventSinkV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make(`thread:background-parent-${status}`); + const projectId = ProjectId.make(`project:background-parent-${status}`); + const runId = RunId.make(`run:background-parent-${status}`); + const rootNodeId = NodeId.make(`node:background-parent-${status}`); + yield* seedParentWithTerminalTask({ + threadId, + projectId, + runId, + rootNodeId, + taskId: NodeId.make(`node:background-settled-${status}`), + deliveryState: "delivered", + now, + }); + const child = yield* seedRestartCancelledChild({ + parentThreadId: threadId, + projectId, + parentRunId: runId, + rootNodeId, + name: `background-child-${status}`, + completionWake: "always", + continuationPending: false, + runStatus: "completed", + now, + }); + const command: OrchestrationV2TurnItem = { + id: TurnItemId.make(`background-command-${status}`), + threadId: child.childThreadId, + runId: child.childRunId, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + type: "command_execution", + status: "running", + title: null, + startedAt: now, + completedAt: null, + updatedAt: now, + input: "test --watch", + }; + yield* eventSink.write({ + commandId: CommandId.make(`command:background-start-${status}`), + events: [ + { + id: EventId.make(`event:background-start-${status}`), + type: "turn-item.updated", + threadId: child.childThreadId, + occurredAt: now, + payload: command, + }, + ], + }); + // Use the same finalization path synchronously so the negative assertion + // cannot pass merely because the terminal-run listener has not run yet. + yield* orchestrator.recoverDelegatedTask(child.childThreadId, child.childRunId); + assert.isTrue(yield* orchestrator.delegatedTaskResultPending(child.childThreadId)); + const pending = yield* orchestrator.getThreadProjection(threadId); + const task = pending.subagents.find((row) => row.id === child.taskId); + assert.equal(task?.status, "running"); + assert.isNull(task?.result); + assert.isUndefined(task?.completionDelivery); + assert.isFalse( + pending.contextTransfers.some( + (transfer) => transfer.sourceThreadId === child.childThreadId, + ), + ); + + const afterSequence = yield* eventSink.latestSequence(); + // No new run event follows this update. The item ending must trigger delivery. + yield* eventSink.write({ + commandId: CommandId.make(`command:background-end-${status}`), + events: [ + { + id: EventId.make(`event:background-end-${status}`), + type: "turn-item.updated", + threadId: child.childThreadId, + occurredAt: now, + payload: { ...command, status, completedAt: now }, + }, + ], + }); + yield* eventSink.stream({ afterSequence, eventType: "subagent.updated" }).pipe( + Stream.filter( + (stored) => + stored.event.type === "subagent.updated" && + stored.event.payload.id === child.taskId && + stored.event.payload.status === "completed", + ), + Stream.take(1), + Stream.runDrain, + ); + const finished = yield* orchestrator.getThreadProjection(threadId); + const result = finished.subagents.find((row) => row.id === child.taskId); + assert.equal(result?.status, "completed"); + assert.equal(result?.result, "Child task completed without an assistant result."); + assert.equal(result?.completionDelivery?.state, "claimed"); + assert.isFalse(yield* orchestrator.delegatedTaskResultPending(child.childThreadId)); + assert.equal( + finished.contextTransfers.filter( + (transfer) => + transfer.sourceThreadId === child.childThreadId && + transfer.type === "subagent_result", + ).length, + 1, + ); + }), + ); +}); + it.layer(layerTest)("delegated tasks across a server restart", (it) => { it.effect("holds a restart-cancelled child for its continuation's result", () => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 5d4b1cfda15f..37ad5034a204 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -58,6 +58,7 @@ import { modelSelectionsEqual } from "@t3tools/shared/model"; import { derivePendingBackgroundWork, pendingBackgroundTurnItems, + turnItemUpdateCanEndBackgroundWork, } from "@t3tools/shared/orchestrationV2PendingBackgroundWork"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; @@ -9563,8 +9564,20 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio Effect.gen(function* () { const childControls = yield* projectionStore.getThreadRecords( childThreadId, - ["runs", "messages", "subagents", "providerThreads", "providerTurns", "attempts"], - { messageRoles: ["user"] }, + [ + "runs", + "messages", + "subagents", + "providerThreads", + "providerTurns", + "attempts", + "turnItems", + ], + { + messageRoles: ["user"], + turnItemTypes: ["command_execution", "dynamic_tool", "subagent"], + turnItemStatuses: ["pending", "running", "waiting"], + }, ); const forkedFrom = childControls.thread.forkedFrom; if ( @@ -10665,6 +10678,37 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio Effect.forkDetach, ); + // Background commands can end after the last run update. Keep draining item + // output while a parent's lock is busy, retaining only settlement identifiers. + yield* eventSink + .stream({ afterSequence: terminalEventsAfterSequence, eventType: "turn-item.updated" }) + .pipe( + Stream.filter( + (stored) => + stored.event.type === "turn-item.updated" && + !String(stored.commandId).startsWith("command:runtime-reconcile:") && + turnItemUpdateCanEndBackgroundWork(stored.event.payload), + ), + Stream.map((stored) => ({ threadId: stored.event.threadId, sequence: stored.sequence })), + Stream.buffer({ capacity: "unbounded" }), + Stream.runForEach(({ threadId, sequence }) => + Effect.gen(function* () { + const parentThreadId = yield* appOwnedSubagentParentThreadId(threadId); + if (parentThreadId === undefined) return; + yield* threadDispatch.withLock(parentThreadId, finalizeAppOwnedSubagent(threadId)); + }).pipe( + Effect.catchCause((cause) => + Effect.logWarning("Failed to settle app-owned subagent after background work", { + threadId, + sequence, + cause, + }), + ), + ), + ), + Effect.forkDetach, + ); + // Settles child results and completion deliveries whose runs ended without // the listener above: before this boot, or in runtime reconciliation, which // it skips. Startup runs this after reconciliation and before the effect @@ -10763,8 +10807,20 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio Effect.gen(function* () { const child = yield* projectionStore.getThreadRecords( childThreadId, - ["runs", "messages", "subagents", "providerThreads", "providerTurns", "attempts"], - { messageRoles: ["user"] }, + [ + "runs", + "messages", + "subagents", + "providerThreads", + "providerTurns", + "attempts", + "turnItems", + ], + { + messageRoles: ["user"], + turnItemTypes: ["command_execution", "dynamic_tool", "subagent"], + turnItemStatuses: ["pending", "running", "waiting"], + }, ); const progress = delegatedTaskProgress(child); // A caller's older read saw a result; newer work since then means it is not final. diff --git a/apps/server/src/orchestration-v2/SubagentProjection.test.ts b/apps/server/src/orchestration-v2/SubagentProjection.test.ts index 23d180102cd3..dbad2570798a 100644 --- a/apps/server/src/orchestration-v2/SubagentProjection.test.ts +++ b/apps/server/src/orchestration-v2/SubagentProjection.test.ts @@ -8,6 +8,7 @@ import { EventId, TurnItemId, type OrchestrationV2Run, + type OrchestrationV2TurnItem, type OrchestrationV2AppThread, ProjectId, ProviderInstanceId, @@ -271,6 +272,58 @@ it("waits for nested work and retains the report across monitor acknowledgements assert.equal(progress.resultRun?.id, report.id); }); +it.each([ + ["pending", "waiting_for_children"], + ["running", "waiting_for_children"], + ["waiting", "waiting_for_children"], + ["completed", "result_available"], + ["failed", "result_available"], + ["interrupted", "result_available"], + ["cancelled", "result_available"], +] as const)("reports a %s background command as %s", (status, expected) => { + const { projection, run } = taskFixture(); + const command: OrchestrationV2TurnItem = { + id: TurnItemId.make("background-command"), + threadId: parentThreadId, + runId: run.id, + nodeId: null, + providerThreadId: null, + providerTurnId: null, + nativeItemRef: null, + parentItemId: null, + ordinal: 1, + type: "command_execution", + status, + title: null, + startedAt: parentCreatedAt, + completedAt: null, + updatedAt: childCreatedAt, + input: "test --watch", + }; + assert.equal(delegatedTaskProgress({ ...projection, turnItems: [command] }).state, expected); + assert.equal( + delegatedTaskProgress({ ...projection, turnItems: [{ ...command, runId: null }] }).state, + expected, + ); + assert.equal( + delegatedTaskProgress({ + ...projection, + runs: [run, { ...run, id: RunId.make("abandoned"), ordinal: 2, status: "rolled_back" }], + turnItems: [{ ...command, runId: RunId.make("abandoned") }], + }).state, + "result_available", + ); + assert.equal( + delegatedTaskProgress({ + ...projection, + turnItems: [ + { ...command, type: "dynamic_tool", toolName: "monitor", input: { persistent: true } }, + ], + }).state, + "result_available", + ); +}); + it("exposes the provider failure rather than a progress message from the failed run", () => { const { projection, run } = taskFixture(); const failedRun = { ...run, status: "failed" as const }; diff --git a/apps/server/src/orchestration-v2/SubagentProjection.ts b/apps/server/src/orchestration-v2/SubagentProjection.ts index f28dc1aa4be5..54333f7bdd34 100644 --- a/apps/server/src/orchestration-v2/SubagentProjection.ts +++ b/apps/server/src/orchestration-v2/SubagentProjection.ts @@ -17,6 +17,7 @@ import type { TurnItemId, } from "@t3tools/contracts"; import { runRanAfter } from "@t3tools/shared/orchestrationV2ThreadError"; +import { pendingBackgroundTurnItems } from "@t3tools/shared/orchestrationV2PendingBackgroundWork"; import * as DateTime from "effect/DateTime"; import { isOrchestrationV2WorkActive } from "@t3tools/contracts"; @@ -206,9 +207,10 @@ export function subagentResultForRun( }; } -/** A finished turn can still own live children or queued completion follow-ups. */ +/** A finished turn can still own background work or queued completion follow-ups. */ export function delegatedTaskProgress(projection: { readonly runs: OrchestrationV2ThreadProjection["runs"]; + readonly turnItems: OrchestrationV2ThreadProjection["turnItems"]; readonly messages: ReadonlyArray< Pick >; @@ -239,7 +241,8 @@ export function delegatedTaskProgress(projection: { task.completionDelivery?.state === "pending" || task.completionDelivery?.state === "claimed", ) || - projection.providerThreads.some((thread) => (thread.pendingBackgroundTasks?.length ?? 0) > 0); + projection.providerThreads.some((thread) => (thread.pendingBackgroundTasks?.length ?? 0) > 0) || + pendingBackgroundTurnItems(projection).length > 0; const resultRun = workRuns .filter((run) => terminal(run.status) && (run.startedAt !== null || run.ordinal === 1)) .toSorted((a, b) => (runRanAfter(a, b) ? -1 : runRanAfter(b, a) ? 1 : 0))[0]; From 76998b80bffdc15b1798ba6a11e11b87f4f12d91 Mon Sep 17 00:00:00 2001 From: ashx-j <281260375+ashx-j@users.noreply.github.com> Date: Wed, 7 Oct 2026 11:36:12 +0200 Subject: [PATCH 2/4] fix(server): coalesce background settlement rechecks by thread --- .../DelegatedCompletionDelivery.test.ts | 39 ++++++++++++++----- .../src/orchestration-v2/Orchestrator.ts | 11 +++++- 2 files changed, 39 insertions(+), 11 deletions(-) diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 087105dbe966..73492685b20e 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -953,17 +953,20 @@ it.layer(layerTest)("delegated background work", (it) => { updatedAt: now, input: "test --watch", }; + const lastCommand = { + ...command, + id: TurnItemId.make(`last-command-${status}`), + ordinal: 2, + }; yield* eventSink.write({ commandId: CommandId.make(`command:background-start-${status}`), - events: [ - { - id: EventId.make(`event:background-start-${status}`), - type: "turn-item.updated", - threadId: child.childThreadId, - occurredAt: now, - payload: command, - }, - ], + events: [command, lastCommand].map((payload) => ({ + id: EventId.make(`event:start-${payload.id}`), + type: "turn-item.updated" as const, + threadId: child.childThreadId, + occurredAt: now, + payload, + })), }); // Use the same finalization path synchronously so the negative assertion // cannot pass merely because the terminal-run listener has not run yet. @@ -980,6 +983,22 @@ it.layer(layerTest)("delegated background work", (it) => { ), ); + // Repeated completion signals may coalesce, but the other command still holds the result. + yield* eventSink.write({ + commandId: CommandId.make(`command:background-repeat-${status}`), + events: Array.from({ length: 20 }, (_, index) => ({ + id: EventId.make(`event:background-repeat-${status}-${index}`), + type: "turn-item.updated" as const, + threadId: child.childThreadId, + occurredAt: now, + payload: { ...command, status: "completed" as const, completedAt: now }, + })), + }); + yield* orchestrator.recoverDelegatedTask(child.childThreadId, child.childRunId); + assert.isTrue(yield* orchestrator.delegatedTaskResultPending(child.childThreadId)); + const stillPending = yield* orchestrator.getThreadProjection(threadId); + assert.isNull(stillPending.subagents.find((row) => row.id === child.taskId)?.result); + const afterSequence = yield* eventSink.latestSequence(); // No new run event follows this update. The item ending must trigger delivery. yield* eventSink.write({ @@ -990,7 +1009,7 @@ it.layer(layerTest)("delegated background work", (it) => { type: "turn-item.updated", threadId: child.childThreadId, occurredAt: now, - payload: { ...command, status, completedAt: now }, + payload: { ...lastCommand, status, completedAt: now }, }, ], }); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 37ad5034a204..40e3d3585fff 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -10679,7 +10679,9 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ); // Background commands can end after the last run update. Keep draining item - // output while a parent's lock is busy, retaining only settlement identifiers. + // output while a parent's lock is busy, with at most one queued recheck per + // thread. Finalization reads current state, so repeated signals can coalesce. + const pendingBackgroundSettlements = new Set(); yield* eventSink .stream({ afterSequence: terminalEventsAfterSequence, eventType: "turn-item.updated" }) .pipe( @@ -10690,9 +10692,16 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio turnItemUpdateCanEndBackgroundWork(stored.event.payload), ), Stream.map((stored) => ({ threadId: stored.event.threadId, sequence: stored.sequence })), + Stream.filter(({ threadId }) => { + if (pendingBackgroundSettlements.has(threadId)) return false; + pendingBackgroundSettlements.add(threadId); + return true; + }), Stream.buffer({ capacity: "unbounded" }), Stream.runForEach(({ threadId, sequence }) => Effect.gen(function* () { + // A signal arriving during this recheck must queue another one. + pendingBackgroundSettlements.delete(threadId); const parentThreadId = yield* appOwnedSubagentParentThreadId(threadId); if (parentThreadId === undefined) return; yield* threadDispatch.withLock(parentThreadId, finalizeAppOwnedSubagent(threadId)); From f102b3b9cbf5feaa7c064939f40c1a5a461988b1 Mon Sep 17 00:00:00 2001 From: ashx-j <281260375+ashx-j@users.noreply.github.com> Date: Wed, 7 Oct 2026 11:37:45 +0200 Subject: [PATCH 3/4] test(server): fix generic record lookup in task status fixture --- apps/server/src/mcp/OrchestratorMcpService.test.ts | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/apps/server/src/mcp/OrchestratorMcpService.test.ts b/apps/server/src/mcp/OrchestratorMcpService.test.ts index 2c73db802c12..b84f3a22c75c 100644 --- a/apps/server/src/mcp/OrchestratorMcpService.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpService.test.ts @@ -118,7 +118,8 @@ describe("OrchestratorMcpService", () => { ? [{ ...parentProjection.subagents[0]!, status: "running" as const }] : [], turnItems: - backgroundWork === "background command" && records.includes("turnItems") + backgroundWork === "background command" && + records.some((record) => record === "turnItems") ? [ { id: TurnItemId.make("background-command"), From bf24eb53ef138074183e21d6f01b6fd8afe804f6 Mon Sep 17 00:00:00 2001 From: ashx-j <281260375+ashx-j@users.noreply.github.com> Date: Wed, 7 Oct 2026 11:53:51 +0200 Subject: [PATCH 4/4] perf(server): reuse lineage reads for background completions --- .../DelegatedCompletionDelivery.test.ts | 49 +++++++++++++++++++ .../src/orchestration-v2/Orchestrator.ts | 10 +++- 2 files changed, 58 insertions(+), 1 deletion(-) diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 73492685b20e..a4337246a318 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -1,6 +1,7 @@ import * as SourceControlProviderRegistry from "../sourceControl/SourceControlProviderRegistry.ts"; import * as NodeServices from "@effect/platform-node/NodeServices"; import { assert, it } from "@effect/vitest"; +import { vi } from "vite-plus/test"; import { CommandId, EventId, @@ -18,6 +19,7 @@ import { type OrchestrationV2TurnItem, } from "@t3tools/contracts"; import * as DateTime from "effect/DateTime"; +import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Stream from "effect/Stream"; @@ -37,6 +39,7 @@ 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 ProjectionStore from "./ProjectionStore.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; import { continueRestartedRun } from "./RestartContinuation.ts"; import * as RuntimeLayer from "./runtimeLayer.ts"; @@ -104,6 +107,7 @@ const layerTestProviderInstanceRegistry = Layer.succeed( ); const layerTest = Layer.mergeAll(RuntimeLayer.layer, RuntimeLayer.layerEventSink).pipe( + Layer.provideMerge(ProjectionStore.layer), Layer.provideMerge(RuntimeLayer.layerProjectService), Layer.provide( Layer.mock(WorkspacePaths.WorkspacePaths)({ @@ -999,11 +1003,55 @@ it.layer(layerTest)("delegated background work", (it) => { const stillPending = yield* orchestrator.getThreadProjection(threadId); assert.isNull(stillPending.subagents.find((row) => row.id === child.taskId)?.result); + const projections = yield* ProjectionStore.ProjectionStoreV2; + const firstOrdinaryRead = yield* Deferred.make(); + const getThread = projections.getThread; + let ordinaryReads = 0; + const readSpy = vi.spyOn(projections, "getThread").mockImplementation((id) => + getThread(id).pipe( + Effect.tap(() => { + if (id !== threadId) return Effect.void; + ordinaryReads += 1; + return Deferred.succeed(firstOrdinaryRead, undefined); + }), + ), + ); + yield* Effect.addFinalizer(() => Effect.sync(() => readSpy.mockRestore())); + const ordinaryCompletion = { + ...command, + id: TurnItemId.make(`ordinary-command-${status}`), + threadId, + runId, + status: "completed" as const, + completedAt: now, + }; + yield* eventSink.write({ + commandId: CommandId.make(`command:ordinary-first-${status}`), + events: [ + { + id: EventId.make(`event:ordinary-first-${status}`), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: ordinaryCompletion, + }, + ], + }); + // The first lookup means its queued key was already cleared. The next + // completion must use the lineage cache, not just queue coalescing. + yield* Deferred.await(firstOrdinaryRead); const afterSequence = yield* eventSink.latestSequence(); // No new run event follows this update. The item ending must trigger delivery. yield* eventSink.write({ commandId: CommandId.make(`command:background-end-${status}`), events: [ + { + id: EventId.make(`event:ordinary-second-${status}`), + type: "turn-item.updated", + threadId, + occurredAt: now, + payload: ordinaryCompletion, + }, { id: EventId.make(`event:background-end-${status}`), type: "turn-item.updated", @@ -1024,6 +1072,7 @@ it.layer(layerTest)("delegated background work", (it) => { Stream.runDrain, ); const finished = yield* orchestrator.getThreadProjection(threadId); + assert.equal(ordinaryReads, 1); const result = finished.subagents.find((row) => row.id === child.taskId); assert.equal(result?.status, "completed"); assert.equal(result?.result, "Child task completed without an assistant result."); diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 40e3d3585fff..dbec6247cb87 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -60,8 +60,10 @@ import { pendingBackgroundTurnItems, turnItemUpdateCanEndBackgroundWork, } from "@t3tools/shared/orchestrationV2PendingBackgroundWork"; +import * as Cache from "effect/Cache"; import * as Context from "effect/Context"; import * as DateTime from "effect/DateTime"; +import * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as FileSystem from "effect/FileSystem"; @@ -10682,6 +10684,12 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio // output while a parent's lock is busy, with at most one queued recheck per // thread. Finalization reads current state, so repeated signals can coalesce. const pendingBackgroundSettlements = new Set(); + // Lineage is immutable. Cache negative lookups too so ordinary command + // completions do not each read a thread record; failed reads remain retryable. + const backgroundSettlementParents = yield* Cache.makeWith(appOwnedSubagentParentThreadId, { + capacity: 1024, + timeToLive: (exit) => (Exit.isSuccess(exit) ? Duration.infinity : Duration.zero), + }); yield* eventSink .stream({ afterSequence: terminalEventsAfterSequence, eventType: "turn-item.updated" }) .pipe( @@ -10702,7 +10710,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio Effect.gen(function* () { // A signal arriving during this recheck must queue another one. pendingBackgroundSettlements.delete(threadId); - const parentThreadId = yield* appOwnedSubagentParentThreadId(threadId); + const parentThreadId = yield* Cache.get(backgroundSettlementParents, threadId); if (parentThreadId === undefined) return; yield* threadDispatch.withLock(parentThreadId, finalizeAppOwnedSubagent(threadId)); }).pipe(