From 78d50b968c8d271e43c282ac0bf922f389ad192c Mon Sep 17 00:00:00 2001 From: Julius Marminge <51714798+juliusmarminge@users.noreply.github.com> Date: Sun, 27 Sep 2026 01:01:37 -0700 Subject: [PATCH] fix(server): parents keep getting woken for every delegated task that finishes A delegated-completion cohort allowed one delivery plus one successor, then left every later child result pending with no wake. A parent that fanned out twelve children was woken twice and then waited silently until the user typed. When a delivery settles with results still pending, reserve one successor that carries all of them. Siblings that finish before a delivery starts still join it, and results that land while it is queued or running wait for the next one, so a cohort never has more than one delivery outstanding. Each child becomes pending at most once, so successors are bounded by the cohort's children. The cap was the only reader of settledDeliveryCount, so the field is dropped from the contract; older rows that carry it decode unchanged. Co-Authored-By: Claude Opus 5.5 (1M context) --- ...OrchestratorMcpToolkit.integration.test.ts | 478 ++++++++++++++++-- .../CheckpointCaptureService.test.ts | 2 - .../DelegatedCompletionDelivery.test.ts | 6 - .../src/orchestration-v2/Orchestrator.ts | 39 +- .../SteeringCompletion.integration.test.ts | 2 - .../orchestration-v2/ThreadDeletion.test.ts | 2 - .../src/orchestration-v2/ThreadDeletion.ts | 1 - packages/contracts/src/orchestrationV2.ts | 3 - 8 files changed, 449 insertions(+), 84 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index 3e67bba6e8b6..6b192caa03a9 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -8,6 +8,7 @@ import { isProviderNativeSubagentThread, MessageId, type ModelSelection, + type OrchestrationV2DelegatedCompletionDelivery, type OrchestrationV2ProviderCapabilities, type OrchestrationV2ProviderSession, type OrchestrationV2ProviderThread, @@ -2701,7 +2702,27 @@ describe("orchestrator MCP toolkit", () => { completionWake: "always", }); const thirdLateTask = taskFromDispatch(thirdLateChild); - if (lateTask.childThreadId === null || thirdLateTask.childThreadId === null) { + const fourthLateTask = taskFromDispatch( + yield* orchestrator.dispatch({ + type: "delegated_task.request", + createdBy: "agent", + creationSource: "mcp", + commandId: CommandId.make("command:mcp-late-parent:fourth-late-task"), + parentThreadId: lateParentThreadId, + parentRunId: lateParentRun.id, + parentNodeId: lateParentRun.rootNodeId, + task: cancellationPrompt, + modelSelection: codexSelection, + runtimeMode: "full-access", + interactionMode: "default", + completionWake: "always", + }), + ); + if ( + lateTask.childThreadId === null || + thirdLateTask.childThreadId === null || + fourthLateTask.childThreadId === null + ) { return yield* Effect.die(new Error("Late completion child thread missing.")); } const beforeFirstDelivery = yield* waitForProjection( @@ -2814,9 +2835,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. + // Results that land while the successor runs wait for it to settle, + // then go out together in one more delivery. Nothing is left + // pending once the parent has been told about every child. const successorGate = yield* Deferred.make(); deliveryTerminalGates.set(lateParentThreadId, successorGate); yield* orchestrator.dispatch({ @@ -2867,55 +2888,129 @@ describe("orchestrator MCP toolkit", () => { item.type === "user_message" && item.messageId === successorDelivery.messageId, ), ).toBe(false); - const thirdLateChildProjection = yield* waitForProjection( - orchestrator, - thirdLateTask.childThreadId, - (projection) => - projection.runs.some((run) => run.status === "running") && - projection.providerTurns.some((turn) => turn.status === "running"), - ); - const thirdLateChildRun = thirdLateChildProjection.runs.find( - (run) => run.status === "running", - ); - if (thirdLateChildRun === undefined) { - return yield* Effect.die(new Error("Third late completion child run missing.")); + for (const [label, childThreadId] of [ + ["third", thirdLateTask.childThreadId], + ["fourth", fourthLateTask.childThreadId], + ] as const) { + const childProjection = yield* waitForProjection( + orchestrator, + childThreadId, + (projection) => + projection.runs.some((run) => run.status === "running") && + projection.providerTurns.some((turn) => turn.status === "running"), + ); + const childRun = childProjection.runs.find((run) => run.status === "running"); + if (childRun === undefined) { + return yield* Effect.die(new Error(`${label} late completion child run missing.`)); + } + yield* orchestrator.dispatch({ + type: "run.interrupt", + commandId: CommandId.make(`command:mcp-late-parent:interrupt-${label}-late-child`), + threadId: childThreadId, + runId: childRun.id, + reason: "Terminalize while the successor delivery is running.", + }); } - yield* orchestrator.dispatch({ - type: "run.interrupt", - commandId: CommandId.make("command:mcp-late-parent:interrupt-third-late-child"), - threadId: thirdLateTask.childThreadId, - runId: thirdLateChildRun.id, - reason: "Terminalize after the bounded successor started.", - }); - yield* waitForProjection( + const pendingDuringSuccessor = yield* waitForProjection( orchestrator, lateParentThreadId, (projection) => - projection.subagents.find((task) => task.id === thirdLateTask.id) - ?.completionDelivery?.state === "pending", + [thirdLateTask.id, fourthLateTask.id].every( + (taskId) => + projection.subagents.find((task) => task.id === taskId)?.completionDelivery + ?.state === "pending", + ), ); + // The running successor still owns the cohort's only reservation. + expect( + pendingDuringSuccessor.runs.find((run) => run.id === lateParentRun.id) + ?.delegatedCompletion?.delivery, + ).toMatchObject({ + generation: successorDelivery.generation, + taskIds: successorDelivery.taskIds, + }); + yield* expectOffersToStay(2); yield* Deferred.succeed(successorGate, undefined); - const exhaustedCohort = yield* waitForProjection( + const batchedReserved = yield* waitForProjection( orchestrator, lateParentThreadId, (projection) => { - const cohort = projection.runs.find( - (run) => run.id === lateParentRun.id, - )?.delegatedCompletion; + const delivery = projection.runs.find((run) => run.id === lateParentRun.id) + ?.delegatedCompletion?.delivery; return ( - cohort?.settledDeliveryCount === 2 && - cohort.delivery === null && + delivery !== undefined && + delivery !== null && + 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" + "completed" ); }, ); + const batchedDelivery = batchedReserved.runs.find((run) => run.id === lateParentRun.id) + ?.delegatedCompletion?.delivery; + if (batchedDelivery === undefined || batchedDelivery === null) { + return yield* Effect.die(new Error("Batched late completion delivery missing.")); + } + expect([...batchedDelivery.taskIds].toSorted()).toEqual( + [thirdLateTask.id, fourthLateTask.id].toSorted(), + ); expect( - exhaustedCohort.runs.find((run) => run.id === lateParentRun.id)?.delegatedCompletion, - ).toMatchObject({ settledDeliveryCount: 2, delivery: null }); - yield* expectOffersToStay(2); + [lateTask.id, thirdLateTask.id, fourthLateTask.id].map( + (taskId) => + batchedReserved.subagents.find((task) => task.id === taskId)?.completionDelivery + ?.state, + ), + ).toEqual(["delivered", "claimed", "claimed"]); + // One offer per delivery: first, successor, and this batch. + yield* waitForContinuationOffers(3); + yield* expectOffersToStay(3); + expect( + batchedReserved.runs.filter((run) => { + const message = batchedReserved.messages.find( + (candidate) => candidate.id === run.userMessageId, + ); + return ( + message?.delegatedCompletion?.parentRunId === lateParentRun.id && + run.status !== "completed" + ); + }), + ).toHaveLength(0); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "agent", + creationSource: "server", + commandId: CommandId.make("command:mcp-late-parent:dispatch-batched-delivery"), + threadId: lateParentThreadId, + messageId: batchedDelivery.messageId, + text: "Delegated tasks reached terminal states.", + attachments: [], + modelSelection: codexSelection, + dispatchMode: { type: "queue_after_active" }, + delegatedCompletion: { + parentRunId: lateParentRun.id, + generation: batchedDelivery.generation, + taskIds: batchedDelivery.taskIds, + }, + }); + const drainedCohort = yield* waitForProjection( + orchestrator, + lateParentThreadId, + (projection) => + projection.runs.find((run) => run.id === lateParentRun.id)?.delegatedCompletion + ?.delivery === null && + projection.runs.some( + (run) => + run.userMessageId === batchedDelivery.messageId && run.status === "completed", + ), + ); + expect( + [earlyLateTask.id, lateTask.id, thirdLateTask.id, fourthLateTask.id].map( + (taskId) => + drainedCohort.subagents.find((task) => task.id === taskId)?.completionDelivery + ?.state, + ), + ).toEqual(["delivered", "delivered", "delivered", "delivered"]); + 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 @@ -3048,6 +3143,313 @@ describe("orchestrator MCP toolkit", () => { removedDelivery.subagents.find((task) => task.id === removeTask.id), ).toMatchObject({ result: expect.any(String), status: "completed" }); yield* expectOffersToStay(0); + + // A wide fan-out keeps waking its parent until every result is + // delivered. Two children finish before the first delivery starts, + // five while it runs, three before the successor starts, and two + // while the successor runs. Three batched deliveries cover all + // twelve, and the cohort never has two deliveries outstanding. + const fanoutParentThreadId = ThreadId.make("thread:mcp-fanout-parent"); + const fanoutParentGate = yield* Deferred.make(); + const firstFanoutDeliveryGate = yield* Deferred.make(); + const secondFanoutDeliveryGate = yield* Deferred.make(); + parentTerminalGates.set(fanoutParentThreadId, fanoutParentGate); + deliveryTerminalGates.set(fanoutParentThreadId, firstFanoutDeliveryGate); + yield* Ref.set(continuationOffers, []); + yield* orchestrator.dispatch({ + type: "thread.create", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("command:mcp-fanout-parent:create"), + threadId: fanoutParentThreadId, + projectId, + title: "Fan-out parent", + modelSelection: codexSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: cwd, + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "user", + creationSource: "web", + commandId: CommandId.make("command:mcp-fanout-parent:start"), + threadId: fanoutParentThreadId, + messageId: MessageId.make("message:mcp-fanout-parent:start"), + text: "Fan out twelve review agents.", + attachments: [], + modelSelection: codexSelection, + dispatchMode: { type: "start_immediately" }, + }); + const fanoutStarted = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + projection.runs.some((run) => run.status === "running") && + projection.providerTurns.some((turn) => turn.status === "running"), + ); + const fanoutRun = fanoutStarted.runs.find((run) => run.status === "running"); + if (fanoutRun === undefined || fanoutRun.rootNodeId === null) { + return yield* Effect.die(new Error("Fan-out parent run missing.")); + } + const fanoutRootNodeId = fanoutRun.rootNodeId; + const fanoutTasks = yield* Effect.forEach( + Array.from({ length: 12 }, (_, index) => index), + (index) => + orchestrator + .dispatch({ + type: "delegated_task.request", + createdBy: "agent", + creationSource: "mcp", + commandId: CommandId.make(`command:mcp-fanout-parent:task-${index}`), + parentThreadId: fanoutParentThreadId, + parentRunId: fanoutRun.id, + parentNodeId: fanoutRootNodeId, + task: cancellationPrompt, + modelSelection: codexSelection, + runtimeMode: "full-access", + interactionMode: "default", + completionWake: "always", + }) + .pipe(Effect.map(taskFromDispatch)), + ); + type FanoutTask = (typeof fanoutTasks)[number]; + const finishFanoutChildren = (tasks: ReadonlyArray) => + Effect.forEach( + tasks, + (task) => + Effect.gen(function* () { + const childThreadId = task.childThreadId; + if (childThreadId === null) { + return yield* Effect.die(new Error("Fan-out child thread missing.")); + } + const child = yield* waitForProjection( + orchestrator, + childThreadId, + (projection) => + projection.runs.some((run) => run.status === "running") && + projection.providerTurns.some((turn) => turn.status === "running"), + ); + const childRun = child.runs.find((run) => run.status === "running"); + if (childRun === undefined) { + return yield* Effect.die(new Error("Fan-out child run missing.")); + } + yield* orchestrator.dispatch({ + type: "run.interrupt", + commandId: CommandId.make(`command:mcp-fanout-parent:finish-${task.id}`), + threadId: childThreadId, + runId: childRun.id, + reason: "Finish this fan-out child.", + }); + }), + { discard: true }, + ); + const fanoutDelivery = (projection: OrchestrationV2ThreadProjection) => + projection.runs.find((run) => run.id === fanoutRun.id)?.delegatedCompletion + ?.delivery ?? null; + const fanoutDeliveryRuns = (projection: OrchestrationV2ThreadProjection) => + projection.runs.filter( + (run) => + projection.messages.find((message) => message.id === run.userMessageId) + ?.delegatedCompletion?.parentRunId === fanoutRun.id, + ); + // The cohort holds one reservation, and at most one of its + // delivery runs is queued or running at a time. + const expectAtMostOneOutstandingDelivery = ( + projection: OrchestrationV2ThreadProjection, + ) => + expect( + fanoutDeliveryRuns(projection).filter( + (run) => + run.status !== "completed" && + run.status !== "failed" && + run.status !== "cancelled" && + run.status !== "interrupted", + ).length, + ).toBeLessThanOrEqual(1); + const deliveryStates = ( + projection: OrchestrationV2ThreadProjection, + tasks: ReadonlyArray, + ) => + tasks.map( + (task) => + projection.subagents.find((candidate) => candidate.id === task.id) + ?.completionDelivery?.state, + ); + const sortedIds = (ids: ReadonlyArray) => [...ids].toSorted(); + const idsOf = (tasks: ReadonlyArray) => + sortedIds(tasks.map((task) => task.id)); + const dispatchFanoutDelivery = ( + label: string, + delivery: OrchestrationV2DelegatedCompletionDelivery, + ) => + orchestrator.dispatch({ + type: "message.dispatch", + createdBy: "agent", + creationSource: "server", + commandId: CommandId.make(`command:mcp-fanout-parent:dispatch-${label}`), + threadId: fanoutParentThreadId, + messageId: delivery.messageId, + text: "Delegated tasks reached terminal states.", + attachments: [], + modelSelection: codexSelection, + dispatchMode: { type: "queue_after_active" }, + delegatedCompletion: { + parentRunId: fanoutRun.id, + generation: delivery.generation, + taskIds: delivery.taskIds, + }, + }); + const deliveryRunStatus = ( + projection: OrchestrationV2ThreadProjection, + delivery: OrchestrationV2DelegatedCompletionDelivery, + ) => projection.runs.find((run) => run.userMessageId === delivery.messageId)?.status; + + const beforeFirstStarts = fanoutTasks.slice(0, 2); + const duringFirst = fanoutTasks.slice(2, 7); + const beforeSecondStarts = fanoutTasks.slice(7, 10); + const duringSecond = fanoutTasks.slice(10, 12); + + yield* finishFanoutChildren(beforeFirstStarts); + const firstReserved = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + deliveryStates(projection, beforeFirstStarts).every( + (state) => state === "claimed", + ) && fanoutDelivery(projection)?.taskIds.length === 2, + ); + const firstFanoutDelivery = fanoutDelivery(firstReserved); + if (firstFanoutDelivery === null) { + return yield* Effect.die(new Error("First fan-out delivery missing.")); + } + expect(sortedIds(firstFanoutDelivery.taskIds)).toEqual(idsOf(beforeFirstStarts)); + yield* dispatchFanoutDelivery("first", firstFanoutDelivery); + yield* Deferred.succeed(fanoutParentGate, undefined); + yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => deliveryRunStatus(projection, firstFanoutDelivery) === "running", + ); + + yield* finishFanoutChildren(duringFirst); + const pendingDuringFirst = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + deliveryStates(projection, duringFirst).every((state) => state === "pending"), + ); + expect(fanoutDelivery(pendingDuringFirst)).toEqual(firstFanoutDelivery); + expectAtMostOneOutstandingDelivery(pendingDuringFirst); + + deliveryTerminalGates.set(fanoutParentThreadId, secondFanoutDeliveryGate); + yield* Deferred.succeed(firstFanoutDeliveryGate, undefined); + const secondReserved = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + deliveryRunStatus(projection, firstFanoutDelivery) === "completed" && + fanoutDelivery(projection)?.generation === firstFanoutDelivery.generation + 1, + ); + expect(sortedIds(fanoutDelivery(secondReserved)?.taskIds ?? [])).toEqual( + idsOf(duringFirst), + ); + expect(deliveryStates(secondReserved, beforeFirstStarts)).toEqual([ + "delivered", + "delivered", + ]); + + // Results that land before the successor starts join it. + yield* finishFanoutChildren(beforeSecondStarts); + const secondJoined = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + deliveryStates(projection, beforeSecondStarts).every( + (state) => state === "claimed", + ) && fanoutDelivery(projection)?.taskIds.length === 8, + ); + const secondFanoutDelivery = fanoutDelivery(secondJoined); + if (secondFanoutDelivery === null) { + return yield* Effect.die(new Error("Second fan-out delivery missing.")); + } + expect(sortedIds(secondFanoutDelivery.taskIds)).toEqual( + idsOf([...duringFirst, ...beforeSecondStarts]), + ); + expectAtMostOneOutstandingDelivery(secondJoined); + yield* dispatchFanoutDelivery("second", secondFanoutDelivery); + yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => deliveryRunStatus(projection, secondFanoutDelivery) === "running", + ); + + yield* finishFanoutChildren(duringSecond); + const pendingDuringSecond = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + deliveryStates(projection, duringSecond).every((state) => state === "pending"), + ); + expectAtMostOneOutstandingDelivery(pendingDuringSecond); + + // Later deliveries complete as soon as they start. + deliveryTerminalGates.delete(fanoutParentThreadId); + yield* Deferred.succeed(secondFanoutDeliveryGate, undefined); + const thirdReserved = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + deliveryRunStatus(projection, secondFanoutDelivery) === "completed" && + fanoutDelivery(projection)?.generation === secondFanoutDelivery.generation + 1, + ); + const thirdFanoutDelivery = fanoutDelivery(thirdReserved); + if (thirdFanoutDelivery === null) { + return yield* Effect.die(new Error("Third fan-out delivery missing.")); + } + expect(sortedIds(thirdFanoutDelivery.taskIds)).toEqual(idsOf(duringSecond)); + yield* dispatchFanoutDelivery("third", thirdFanoutDelivery); + const fanoutDrained = yield* waitForProjection( + orchestrator, + fanoutParentThreadId, + (projection) => + deliveryRunStatus(projection, thirdFanoutDelivery) === "completed" && + fanoutDelivery(projection) === null, + ); + expect(deliveryStates(fanoutDrained, fanoutTasks)).toEqual( + fanoutTasks.map(() => "delivered"), + ); + expect( + fanoutTasks.map( + (task) => + fanoutDrained.subagents.find((candidate) => candidate.id === task.id)?.status, + ), + ).toEqual(fanoutTasks.map(() => "interrupted")); + const drainedDeliveryRuns = fanoutDeliveryRuns(fanoutDrained); + expect(drainedDeliveryRuns.map((run) => run.status)).toEqual([ + "completed", + "completed", + "completed", + ]); + const deliveredTaskIds = drainedDeliveryRuns.flatMap( + (run) => + fanoutDrained.messages.find((message) => message.id === run.userMessageId) + ?.delegatedCompletion?.taskIds ?? [], + ); + expect(sortedIds(deliveredTaskIds)).toEqual(idsOf(fanoutTasks)); + const offeredDeliveries = new Set( + (yield* Ref.get(continuationOffers)).map( + (offer) => offer.delegatedCompletion?.messageId, + ), + ); + expect(offeredDeliveries).toEqual( + new Set([ + firstFanoutDelivery.messageId, + secondFanoutDelivery.messageId, + thirdFanoutDelivery.messageId, + ]), + ); }).pipe(Effect.provide(testLayer)); }), ), diff --git a/apps/server/src/orchestration-v2/CheckpointCaptureService.test.ts b/apps/server/src/orchestration-v2/CheckpointCaptureService.test.ts index 17ee42f99488..824f4b36659b 100644 --- a/apps/server/src/orchestration-v2/CheckpointCaptureService.test.ts +++ b/apps/server/src/orchestration-v2/CheckpointCaptureService.test.ts @@ -62,13 +62,11 @@ it.layer(ProjectionStoreTestLayer)("CheckpointCaptureServiceV2", (it) => { const staleDelegatedCompletion = { disposition: "open" as const, nextGeneration: 1, - settledDeliveryCount: 0, delivery: null, }; const newerCohort = { disposition: "open" as const, nextGeneration: 2, - settledDeliveryCount: 1, delivery: { generation: 1, messageId: deliveryMessageId, diff --git a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts index 8da8378a8b6e..7152e095542f 100644 --- a/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts +++ b/apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts @@ -232,7 +232,6 @@ const seedParentWithTerminalTask = (input: { delegatedCompletion: { disposition: "open", nextGeneration: 2, - settledDeliveryCount: 1, delivery: input.deliveryTaskIds === undefined ? null @@ -361,10 +360,6 @@ it.layer(TestLayer)("delegated completion delivery repairs", (it) => { const cohort = accepted.runs.find((row) => row.id === runId)?.delegatedCompletion; assert.deepEqual(cohort?.delivery?.taskIds, pendingIds); assert.equal(cohort?.delivery?.generation, 2); - assert.equal( - cohort?.settledDeliveryCount, - projection.runs.find((row) => row.id === runId)?.delegatedCompletion?.settledDeliveryCount, - ); for (const id of pendingIds) { assert.deepEqual(accepted.subagents.find((row) => row.id === id)?.completionDelivery, { state: "claimed", @@ -473,7 +468,6 @@ it.layer(TestLayer)("delegated completion delivery repairs", (it) => { assert.deepEqual(parentRun?.delegatedCompletion, { disposition: "open", nextGeneration: 2, - settledDeliveryCount: 1, delivery: null, }); assert.isFalse( diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 319f6f02ff26..f959f1b7e987 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -1888,7 +1888,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const nextCohort = { disposition: nextDisposition, nextGeneration: cohort?.nextGeneration ?? 1, - settledDeliveryCount: cohort?.settledDeliveryCount ?? 0, delivery: null, } as const; yield* emitEvent({ @@ -8074,8 +8073,9 @@ 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 + // successor that listener reserves, rather than creating a competing + // delivery. return { task: { ...input.updatedTask, @@ -8114,24 +8114,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }; } - 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({ @@ -8145,7 +8127,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio const nextCohort = { disposition: "open" as const, nextGeneration: generation + 1, - settledDeliveryCount, delivery: { generation, messageId, @@ -8523,12 +8504,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio nextTaskStates.get(task.id)?.state === "pending"), ) .map((task) => task.id); - const settledDeliveryCount = (cohort.settledDeliveryCount ?? 0) + 1; + // Results that arrived while this delivery was outstanding go out + // together in one successor. Each child becomes pending once, so a + // cohort's successors are bounded by its children. const canReserveFollowUp = cohort.disposition === "open" && projection.thread.archivedAt === null && projection.thread.deletedAt === null && - settledDeliveryCount < 2 && pendingTaskIds.length > 0; const nextDelivery = canReserveFollowUp ? { @@ -8557,7 +8539,6 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ...parentRun, delegatedCompletion: { ...cohort, - settledDeliveryCount, nextGeneration: nextDelivery === null ? cohort.nextGeneration : cohort.nextGeneration + 1, delivery: nextDelivery, }, @@ -8631,9 +8612,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) => @@ -8681,8 +8660,8 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio payload: { ...task, completionDelivery: { state, observedByRunId: null }, updatedAt: now }, }); } - // Provider acceptance drains this batch but does not acknowledge its results - // or spend an idle-wake allowance. task_status owns acknowledgment. + // Provider acceptance drains this batch but does not acknowledge its results. + // task_status owns acknowledgment. yield* emitEvent({ type: "run.updated", threadId: command.threadId, diff --git a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts index 50ee53d30fcc..750d2b4312bc 100644 --- a/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts +++ b/apps/server/src/orchestration-v2/SteeringCompletion.integration.test.ts @@ -217,7 +217,6 @@ for (const mailbox of [false, true]) { delegatedCompletion: { disposition: "open", nextGeneration: 2, - settledDeliveryCount: 0, delivery: { generation: 1, messageId, taskIds: [taskId] }, }, }, @@ -331,7 +330,6 @@ for (const mailbox of [false, true]) { assert.equal(delivered.subagents[0]?.completionDelivery?.state, "delivered"); assert.equal(delivered.subagents[0]?.completionDelivery?.observedByRunId, null); assert.equal(delivered.runs[0]?.delegatedCompletion?.delivery, null); - assert.equal(delivered.runs[0]?.delegatedCompletion?.settledDeliveryCount, 0); assert.equal( delivered.turnItems.filter((item) => item.type === "notification").length, 1, diff --git a/apps/server/src/orchestration-v2/ThreadDeletion.test.ts b/apps/server/src/orchestration-v2/ThreadDeletion.test.ts index ce98326c09a1..4f79fe01b764 100644 --- a/apps/server/src/orchestration-v2/ThreadDeletion.test.ts +++ b/apps/server/src/orchestration-v2/ThreadDeletion.test.ts @@ -149,7 +149,6 @@ it.effect("cancels active work without reviving a run while disposing delegated delegatedCompletion: { disposition: "open", nextGeneration: 3, - settledDeliveryCount: 1, delivery: { generation: 2, messageId: queuedRun.userMessageId, taskIds: [taskId] }, }, } @@ -216,7 +215,6 @@ it.effect("cancels active work without reviving a run while disposing delegated assert.deepEqual(deleted.runs.find((run) => run.id === parentRun.id)?.delegatedCompletion, { disposition: "disposed", nextGeneration: 3, - settledDeliveryCount: 1, delivery: null, }); assert.equal(deleted.subagents[0]?.completionDelivery?.state, "disposed"); diff --git a/apps/server/src/orchestration-v2/ThreadDeletion.ts b/apps/server/src/orchestration-v2/ThreadDeletion.ts index c65a4650020d..86d1e8050267 100644 --- a/apps/server/src/orchestration-v2/ThreadDeletion.ts +++ b/apps/server/src/orchestration-v2/ThreadDeletion.ts @@ -154,7 +154,6 @@ export const planThreadDeletion = Effect.fn("ThreadDeletion.planThreadDeletion") delegatedCompletion: { disposition: "disposed", nextGeneration: cohort?.nextGeneration ?? 1, - settledDeliveryCount: cohort?.settledDeliveryCount ?? 0, delivery: null, }, }, diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 821e5fc60687..679833e8d0cd 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -480,9 +480,6 @@ 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. - settledDeliveryCount: Schema.optional(NonNegativeInt), delivery: Schema.NullOr(OrchestrationV2DelegatedCompletionDelivery), }); export type OrchestrationV2DelegatedCompletionCohort =