Skip to content
Open
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
80 changes: 80 additions & 0 deletions apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -788,6 +788,86 @@ describe("orchestrator MCP toolkit", () => {
yield* invoke("t3_thread_organize", { action: "unpin" });
expect((yield* orchestrator.getThreadShell(parentThreadId))?.pinnedAt).toBeNull();

// Archiving detaches the provider, so the caller's own running turn
// must refuse it instead of failing mid-turn.
const selfArchiveCall = yield* invoke("t3_thread_organize", { action: "archive" });
expect(selfArchiveCall.isError).toBe(true);
expect(declaredFailure(selfArchiveCall)).toEqual({
_tag: "OrchestratorMcpFailure",
code: "invalid_request",
message:
"A thread cannot be archived while a turn is running. Archive it after the turn ends, from another thread or in the app.",
});
const afterSelfArchive = yield* orchestrator.getThreadProjection(parentThreadId);
expect(afterSelfArchive.thread.archivedAt).toBeNull();
expect(afterSelfArchive.runs[0]?.status).toBe("running");
// A detach would have dropped the session from the projection.
expect(parent.providerSessions).not.toHaveLength(0);
expect(afterSelfArchive.providerSessions.map((session) => session.id)).toEqual(
expect.arrayContaining(parent.providerSessions.map((session) => session.id)),
);

const idleThreadId = ThreadId.make("thread:mcp-idle-archive");
yield* orchestrator.dispatch({
type: "thread.create",
createdBy: "user",
creationSource: "web",
commandId: CommandId.make("command:mcp-idle-archive:create"),
threadId: idleThreadId,
projectId,
title: "Idle thread to archive",
modelSelection: codexSelection,
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: cwd,
});
const archiveTargetGate = yield* Deferred.make<void>();
parentTerminalGates.set(idleThreadId, archiveTargetGate);
yield* orchestrator.dispatch({
type: "message.dispatch",
createdBy: "user",
creationSource: "web",
commandId: CommandId.make("command:mcp-idle-archive:start"),
threadId: idleThreadId,
messageId: MessageId.make("message:mcp-idle-archive:start"),
text: "Finish after the archive rejection.",
attachments: [],
modelSelection: codexSelection,
dispatchMode: { type: "start_immediately" },
});
const archiveEvents = yield* EventSink.EventSinkV2;
const awaitArchiveTargetStatus = (status: "running" | "completed") =>
archiveEvents.stream({ threadId: idleThreadId }).pipe(
Stream.filter(
({ event }) => event.type === "run.updated" && event.payload.status === status,
),
Stream.take(1),
Stream.runDrain,
);
yield* awaitArchiveTargetStatus("running");
const otherArchiveCall = yield* invoke("t3_thread_organize", {
threadId: idleThreadId,
action: "archive",
});
expect(declaredFailure(otherArchiveCall)).toEqual(declaredFailure(selfArchiveCall));
expect(
(yield* orchestrator.getThreadProjection(idleThreadId)).thread.archivedAt,
).toBeNull();
yield* Deferred.succeed(archiveTargetGate, undefined);
yield* awaitArchiveTargetStatus("completed");
const completedArchiveTarget = yield* orchestrator.getThreadProjection(idleThreadId);
expect(completedArchiveTarget.runs[0]?.checkpointId).not.toBeNull();
expect(completedArchiveTarget.runs[0]?.status).toBe("completed");
const idleArchiveCall = yield* invoke("t3_thread_organize", {
threadId: idleThreadId,
action: "archive",
});
expect(idleArchiveCall.structuredContent).toHaveProperty("sequence");
expect(
(yield* orchestrator.getThreadProjection(idleThreadId)).thread.archivedAt,
).not.toBeNull();

if (parentRun === undefined || parentRun.rootNodeId === null) {
return yield* Effect.die(new Error("Parent run missing."));
}
Expand Down
12 changes: 11 additions & 1 deletion apps/server/src/mcp/toolkits/thread/handlers.ts
Original file line number Diff line number Diff line change
Expand Up @@ -315,7 +315,17 @@ export const layer = McpToolAccess.toLayer(ThreadToolkit, {
default:
command = { ...common, type: `thread.${input.action}` };
}
const result = yield* threads.dispatch(command).pipe(Effect.mapError(dispatchFailure));
const result = yield* threads.dispatch(command).pipe(
Effect.mapError((error) =>
error._tag === "OrchestratorThreadTurnRunningError"
? new OrchestratorMcpFailure({
code: "invalid_request",
message:
"A thread cannot be archived while a turn is running. Archive it after the turn ends, from another thread or in the app.",
})
: dispatchFailure(error),
),
);
return { sequence: result.sequence };
}),
),
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -266,6 +266,17 @@ it.effect(
turnItemTypes: ["user_message"],
});
assert.isAbove(fresh.turnItems.at(-1)!.ordinal, 900);
// Archive refuses a preparing run, so end the deferred one first.
const deferred = (yield* projections.getThreadRecords(threadId, ["runs"])).runs.at(-1)!;
assert.equal(deferred.status, "preparing");
yield* projections.apply({
id: EventId.make("cancel-deferred-run"),
type: "run.updated",
threadId,
runId: deferred.id,
occurredAt: now,
payload: { ...deferred, status: "cancelled", completedAt: now },
});
yield* orchestrator.dispatch({
type: "thread.archive",
commandId: CommandId.make("archive-with-old-history"),
Expand Down
36 changes: 32 additions & 4 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -190,6 +190,16 @@ export class OrchestratorSubagentThreadReadOnlyError extends Schema.TaggedError<
}
}

