Skip to content
Closed
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
63 changes: 53 additions & 10 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>();
deliveryTerminalGates.set(lateParentThreadId, successorGate);
yield* orchestrator.dispatch({
Expand Down Expand Up @@ -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,
Expand All @@ -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) => {
Expand All @@ -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
Expand Down
251 changes: 248 additions & 3 deletions apps/server/src/orchestration-v2/DelegatedCompletionDelivery.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -136,6 +136,8 @@ const seedParentWithTerminalTask = (input: {
readonly deliveryState: "delivered" | "claimed" | "acknowledged" | "disposed";
readonly completionWake?: "always" | "settled_only";
readonly deliveryTaskIds?: ReadonlyArray<NodeId>;
readonly settledDeliveryCount?: number;
readonly deliveryGeneration?: number;
readonly now: DateTime.Utc;
}) =>
Effect.gen(function* () {
Expand Down Expand Up @@ -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,
},
Expand Down Expand Up @@ -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<NodeId>;
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* () {
Expand Down Expand Up @@ -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;
Expand Down
Loading
Loading