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
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand All @@ -42,13 +43,71 @@ 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]),
{ databaseLayer: database, runEffectWorker: false },
),
);

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",
() =>
Expand Down
7 changes: 7 additions & 0 deletions apps/server/src/orchestration-v2/Orchestrator.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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.`,
});
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.
const state = preparedRunState(command, projection);
if (state === null) {
return yield* new OrchestratorDispatchError({
Expand Down
71 changes: 71 additions & 0 deletions apps/server/src/orchestration-v2/ThreadLaunchService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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<void>();
const allowSetup = yield* Deferred.make<void>();
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<void>();
Expand Down
Loading