From e1ff11d1e38cae52a7633b01da4258f81ae5e8f0 Mon Sep 17 00:00:00 2001 From: "github-actions[bot]" <41898282+github-actions[bot]@users.noreply.github.com> Date: Thu, 24 Sep 2026 11:27:56 +1000 Subject: [PATCH] fix(server): stop stranding late delegated task reports A parent-run cohort allowed one completion delivery plus one successor. Any child that finished after the successor started was marked pending and never delivered, so the parent kept working without learning that the child had failed or finished. Coalescing already keeps at most one outstanding delivery per cohort, so drop the lifetime cap and let each settled delivery reserve the next one for results that arrived while it ran. The cap also bounded retries of a cancelled delivery, so bound that directly: a cancelled batch is retried once, and after a second cancellation it waits for a fresh child result instead of re-arming the parent with the same results. Co-Authored-By: Claude Opus 5.5 (1M context) --- ...OrchestratorMcpToolkit.integration.test.ts | 63 ++++- .../DelegatedCompletionDelivery.test.ts | 251 +++++++++++++++++- .../src/orchestration-v2/Orchestrator.ts | 60 +++-- packages/contracts/src/orchestrationV2.ts | 4 +- 4 files changed, 338 insertions(+), 40 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index ed32fb6d89b7..fb7ecc698625 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -2813,9 +2813,9 @@ describe("orchestrator MCP toolkit", () => { return yield* Effect.die(new Error("Late completion successor delivery missing.")); } - // The bounded successor can coalesce a late result only once. A - // third terminal after that successor has started remains - // inspectable, but cannot recursively create a third parent run. + // A terminal that lands after the successor has started is not + // stranded: once the successor settles, the cohort reserves one + // more delivery for it rather than leaving it pending forever. const successorGate = yield* Deferred.make(); deliveryTerminalGates.set(lateParentThreadId, successorGate); yield* orchestrator.dispatch({ @@ -2884,7 +2884,7 @@ describe("orchestrator MCP toolkit", () => { commandId: CommandId.make("command:mcp-late-parent:interrupt-third-late-child"), threadId: thirdLateTask.childThreadId, runId: thirdLateChildRun.id, - reason: "Terminalize after the bounded successor started.", + reason: "Terminalize after the successor delivery started.", }); yield* waitForProjection( orchestrator, @@ -2894,7 +2894,7 @@ describe("orchestrator MCP toolkit", () => { ?.completionDelivery?.state === "pending", ); yield* Deferred.succeed(successorGate, undefined); - const exhaustedCohort = yield* waitForProjection( + const thirdReserved = yield* waitForProjection( orchestrator, lateParentThreadId, (projection) => { @@ -2903,18 +2903,61 @@ describe("orchestrator MCP toolkit", () => { )?.delegatedCompletion; return ( cohort?.settledDeliveryCount === 2 && - cohort.delivery === null && + cohort.delivery?.generation === successorDelivery.generation + 1 && projection.runs.find((run) => run.id === activeSuccessorRun.id)?.status === "completed" && projection.subagents.find((task) => task.id === thirdLateTask.id) - ?.completionDelivery?.state === "pending" + ?.completionDelivery?.state === "claimed" + ); + }, + ); + const thirdDelivery = thirdReserved.runs.find((run) => run.id === lateParentRun.id) + ?.delegatedCompletion?.delivery; + if (thirdDelivery === undefined || thirdDelivery === null) { + return yield* Effect.die(new Error("Third late completion delivery missing.")); + } + expect(thirdDelivery.taskIds).toEqual([thirdLateTask.id]); + expect(yield* waitForContinuationOffers(3)).toHaveLength(3); + + // Settle the third delivery so the cohort is drained before the + // Queue Remove scenario below starts its own cohort. + deliveryTerminalGates.delete(lateParentThreadId); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "agent", + creationSource: "server", + commandId: CommandId.make("command:mcp-late-parent:dispatch-third-delivery"), + threadId: lateParentThreadId, + messageId: thirdDelivery.messageId, + text: "Delegated task reached a terminal state.", + attachments: [], + modelSelection: codexSelection, + dispatchMode: { type: "queue_after_active" }, + delegatedCompletion: { + parentRunId: lateParentRun.id, + generation: thirdDelivery.generation, + taskIds: thirdDelivery.taskIds, + }, + }); + const drainedCohort = yield* waitForProjection( + orchestrator, + lateParentThreadId, + (projection) => { + const cohort = projection.runs.find( + (run) => run.id === lateParentRun.id, + )?.delegatedCompletion; + return ( + cohort?.settledDeliveryCount === 3 && + cohort.delivery === null && + projection.subagents.find((task) => task.id === thirdLateTask.id) + ?.completionDelivery?.state === "delivered" ); }, ); expect( - exhaustedCohort.runs.find((run) => run.id === lateParentRun.id)?.delegatedCompletion, - ).toMatchObject({ settledDeliveryCount: 2, delivery: null }); - yield* expectOffersToStay(2); + drainedCohort.runs.find((run) => run.id === lateParentRun.id)?.delegatedCompletion, + ).toMatchObject({ settledDeliveryCount: 3, delivery: null }); + yield* expectOffersToStay(3); // Queue Remove is a durable disposal action, not a local queue // edit. Start a fresh parent-run cohort so removing this delivery diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 567fea60f7c4..ba9659b9a5e7 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -136,6 +136,8 @@ const seedParentWithTerminalTask = (input: { readonly deliveryState: "delivered" | "claimed" | "acknowledged" | "disposed"; readonly completionWake?: "always" | "settled_only"; readonly deliveryTaskIds?: ReadonlyArray; + readonly settledDeliveryCount?: number; + readonly deliveryGeneration?: number; readonly now: DateTime.Utc; }) => Effect.gen(function* () { @@ -230,13 +232,13 @@ const seedParentWithTerminalTask = (input: { contextHandoffId: null, delegatedCompletion: { disposition: "open", - nextGeneration: 2, - settledDeliveryCount: 1, + nextGeneration: (input.deliveryGeneration ?? 1) + 1, + settledDeliveryCount: input.settledDeliveryCount ?? 1, delivery: input.deliveryTaskIds === undefined ? null : { - generation: 1, + generation: input.deliveryGeneration ?? 1, messageId: MessageId.make(`message:delegated-delivery:${input.threadId}`), taskIds: input.deliveryTaskIds, }, @@ -283,6 +285,87 @@ const seedParentWithTerminalTask = (input: { }); }); +// Writes a cancelled delivery run for `generation` and, unless `settles` is +// false, waits for the terminal-run reactor to settle the parent cohort. +const cancelDeliveryRun = (input: { + readonly threadId: ThreadId; + readonly runId: RunId; + readonly generation: number; + readonly messageId: MessageId; + readonly taskIds: ReadonlyArray; + readonly settles?: boolean; + readonly now: DateTime.Utc; +}) => + Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + const sink = yield* EventSinkV2; + const { threadId, runId, generation, now } = input; + const parentRun = (yield* orchestrator.getThreadProjection(threadId)).runs.find( + (row) => row.id === runId, + )!; + const afterSequence = yield* sink.latestSequence({ threadId }); + const deliveryRunId = RunId.make(`${runId}:delivery:${generation}`); + yield* sink.write({ + commandId: CommandId.make(`command:cancel-delivery:${threadId}:${generation}`), + events: [ + { + id: EventId.make(`event:delivery-message:${threadId}:${generation}`), + type: "message.updated", + threadId, + runId: deliveryRunId, + occurredAt: now, + payload: { + id: input.messageId, + threadId, + runId: deliveryRunId, + nodeId: null, + role: "user", + text: "Background task finished", + attachments: [], + streaming: false, + createdBy: "agent", + creationSource: "server", + createdAt: now, + updatedAt: now, + delegatedCompletion: { parentRunId: runId, generation, taskIds: input.taskIds }, + }, + }, + { + id: EventId.make(`event:delivery-run:${threadId}:${generation}`), + type: "run.updated", + threadId, + runId: deliveryRunId, + providerInstanceId: modelSelection.instanceId, + occurredAt: now, + payload: { + ...parentRun, + id: deliveryRunId, + ordinal: generation + 1, + userMessageId: input.messageId, + rootNodeId: null, + status: "cancelled", + completedAt: now, + delegatedCompletion: undefined, + }, + }, + ], + }); + if (input.settles === false) return undefined; + const settledDeliveryCount = parentRun.delegatedCompletion?.settledDeliveryCount ?? 0; + const settled = yield* sink.stream({ threadId, afterSequence, eventType: "run.updated" }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && + stored.event.payload.id === runId && + stored.event.payload.delegatedCompletion?.settledDeliveryCount === + settledDeliveryCount + 1, + ), + Stream.runHead, + ); + assert.isTrue(settled._tag === "Some"); + return yield* orchestrator.getThreadProjection(threadId); + }); + it.layer(TestLayer)("delegated completion delivery repairs", (it) => { it.effect("acceptance batches pending siblings without acknowledging their results", () => Effect.gen(function* () { @@ -381,6 +464,168 @@ it.layer(TestLayer)("delegated completion delivery repairs", (it) => { }), ); + it.effect("keeps delivering late siblings after earlier deliveries settled", () => + Effect.gen(function* () { + const orchestrator = yield* OrchestratorV2; + const sink = yield* EventSinkV2; + const now = yield* DateTime.now; + const threadId = ThreadId.make("mailbox-late-batch"); + const runId = RunId.make("mailbox-late-parent"); + const taskId = NodeId.make("mailbox-late-first"); + const lateTaskId = NodeId.make("mailbox-late-second"); + const messageId = MessageId.make(`message:delegated-delivery:${threadId}`); + yield* seedParentWithTerminalTask({ + threadId, + runId, + projectId: ProjectId.make("mailbox-late-project"), + rootNodeId: NodeId.make("mailbox-late-root"), + taskId, + deliveryState: "claimed", + completionWake: "always", + deliveryTaskIds: [taskId], + settledDeliveryCount: 2, + now, + }); + const task = (yield* orchestrator.getThreadProjection(threadId)).subagents[0]!; + yield* sink.write({ + events: [ + { + id: EventId.make("mailbox-late-message"), + type: "message.updated", + threadId, + runId, + occurredAt: now, + payload: { + id: messageId, + threadId, + runId, + nodeId: task.parentNodeId, + role: "user", + text: "Background task finished", + attachments: [], + streaming: false, + createdBy: "agent", + creationSource: "server", + createdAt: now, + updatedAt: now, + delegatedCompletion: { parentRunId: runId, generation: 1, taskIds: [taskId] }, + }, + }, + { + id: EventId.make(`event:${lateTaskId}`), + type: "subagent.updated", + threadId, + runId, + nodeId: lateTaskId, + occurredAt: now, + payload: { + ...task, + id: lateTaskId, + status: "failed", + result: "Claude API rate limit reached. Try again later.", + completionDelivery: { state: "pending", observedByRunId: null }, + }, + }, + ], + }); + yield* orchestrator.dispatch({ + type: "notification.delivery.accept", + commandId: CommandId.make("accept-late-first"), + threadId, + messageId, + }); + const accepted = yield* orchestrator.getThreadProjection(threadId); + const cohort = accepted.runs.find((row) => row.id === runId)?.delegatedCompletion; + assert.deepEqual(cohort?.delivery?.taskIds, [lateTaskId]); + assert.deepEqual( + accepted.subagents.find((row) => row.id === lateTaskId)?.completionDelivery, + { + state: "claimed", + observedByRunId: null, + }, + ); + }), + ); + + it.effect("retries a cancelled delivery once without re-arming the parent forever", () => + Effect.gen(function* () { + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:delegated-delivery-cancel-loop"); + const runId = RunId.make("run:delegated-delivery-cancel-loop"); + const taskId = NodeId.make("node:delegated-delivery-cancel-loop-task"); + yield* seedParentWithTerminalTask({ + threadId, + runId, + projectId: ProjectId.make("project:delegated-delivery-cancel-loop"), + rootNodeId: NodeId.make("node:delegated-delivery-cancel-loop-root"), + taskId, + deliveryState: "claimed", + completionWake: "always", + deliveryTaskIds: [taskId], + settledDeliveryCount: 0, + now, + }); + const cancel = (generation: number, messageId: MessageId) => + cancelDeliveryRun({ threadId, runId, generation, messageId, taskIds: [taskId], now }); + + const retried = yield* cancel(1, MessageId.make(`message:delegated-delivery:${threadId}`)); + const retry = retried!.runs.find((row) => row.id === runId)?.delegatedCompletion?.delivery; + assert.equal(retry?.generation, 2); + assert.deepEqual(retry?.taskIds, [taskId]); + + const exhausted = yield* cancel(2, retry!.messageId); + assert.isNull(exhausted!.runs.find((row) => row.id === runId)?.delegatedCompletion?.delivery); + assert.deepEqual(exhausted!.subagents.find((row) => row.id === taskId)?.completionDelivery, { + state: "pending", + observedByRunId: null, + }); + }), + ); + + it.effect("retries a fresh batch even after an unrelated delivery was cancelled", () => + Effect.gen(function* () { + const now = yield* DateTime.now; + const threadId = ThreadId.make("thread:delegated-delivery-unrelated-cancel"); + const runId = RunId.make("run:delegated-delivery-unrelated-cancel"); + const taskId = NodeId.make("node:delegated-delivery-unrelated-cancel-task"); + yield* seedParentWithTerminalTask({ + threadId, + runId, + projectId: ProjectId.make("project:delegated-delivery-unrelated-cancel"), + rootNodeId: NodeId.make("node:delegated-delivery-unrelated-cancel-root"), + taskId, + deliveryState: "claimed", + completionWake: "always", + deliveryTaskIds: [taskId], + deliveryGeneration: 2, + settledDeliveryCount: 0, + now, + }); + // An earlier delivery for a different task was cancelled after the + // parent acknowledged it, so it no longer owns the cohort. + yield* cancelDeliveryRun({ + threadId, + runId, + generation: 1, + messageId: MessageId.make(`message:delegated-delivery-stale:${threadId}`), + taskIds: [NodeId.make("node:delegated-delivery-unrelated-cancel-stale")], + settles: false, + now, + }); + const retried = yield* cancelDeliveryRun({ + threadId, + runId, + generation: 2, + messageId: MessageId.make(`message:delegated-delivery:${threadId}`), + taskIds: [taskId], + now, + }); + const retry = retried!.runs.find((row) => row.id === runId)?.delegatedCompletion?.delivery; + assert.equal(retry?.generation, 3); + assert.deepEqual(retry?.taskIds, [taskId]); + }), + ); + it.effect("builds completion text and metadata from the same live cohort", () => Effect.gen(function* () { const orchestrator = yield* OrchestratorV2; diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 05565a40e2eb..a29f5799c3fb 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -8011,8 +8011,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } if (deliveryRun !== undefined) { // The terminal-run listener owns reconciliation of a completed wake. - // A sibling that wins the parent lock first remains pending for its - // one successor rather than creating a competing delivery. + // A sibling that wins the parent lock first remains pending for the + // next delivery rather than creating a competing one. return { task: { ...input.updatedTask, @@ -8051,24 +8051,11 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }; } + // Coalescing bounds parent wakes: a cohort holds at most one outstanding + // delivery, and siblings that finish while it runs join the next one. + // There is no lifetime cap; finalizeDelegatedCompletionDelivery bounds + // retries of a cancelled delivery instead. const settledDeliveryCount = cohort?.settledDeliveryCount ?? 0; - if (settledDeliveryCount >= 2) { - // A cohort permits one initial delivery and one successor. Keep the - // result pending and inspectable instead of recursively re-arming the - // parent for every child that finishes after that bounded handoff. - return { - task: { - ...input.updatedTask, - completionDelivery: { - state: "pending" as const, - observedByRunId: null, - }, - }, - parentRun: undefined, - message: undefined, - offer: false, - }; - } const generation = cohort?.nextGeneration ?? 1; const messageId = yield* mapDelegatedCompletionError( idAllocator.allocate.message({ @@ -8461,12 +8448,37 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ) .map((task) => task.id); const settledDeliveryCount = (cohort.settledDeliveryCount ?? 0) + 1; + // A cancelled delivery returns its batch to pending. Retry that batch + // once; if the previous delivery carried the same results and was also + // cancelled, wait for a fresh child result rather than re-arming the + // parent with the same results forever. + const previousDeliveryMessage = projection.messages + .filter( + (candidate) => + candidate.delegatedCompletion?.parentRunId === parentRun.id && + candidate.delegatedCompletion.generation < delivery.generation, + ) + .toSorted( + (left, right) => + (right.delegatedCompletion?.generation ?? 0) - + (left.delegatedCompletion?.generation ?? 0), + )[0]; + const retriedCancelledBatch = + deliveryRun.status === "cancelled" && + previousDeliveryMessage !== undefined && + projection.runs.find((candidate) => candidate.userMessageId === previousDeliveryMessage.id) + ?.status === "cancelled" && + pendingTaskIds.every( + (taskId) => + delivery.taskIds.includes(taskId) && + previousDeliveryMessage.delegatedCompletion?.taskIds.includes(taskId) === true, + ); const canReserveFollowUp = cohort.disposition === "open" && projection.thread.archivedAt === null && projection.thread.deletedAt === null && - settledDeliveryCount < 2 && - pendingTaskIds.length > 0; + pendingTaskIds.length > 0 && + !retriedCancelledBatch; const nextDelivery = canReserveFollowUp ? { generation: cohort.nextGeneration, @@ -8568,9 +8580,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } const parentIsLive = hasLiveRun(projection); const pendingTaskIds = - projection.thread.archivedAt === null && - projection.thread.deletedAt === null && - (cohort.settledDeliveryCount ?? 0) < 2 + projection.thread.archivedAt === null && projection.thread.deletedAt === null ? projection.subagents .filter( (task) => @@ -8619,7 +8629,7 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); } // Provider acceptance drains this batch but does not acknowledge its results - // or spend an idle-wake allowance. task_status owns acknowledgment. + // or count as a settled delivery run. task_status owns acknowledgment. yield* emitEvent({ type: "run.updated", threadId: command.threadId, diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 3ba1244729bd..16209a7dd748 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -460,8 +460,8 @@ export type OrchestrationV2DelegatedCompletionDelivery = export const OrchestrationV2DelegatedCompletionCohort = Schema.Struct({ disposition: Schema.Literals(["open", "stopped", "disposed"]), nextGeneration: PositiveInt, - // Optional for compatibility with cohorts persisted before bounded - // follow-up delivery was introduced. Missing means no delivery has settled. + // Delivery runs this cohort has settled; provider mailbox acceptance is not + // counted. Optional for cohorts persisted before follow-up delivery. settledDeliveryCount: Schema.optional(NonNegativeInt), delivery: Schema.NullOr(OrchestrationV2DelegatedCompletionDelivery), });