Skip to content
Merged
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
236 changes: 235 additions & 1 deletion apps/server/src/orchestration-v2/ThreadLaunchService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ import {
type ServerProvider,
ThreadId,
} from "@t3tools/contracts";
import * as Cause from "effect/Cause";
import * as Deferred from "effect/Deferred";
import * as DateTime from "effect/DateTime";
import * as Duration from "effect/Duration";
Expand Down Expand Up @@ -59,6 +60,7 @@ import * as ThreadTitleRegeneration from "./ThreadTitleRegenerationService.ts";
import { makeOrchestratorV2ReplayLayerWithRegistry } from "./testkit/ProviderReplayHarness.ts";

const projectId = ProjectId.make("project:launch-test");
const otherProjectId = ProjectId.make("project:launch-other");
const encodeThreadProjection = Schema.encodeEffect(OrchestrationV2ThreadProjectionJson);
const modelSelection = {
instanceId: ProviderInstanceId.make("codex"),
Expand All @@ -78,6 +80,12 @@ const project = {
deletedAt: null,
} as const;

const otherProject = {
...project,
id: otherProjectId,
title: "Other",
} as const;

const adapter = {
instanceId: modelSelection.instanceId,
driver: ProviderDriverKind.make("codex"),
Expand Down Expand Up @@ -136,7 +144,14 @@ function makeHarness(options: HarnessOptions = {}) {
bootstrap: () => Effect.die("unused"),
update: () => Effect.die("unused"),
delete: () => Effect.die("unused"),
getById: (id) => Effect.succeed(id === projectId ? Option.some(project) : Option.none()),
getById: (id) =>
Effect.succeed(
id === projectId
? Option.some(project)
: id === otherProjectId
? Option.some(otherProject)
: Option.none(),
),
getByWorkspaceRoot: () => Effect.succeed(Option.some(project)),
snapshot: Effect.die("unused"),
}),
Expand Down Expand Up @@ -1253,6 +1268,225 @@ for (const failurePoint of ["worktree", "setup"] as const) {
);
}

it.effect("replays a server-allocated launch", () =>
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 { threadId: _unusedThreadId, ...rest } = launchInput({
command: "command:launch:allocated-retry",
thread: "unused",
message: "Only once",
});
const first = yield* launches.launch(rest);
yield* Deferred.await(setupEntered);
const retry = yield* launches.launch(rest);
assert.equal(first.threadId, retry.threadId);
assert.isFalse(first.resumed);
assert.isTrue(retry.resumed);
assert.equal(harness.runSetup.mock.calls.length, 1);
assert.equal(retry.projection.messages.length, 1);
assert.equal(retry.projection.runs.length, 1);
assert.equal(retry.projection.messages[0]?.id, first.projection.messages[0]?.id);
assert.equal(retry.projection.runs[0]?.id, first.projection.runs[0]?.id);
yield* Deferred.succeed(allowSetup, undefined);
yield* threads.streamStoredEventsFrom({ threadId: first.threadId }).pipe(
Stream.filter(
(stored) =>
stored.commandId === CommandId.make(`${rest.commandId}:release`) &&
stored.event.type === "run.updated",
),
Stream.runHead,
);
const settled = yield* launches.launch(rest);
assert.equal(settled.threadId, first.threadId);
assert.isTrue(settled.resumed);
assert.equal(settled.projection.messages[0]?.id, first.projection.messages[0]?.id);
assert.equal(settled.projection.runs[0]?.id, first.projection.runs[0]?.id);
assert.equal(harness.runSetup.mock.calls.length, 1);
}).pipe(Effect.provide(harness.layer));
}),
);

it.effect("rejects a server-allocated launch replay with a mismatching thread id", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const { threadId: _unusedThreadId, ...rest } = launchInput({
command: "command:launch:allocated-mismatch",
thread: "unused",
message: "Mismatch",
});
const first = yield* launches.launch(rest);
const failed = yield* launches
.launch({
...rest,
threadId: ThreadId.make("thread:launch:allocated-mismatch"),
})
.pipe(Effect.flip);
assert.notEqual(first.threadId, ThreadId.make("thread:launch:allocated-mismatch"));
assert.equal(failed._tag, "ThreadLaunchError");
assert.equal(failed.operation, "create-thread");
assert.include(String(failed.cause), "cannot be replayed");
}).pipe(Effect.provide(harness.layer));
});

it.effect("rejects a server-allocated launch receipt from another project", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const { threadId: _unusedThreadId, ...rest } = launchInput({
command: "command:launch:allocated-wrong-project",
thread: "unused",
message: "Wrong project",
});
const first = yield* launches.launch(rest);
const failed = yield* launches
.launch({
...rest,
projectId: otherProjectId,
})
.pipe(Effect.flip);
assert.equal(failed._tag, "ThreadLaunchError");
assert.equal(failed.operation, "resolve-project");
assert.equal(failed.threadId, first.threadId);
assert.equal(failed.cause, "Project identity changed.");
}).pipe(Effect.provide(harness.layer));
});

