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..0f1d12bb6861 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.control-reads.test.ts @@ -23,6 +23,7 @@ import * as Layer from "effect/Layer"; import * as SqlClient from "effect/unstable/sql/SqlClient"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; import { CodexProviderCapabilitiesV2 } from "./Adapters/CodexAdapterV2.ts"; +import * as EffectOutbox from "./EffectOutbox.ts"; import * as Orchestrator from "./Orchestrator.ts"; import * as ProjectionStore from "./ProjectionStore.ts"; import type { ProviderAdapterV2Shape } from "./ProviderAdapter.ts"; @@ -42,6 +43,7 @@ const database = SqlitePersistenceMemory; const testLayer = Layer.mergeAll( database, ProjectionStore.layer.pipe(Layer.provide(database)), + EffectOutbox.layer.pipe(Layer.provide(database)), makeOrchestratorV2ReplayLayerWithRegistry( { name: "control-reads" }, ProviderAdapterRegistry.makeLayer([adapter]), @@ -49,6 +51,63 @@ const testLayer = Layer.mergeAll( ), ); +for (const terminalCommand of ["thread.archive", "thread.delete"] as const) { + it.effect(`rejects prepared-run.release after ${terminalCommand}`, () => + Effect.gen(function* () { + const orchestrator = yield* Orchestrator.OrchestratorV2; + const projections = yield* ProjectionStore.ProjectionStoreV2; + const outbox = yield* EffectOutbox.EffectOutboxV2; + const threadId = ThreadId.make("thread:terminal-preparation"); + yield* orchestrator.dispatch({ + type: "thread.create", + commandId: CommandId.make("create-terminal-preparation"), + threadId, + projectId: ProjectId.make("project:terminal-preparation"), + title: "Preparing thread", + modelSelection, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: null, + createdBy: "user", + creationSource: "web", + }); + yield* orchestrator.dispatch({ + type: "message.dispatch", + commandId: CommandId.make("prepare-terminal-message"), + threadId, + messageId: MessageId.make("preparing-input"), + text: "Prepare workspace", + attachments: [], + dispatchMode: { type: "defer_start" }, + createdBy: "user", + creationSource: "web", + }); + yield* orchestrator.dispatch({ + type: terminalCommand, + commandId: CommandId.make("close-terminal-preparation"), + threadId, + }); + const before = yield* projections.getThreadProjection(threadId); + const runId = before.runs[0]!.id; + const commandId = CommandId.make("release-terminal-preparation"); + const error = yield* orchestrator + .dispatch({ + type: "prepared-run.release", + commandId, + threadId, + runId, + }) + .pipe(Effect.flip); + + assert.equal(error._tag, "OrchestratorDispatchError"); + assert.equal(error.cause, `Thread ${threadId} is not active.`); + assert.deepEqual(yield* projections.getThreadProjection(threadId), before); + assert.deepEqual(yield* outbox.listByCommandId(commandId), []); + }).pipe(Effect.provide(testLayer)), + ); +} + it.effect( "dispatches metadata, queue resume and request controls without hydrating unrelated history", () => diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index df9805d504c9..2dd672b39811 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -7360,6 +7360,13 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ["runs", "attempts", "nodes", "providerThreads", "turnItems"], { turnItemTypes: ["command_execution"], turnItemRunId: command.runId }, ); + if (projection.thread.archivedAt !== null || projection.thread.deletedAt !== null) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: `Thread ${command.threadId} is not active.`, + }); + } const state = preparedRunState(command, projection); if (state === null) { return yield* new OrchestratorDispatchError({ diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 7b33d6f6ec91..cb8ac9bf9fd5 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -1959,6 +1959,77 @@ it.effect("shared intake preserves durable attachment bytes after a lost launch }).pipe(Effect.provide(Layer.mergeAll(harness.layer, files))); }); +it.effect("fails preparation completed while archived and accepts a new send after unarchive", () => + Effect.gen(function* () { + const setupEntered = yield* Deferred.make(); + const allowSetup = yield* Deferred.make(); + const harness = makeHarness({ + runSetup: () => + Deferred.succeed(setupEntered, undefined).pipe( + Effect.andThen(Deferred.await(allowSetup)), + Effect.as({ status: "no-script" as const }), + ), + }); + yield* Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const outbox = yield* EffectOutbox.EffectOutboxV2; + const input = launchInput({ + command: "launch:archive-during-setup", + thread: "thread:archive-during-setup", + message: "Start", + }); + const launched = yield* launches.launch(input); + yield* Deferred.await(setupEntered); + yield* threads.dispatch({ + type: "thread.archive", + commandId: CommandId.make("archive:during-setup"), + threadId: launched.threadId, + }); + yield* Deferred.succeed(allowSetup, undefined); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + (stored.commandId === CommandId.make(`${input.commandId}:fail`) || + stored.commandId === CommandId.make(`${input.commandId}:release`)) && + stored.event.type === "run.updated", + ), + Stream.runHead, + ); + const failed = yield* threads.getThreadProjection(launched.threadId); + assert.equal(failed.runs[0]?.status, "failed"); + assert.isEmpty(failed.providerTurns); + assert.isEmpty(failed.checkpointScopes); + assert.isEmpty(yield* outbox.listByCommandId(CommandId.make(`${input.commandId}:release`))); + + yield* threads.dispatch({ + type: "thread.unarchive", + commandId: CommandId.make("unarchive:after-setup"), + threadId: launched.threadId, + }); + const next = yield* threads.sendToThread({ + projectId, + commandId: CommandId.make("send:after-unarchive"), + threadId: launched.threadId, + messageId: MessageId.make("message:after-unarchive"), + text: "Try again", + attachments: [], + mode: "auto", + createdBy: "user", + creationSource: "web", + }); + assert.equal(next.delivery, "started"); + assert.equal(next.run.status, "starting"); + assert.notEqual(next.run.id, failed.runs[0]?.id); + assert.isTrue( + (yield* outbox.listByCommandId(CommandId.make("send:after-unarchive"))).some( + (entry) => entry.request.type === "provider-turn.start", + ), + ); + }).pipe(Effect.provide(harness.layer)); + }), +); + it.effect("cancels tracked setup before provider work is released", () => Effect.gen(function* () { const entered = yield* Deferred.make();