/** Archive detaches the thread's provider sessions, which would fail a turn in flight. */
export class OrchestratorThreadTurnRunningError extends Schema.TaggedError<OrchestratorThreadTurnRunningError>()(
"OrchestratorThreadTurnRunningError",
{ commandId: CommandId, threadId: ThreadId },
) {
override get message(): string {
return "This thread cannot be archived while a turn is running. Archive it after the turn ends.";
}
}

/** The command's thread runs above the modes its sender may touch (see `DispatchModeLimit`). */
export class OrchestratorThreadAboveModeLimitError extends Schema.TaggedError<OrchestratorThreadAboveModeLimitError>()(
"OrchestratorThreadAboveModeLimitError",
Expand Down Expand Up @@ -255,6 +265,7 @@ export const OrchestratorV2Error = Schema.Union([
OrchestratorCommandPreviouslyRejectedError,
OrchestratorCommandIdConflictError,
OrchestratorSubagentThreadReadOnlyError,
OrchestratorThreadTurnRunningError,
OrchestratorThreadAboveModeLimitError,
]);
export type OrchestratorV2Error = typeof OrchestratorV2Error.Type;
Expand Down Expand Up @@ -2478,6 +2489,22 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
cause: `Thread ${command.threadId} is already archived.`,
});
}
if (command.type === "thread.archive") {
// Archive detaches every provider session below, which would fail a turn
// in flight. Queued runs are cancelled by the archive itself.
const { runs } = yield* loadProjectionForCommand(command, ["runs"]);
if (
runs.some(
(run) =>
run.status === "preparing" || run.status === "starting" || run.status === "running",
)
) {
return yield* new OrchestratorThreadTurnRunningError({
commandId: command.commandId,
threadId: command.threadId,
});
}
}
if (command.type === "thread.unarchive" && thread.archivedAt === null) {
return yield* new OrchestratorDispatchError({
commandId: command.commandId,
Expand Down Expand Up @@ -3289,10 +3316,11 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio
// Settle joins archive here: both mean "done with this
// thread", so a live provider session must not keep running background
// work (PR monitors, dev servers, subagent fleets) after any of them
// lands. The settle guard above already rejects active or blocked runs,
// so for settle this only ever stops an idle session; commands are
// decided serially against the projection, so a turn start that
// re-engages the thread cannot race this detach.
// lands. The guards above reject a preparing, starting or running run
// (and settle also rejects any other active or blocked run), so this
// never detaches a provider mid-turn; commands are decided serially
// against the projection, so a turn start that re-engages the thread
// cannot race this detach.
const detachSessionIds = new Set(
command.type === "thread.archive" || command.type === "thread.settle"
? (providerContext?.providerSessions ?? []).map((session) => session.id)
Expand Down
152 changes: 152 additions & 0 deletions apps/server/src/orchestration-v2/runtimeLayer.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -263,6 +263,41 @@ const moveProject = (projectId: ProjectId, workspaceRoot: string, updatedAt: str
}),
);

/** Stage run rows for command-policy tests. This does not simulate provider or checkpoint effects. */
const rewriteRuns = (
threadId: ThreadId,
commandId: string,
rewrite: (run: OrchestrationV2Run, now: DateTime.Utc) => OrchestrationV2Run | undefined,
) =>
Effect.gen(function* () {
const orchestrator = yield* Orchestrator.OrchestratorV2;
const eventSink = yield* EventSink.EventSinkV2;
const projection = yield* orchestrator.getThreadProjection(threadId);
const now = yield* DateTime.now;
yield* eventSink.commitCommand({
commandId: CommandId.make(commandId),
threadId,
commandType: "provider-runtime.reconcile",
acceptedAt: now,
events: projection.runs.flatMap((run) => {
const payload = rewrite(run, now);
return payload === undefined
? []
: [
{
id: EventId.make(`${commandId}:${run.id}`),
type: "run.updated" as const,
threadId,
runId: run.id,
occurredAt: now,
payload,
},
];
}),
effects: [],
});
});

const layerTest = Layer.mergeAll(
RuntimeLayer.layer,
RuntimeLayer.layerEventSink,
Expand Down Expand Up @@ -3548,6 +3583,13 @@ it.layer(layerTest)("OrchestrationV2LayerLive lifecycle", (it) => {
assert.isDefined(activeRun);
assert.isDefined(queuedRun);

// Archive refuses a starting turn, so simulate a restart: the active run
// ends and recovery holds the queue.
yield* rewriteRuns(threadId, "runtime-layer-archive-queued-restart", (run, now) =>
run.status === "queued"
? { ...run, queueHeld: true }
: { ...run, status: "cancelled", completedAt: now },
);
yield* orchestrator.dispatch({
type: "thread.archive",
commandId: CommandId.make("runtime-layer-archive-queued-archive"),
Expand Down Expand Up @@ -3580,6 +3622,116 @@ it.layer(layerTest)("OrchestrationV2LayerLive lifecycle", (it) => {
}),
);

const startArchiveTestThread = (name: string) =>
Effect.gen(function* () {
const orchestrator = yield* Orchestrator.OrchestratorV2;
const threadId = ThreadId.make(`runtime-layer-archive-${name}-thread`);
yield* orchestrator.dispatch({
type: "thread.create",
createdBy: "user",
creationSource: "web",
commandId: CommandId.make(`runtime-layer-archive-${name}-create`),
threadId,
projectId: ProjectId.make(`runtime-layer-archive-${name}-project`),
title: `Archive ${name}`,
modelSelection,
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: `/tmp/runtime-layer-archive-${name}`,
});
yield* orchestrator.dispatch({
type: "message.dispatch",
createdBy: "user",
creationSource: "web",
commandId: CommandId.make(`runtime-layer-archive-${name}-message`),
threadId,
messageId: MessageId.make(`runtime-layer-archive-${name}-message`),
text: "Keep the provider occupied.",
attachments: [],
modelSelection,
dispatchMode: { type: "start_immediately" },
});
const started = yield* orchestrator.getThreadProjection(threadId);
assert.deepEqual(
started.runs.map((run) => run.status),
["starting"],
);
return threadId;
});

it.effect.each(["preparing", "starting", "running"] as const)(
"refuses to archive a thread while its turn is %s, then archives it once the turn ends",
(status) =>
Effect.gen(function* () {
const orchestrator = yield* Orchestrator.OrchestratorV2;
const eventSink = yield* EventSink.EventSinkV2;
const outbox = yield* EffectOutbox.EffectOutboxV2;
const threadId = yield* startArchiveTestThread(status);
if (status !== "starting") {
yield* rewriteRuns(threadId, `runtime-layer-archive-${status}-stage`, (run) => ({
...run,
status,
}));
}
const previousSequence = yield* orchestrator.getThreadEventSequence(threadId);

const commandId = CommandId.make(`runtime-layer-archive-${status}-archive`);
const error = yield* orchestrator
.dispatch({ type: "thread.archive", commandId, threadId })
.pipe(Effect.flip);
assert.instanceOf(error, Orchestrator.OrchestratorThreadTurnRunningError);
assert.equal(
error.message,
"This thread cannot be archived while a turn is running. Archive it after the turn ends.",
);
assert.equal(yield* orchestrator.getThreadEventSequence(threadId), previousSequence);
assert.deepEqual(
yield* eventSink.readByCommandId({ commandId }).pipe(Stream.runCollect),
[],
);
assert.deepEqual(yield* outbox.listByCommandId(commandId), []);
const refused = yield* orchestrator.getThreadProjection(threadId);
assert.isNull(refused.thread.archivedAt);
assert.deepEqual(
refused.runs.map((run) => run.status),
[status],
);

yield* rewriteRuns(threadId, `runtime-layer-archive-${status}-complete`, (run, now) => ({
...run,
status: "completed",
completedAt: now,
}));
yield* orchestrator.dispatch({
type: "thread.archive",
commandId: CommandId.make(`runtime-layer-archive-${status}-archive-after-turn`),
threadId,
});
assert.isNotNull((yield* orchestrator.getThreadProjection(threadId)).thread.archivedAt);
}),
);

it.effect("archives a thread whose run is waiting", () =>
Effect.gen(function* () {
const orchestrator = yield* Orchestrator.OrchestratorV2;
const threadId = yield* startArchiveTestThread("waiting");
yield* rewriteRuns(threadId, "runtime-layer-archive-waiting-stage", (run) => ({
...run,
status: "waiting",
}));

yield* orchestrator.dispatch({
type: "thread.archive",
commandId: CommandId.make("runtime-layer-archive-waiting-archive"),
threadId,
});
const archived = yield* orchestrator.getThreadProjection(threadId);
assert.isNotNull(archived.thread.archivedAt);
assert.equal(archived.runs[0]?.status, "waiting");
}),
);

it.effect.each([false, true])(
"promotes only one queued run after each terminal run (notification: %s)",
(automatic) =>
Expand Down
2 changes: 1 addition & 1 deletion apps/web/src/components/threadActionMenu.logic.ts
Original file line number Diff line number Diff line change
Expand Up @@ -87,7 +87,7 @@ export interface ThreadActionMenuState {
readonly isSnoozed: boolean;
readonly canSnoozeNow: boolean;
readonly isRegeneratingTitle: boolean;
/** Archive rejects a thread with an attached provider, so disable it here rather than let the action fail. */
/** The server rejects archive while a run is preparing, starting or running (idle attached providers still archive), so disable it here. */
readonly isRunning: boolean;
readonly supports: {
readonly settlement: boolean;
Expand Down
Loading