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), });