From 96d01d05cd7af0a7577f84726ce879fe223efd30 Mon Sep 17 00:00:00 2001 From: Mnigos Date: Sat, 3 Oct 2026 12:51:58 +0200 Subject: [PATCH] fix(server): a thread cannot be archived while its turn is running --- ...OrchestratorMcpToolkit.integration.test.ts | 79 +++++++++ .../src/mcp/toolkits/thread/handlers.ts | 12 +- .../Orchestrator.control-reads.test.ts | 11 ++ .../src/orchestration-v2/Orchestrator.ts | 36 ++++- .../src/orchestration-v2/runtimeLayer.test.ts | 152 ++++++++++++++++++ .../src/components/threadActionMenu.logic.ts | 2 +- 6 files changed, 286 insertions(+), 6 deletions(-) diff --git a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts index 142e0cc55f12..a7b40ed37aef 100644 --- a/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts +++ b/apps/server/src/mcp/OrchestratorMcpToolkit.integration.test.ts @@ -702,6 +702,85 @@ 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.structuredContent).toMatchObject({ + _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(); + 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(otherArchiveCall.structuredContent).toEqual(selfArchiveCall.structuredContent); + 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.")); } diff --git a/apps/server/src/mcp/toolkits/thread/handlers.ts b/apps/server/src/mcp/toolkits/thread/handlers.ts index c0f5d135b6cc..faf61fddce84 100644 --- a/apps/server/src/mcp/toolkits/thread/handlers.ts +++ b/apps/server/src/mcp/toolkits/thread/handlers.ts @@ -297,7 +297,17 @@ export const ThreadToolkitHandlersLive = ThreadToolkit.toLayer({ default: command = { ...common, type: `thread.${input.action}` }; } - const result = yield* threads.dispatch(command).pipe(Effect.mapError(unavailable)); + 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.", + }) + : unavailable(), + ), + ); return { sequence: result.sequence }; }), }); diff --git a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts index ad41ae6ffde6..2732e0df5dfa 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -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"), diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 2707fe5b8344..fca830e14ccb 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -183,6 +183,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", + { 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."; + } +} + export class OrchestratorCommandPreviouslyRejectedError extends Schema.TaggedError()( "OrchestratorCommandPreviouslyRejectedError", { @@ -232,6 +242,7 @@ export const OrchestratorV2Error = Schema.Union([ OrchestratorCommandPreviouslyRejectedError, OrchestratorCommandIdConflictError, OrchestratorSubagentThreadReadOnlyError, + OrchestratorThreadTurnRunningError, ]); export type OrchestratorV2Error = typeof OrchestratorV2Error.Type; @@ -2370,6 +2381,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, @@ -3172,10 +3199,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) diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 2c5235c767d7..8cd1abe4511c 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -212,6 +212,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 TestLayer = Layer.mergeAll( OrchestrationV2LayerLive, OrchestrationV2EventSinkLayerLive, @@ -2864,6 +2899,13 @@ it.layer(TestLayer)("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"), @@ -2896,6 +2938,116 @@ it.layer(TestLayer)("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) => diff --git a/apps/web/src/components/threadActionMenu.logic.ts b/apps/web/src/components/threadActionMenu.logic.ts index a1ed9c277a4f..aa140f836d35 100644 --- a/apps/web/src/components/threadActionMenu.logic.ts +++ b/apps/web/src/components/threadActionMenu.logic.ts @@ -48,7 +48,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;