it.effect("rejects a server-allocated launch retry after the thread is deleted", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const { threadId: _unusedThreadId, ...rest } = launchInput({
command: "command:launch:allocated-deleted",
thread: "unused",
message: "Deleted before retry",
});
const first = yield* launches.launch(rest);
yield* threads.dispatch({
type: "thread.delete",
commandId: CommandId.make("command:launch:allocated-deleted:delete"),
threadId: first.threadId,
});
const failed = yield* launches.launch(rest).pipe(Effect.flip);
assert.equal(failed._tag, "ThreadLaunchError");
assert.equal(failed.operation, "create-thread");
assert.equal(failed.threadId, first.threadId);
assert.equal(failed.cause, "Thread not found.");
const shells = yield* threads.listProjectThreads({ projectId, includeSubagents: true });
assert.equal(shells.length, 0);
}).pipe(Effect.provide(harness.layer));
});

it.effect("does not treat an unrelated accepted command receipt as a launch", () => {
const harness = makeHarness();
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const threadId = ThreadId.make("thread:launch:unrelated-receipt");
yield* threads.dispatch({
type: "thread.create",
commandId: CommandId.make("command:launch:unrelated-receipt:create"),
threadId,
projectId,
title: "Existing",
modelSelection,
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
createdBy: "user",
creationSource: "web",
});
yield* threads.dispatch({
type: "thread.metadata.update",
commandId: CommandId.make("command:launch:unrelated-receipt"),
threadId,
expectedEmpty: true,
});
const { threadId: _unusedThreadId, ...rest } = launchInput({
command: "command:launch:unrelated-receipt",
thread: "unused",
message: "Should not become a launch",
});
const failed = yield* launches.launch(rest).pipe(Effect.flip);
assert.equal(failed._tag, "ThreadLaunchError");
assert.equal(failed.operation, "create-thread");
assert.include(String(failed.cause), "cannot be replayed");
const projection = yield* threads.getThreadProjection(threadId);
assert.equal(projection.messages.length, 0);
assert.equal(projection.runs.length, 0);
}).pipe(Effect.provide(harness.layer));
});

it.effect("bounds concurrent first launches to one thread per command", () =>
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 { threadId: _unusedThreadId, ...rest } = launchInput({
command: "command:launch:concurrent-allocated",
thread: "unused",
message: "Race me",
});
// The command receipt is reserved atomically with the winning create, so
// a loser either replays the winner's stored events or surfaces a
// transient replay conflict that the next attempt resolves — the race
// can never persist a second thread.
const results = yield* Effect.all(
[launches.launch(rest).pipe(Effect.exit), launches.launch(rest).pipe(Effect.exit)],
{ concurrency: "unbounded" },
);
const winner = results.find(Exit.isSuccess);
assert.isDefined(winner);
const threadId = winner!.value.threadId;
for (const result of results) {
if (Exit.isSuccess(result)) {
assert.equal(result.value.threadId, threadId);
continue;
}
const error = Cause.findErrorOption(result.cause).pipe(Option.getOrThrow);
assert.equal(error._tag, "ThreadLaunchError");
assert.equal(error.operation, "create-thread");
assert.include(String(error.cause), "cannot be replayed");
const retried = yield* launches.launch(rest);
assert.equal(retried.threadId, threadId);
}
const projectThreads = yield* threads.listProjectThreads({
projectId,
includeSubagents: false,
});
assert.equal(projectThreads.length, 1);
const projection = yield* threads.getThreadProjection(threadId);
assert.equal(projection.messages.length, 1);
assert.equal(projection.runs.length, 1);
yield* Deferred.await(setupEntered);
assert.equal(harness.runSetup.mock.calls.length, 1);
yield* Deferred.succeed(allowSetup, undefined);
}).pipe(Effect.provide(harness.layer));
}),
);

it.effect("deduplicates retried launch side effects in-process", () =>
Effect.gen(function* () {
const setupEntered = yield* Deferred.make<void>();
Expand Down
28 changes: 28 additions & 0 deletions apps/server/src/orchestration-v2/ThreadLaunchService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -608,12 +608,40 @@ const make = Effect.gen(function* () {

const launchReceipt = yield* readReceipt(input, input.commandId);
return yield* Effect.gen(function* () {
// A retried launch has no client-supplied id to replay against, so
// recover the thread id its accepted create was recorded under before
// allocating another one; a fresh id would only collide with the
// recorded receipt.
const reusableLaunchReceipt =
input.threadId === undefined &&
Option.isSome(launchReceipt) &&
launchReceipt.value.status === "accepted" &&
launchReceipt.value.commandType === "thread.create"
? launchReceipt.value
: undefined;
const candidateThreadId =
input.threadId ??
reusableLaunchReceipt?.threadId ??
(yield* ids.allocate
.thread({ projectId: input.projectId })
.pipe(Effect.mapError(mapError(input, "create-thread"))));

if (reusableLaunchReceipt !== undefined) {
const shell = yield* threads
.getThreadShell(candidateThreadId)
.pipe(Effect.mapError(mapError(input, "create-thread", candidateThreadId)));
if (shell === null) {
return yield* mapError(input, "create-thread", candidateThreadId)("Thread not found.");
}
if (shell.projectId !== input.projectId) {
return yield* mapError(
input,
"resolve-project",
candidateThreadId,
)("Project identity changed.");
}
}

if (input.reuseExistingThread === true && Option.isNone(launchReceipt)) {
yield* validateReusableThread(input, candidateThreadId);
}
Expand Down
Loading