Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
478 changes: 440 additions & 38 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts

Large diffs are not rendered by default.

Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -232,7 +232,6 @@ const seedParentWithTerminalTask = (input: {
delegatedCompletion: {
disposition: "open",
nextGeneration: 2,
settledDeliveryCount: 1,
delivery:
input.deliveryTaskIds === undefined
? null
Expand Down Expand Up @@ -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",
Expand Down Expand Up @@ -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(
Expand Down
39 changes: 9 additions & 30 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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({
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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({
Expand All @@ -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,
Expand Down Expand Up @@ -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
? {
Expand Down Expand Up @@ -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,
},
Expand Down Expand Up @@ -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) =>
Expand Down Expand Up @@ -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,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -217,7 +217,6 @@ for (const mailbox of [false, true]) {
delegatedCompletion: {
disposition: "open",
nextGeneration: 2,
settledDeliveryCount: 0,
delivery: { generation: 1, messageId, taskIds: [taskId] },
},
},
Expand Down Expand Up @@ -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,
Expand Down
2 changes: 0 additions & 2 deletions apps/server/src/orchestration-v2/ThreadDeletion.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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] },
},
}
Expand Down Expand Up @@ -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");
Expand Down
1 change: 0 additions & 1 deletion apps/server/src/orchestration-v2/ThreadDeletion.ts
Original file line number Diff line number Diff line change
Expand Up @@ -154,7 +154,6 @@ export const planThreadDeletion = Effect.fn("ThreadDeletion.planThreadDeletion")
delegatedCompletion: {
disposition: "disposed",
nextGeneration: cohort?.nextGeneration ?? 1,
settledDeliveryCount: cohort?.settledDeliveryCount ?? 0,
delivery: null,
},
},
Expand Down
3 changes: 0 additions & 3 deletions packages/contracts/src/orchestrationV2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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 =
Expand Down
Loading