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
368 changes: 335 additions & 33 deletions apps/server/src/orchestration-v2/ThreadLaunchService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -34,12 +34,14 @@ import * as Effect from "effect/Effect";
import * as Exit from "effect/Exit";
import * as Layer from "effect/Layer";
import * as Option from "effect/Option";
import * as Path from "effect/Path";
import * as Ref from "effect/Ref";
import * as Stream from "effect/Stream";
import * as Schema from "effect/Schema";
import * as TestClock from "effect/testing/TestClock";

import * as GitWorkflow from "../git/GitWorkflowService.ts";
import * as GitVcsDriver from "../vcs/GitVcsDriver.ts";
import * as SqlitePersistence from "../persistence/Sqlite.ts";
import * as ProjectStore from "./ProjectStore.ts";
import * as ProjectService from "../project/ProjectService.ts";
Expand Down Expand Up @@ -1235,6 +1237,289 @@ it.effect("keeps an explicit branch name instead of generating one", () =>
}),
);

it.effect.each([false, true])(
"bulk launches get independent real worktrees with existing branch %s",
(branchExists) =>
Effect.gen(function* () {
const driver = yield* GitVcsDriver.GitVcsDriver;
const fs = yield* FileSystem.FileSystem;
const path = yield* Path.Path;
const cwd = yield* fs.makeTempDirectoryScoped({ prefix: "launch-worktree-collisions-" });
const worktreesDirectory = yield* fs.makeTempDirectoryScoped({ prefix: "launch-worktrees-" });
const git = (args: ReadonlyArray<string>) =>
driver.execute({ operation: "ThreadLaunchService.test", cwd, args });
yield* git(["init", "--initial-branch=main"]);
yield* git([
"-c",
"user.name=Test",
"-c",
"user.email=test@example.com",
"commit",
"--allow-empty",
"-m",
"initial",
]);
const existingPath = path.join(worktreesDirectory, "existing");
if (branchExists) {
yield* driver.createWorktree({
cwd,
refName: "main",
newRefName: "feature/shared",
path: existingPath,
});
yield* fs.writeFileString(path.join(existingPath, "keep.txt"), "existing uncommitted work");
}
const harness = makeHarness({
createWorktree: (input, options) =>
driver.createWorktree({ ...input, cwd }, { ...options, worktreesDirectory }),
});
yield* Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const launched = yield* Effect.all(
Array.from({ length: 12 }, (_, index) => {
const command = `command:launch:bulk-${index}`;
return launches
.launch(
launchInput({
command,
thread: `thread:launch:bulk-${index}`,
message: "Independent work",
workspace: { type: "worktree", baseRef: "main", branch: "feature/shared" },
}),
)
.pipe(
Effect.flatMap((result) =>
threads.streamStoredEventsFrom({ threadId: result.threadId }).pipe(
Stream.filter(
(stored) =>
(stored.commandId === CommandId.make(`${command}:release`) ||
stored.commandId === CommandId.make(`${command}:fail`)) &&
stored.event.type === "run.updated",
),
Stream.runHead,
Effect.andThen(threads.getThreadProjection(result.threadId)),
),
),
);
}),
{ concurrency: "unbounded" },
);
assert.equal(new Set(launched.map((projection) => projection.thread.branch)).size, 12);
assert.equal(
new Set(launched.map((projection) => projection.thread.worktreePath)).size,
12,
);
for (const projection of launched) {
assert.equal(
projection.runs[0]?.status,
"starting",
projection.turnItems.find((item) => item.type === "error")?.failure.message,
);
assert.isTrue(yield* fs.exists(projection.thread.worktreePath!));
const base = yield* git(["config", `branch.${projection.thread.branch}.gh-merge-base`]);
assert.equal(base.stdout.trim(), "main");
}
assert.deepEqual(
launched.map((projection) => projection.thread.branch).toSorted(),
Array.from({ length: 12 }, (_, index) =>
branchExists || index > 0
? `feature/shared-${branchExists ? index + 1 : index}`
: "feature/shared",
).toSorted(),
);
if (branchExists) {
assert.equal(
yield* fs.readFileString(path.join(existingPath, "keep.txt")),
"existing uncommitted work",
);
}
assert.equal(harness.removeWorktree.mock.calls.length, 0);
}).pipe(Effect.provide(harness.layer));
}).pipe(
Effect.provide(
GitVcsDriver.layer.pipe(
Layer.provide(
ServerConfig.layerTest(process.cwd(), { prefix: "launch-collision-test-" }),
),
Layer.provideMerge(NodeServices.layer),
),
),
),
);

it.effect.each([
{ branch: "t3/abcd1234", flat: false },
{ branch: "t3-abcd1234", flat: true },
])(
"keeps temporary naming and background rename after a collision on $branch",
({ branch, flat }) => {
let attempts = 0;
const harness = makeHarness({
hasCommit: () => Effect.succeed(flat),
createWorktree: (input) =>
++attempts === 1
? Effect.fail(
new GitCommandError({
operation: "GitVcsDriver.createWorktree",
command: "git",
cwd: project.workspaceRoot,
detail: "git worktree add failed",
reason: "branch_already_exists",
exitCode: 128,
}),
)
: Effect.succeed({
worktree: { path: "/repo-worktrees/fresh", refName: input.newRefName! },
}),
});
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const command = `command:launch:temp-collision-${flat}`;
const launched = yield* launches.launch(
launchInput({
command,
thread: `thread:launch:temp-collision-${flat}`,
message: "Name this worktree",
workspace: { type: "worktree", baseRef: "main", branch },
}),
);
yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe(
Stream.filter((stored) => stored.commandId === CommandId.make(`${command}:branch-rename`)),
Stream.runHead,
);
assert.equal(harness.createWorktree.mock.calls.length, 2);
const replacement = harness.createWorktree.mock.calls[1]![0].newRefName!;
assert.notEqual(replacement, branch);
assert.match(replacement, flat ? /^t3-[0-9a-f]{8}$/u : /^t3\/[0-9a-f]{8}$/u);
assert.equal(harness.renameBranch.mock.calls[0]![0].oldBranch, replacement);
assert.equal(
(yield* threads.getThreadProjection(launched.threadId)).thread.branch,
"generated-branch",
);
assert.equal(harness.removeWorktree.mock.calls.length, 0);
}).pipe(Effect.provide(harness.layer));
},
);

it.effect.each(["path_already_exists", "authentication_failed", "branch_already_exists"] as const)(
"handles %s without unbounded attempts",
(reason) => {
let attempts = 0;
const harness = makeHarness({
createWorktree: (input) => {
attempts += 1;
return reason === "path_already_exists" && attempts > 1
? Effect.succeed({
worktree: { path: "/repo-worktrees/fresh", refName: input.newRefName! },
})
: Effect.fail(
new GitCommandError({
operation: "GitVcsDriver.createWorktree",
command: "git",
cwd: project.workspaceRoot,
detail: "git worktree add failed",
reason,
exitCode: 128,
}),
);
},
});
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const command = `command:launch:collision-limit-${reason}`;
const launched = yield* launches.launch(
launchInput({
command,
thread: `thread:launch:collision-limit-${reason}`,
message: "Prepare worktree",
workspace: { type: "worktree", baseRef: "main", branch: "my-feature" },
}),
);
const succeeds = reason === "path_already_exists";
yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe(
Stream.filter(
(stored) =>
stored.commandId === CommandId.make(`${command}:${succeeds ? "release" : "fail"}`) &&
stored.event.type === "run.updated",
),
Stream.runHead,
);
assert.equal(attempts, succeeds ? 2 : reason === "authentication_failed" ? 1 : 101);
const projection = yield* threads.getThreadProjection(launched.threadId);
assert.equal(projection.runs[0]?.status, succeeds ? "starting" : "failed");
assert.equal(harness.runSetup.mock.calls.length, succeeds ? 1 : 0);
assert.equal(harness.removeWorktree.mock.calls.length, 0);
}).pipe(Effect.provide(harness.layer));
},
);

it.effect("Retry recovers an existing failed run whose requested branch is taken", () => {
let fails = true;
const harness = makeHarness({
createWorktree: (input) =>
fails || input.newRefName === "my-feature"
? Effect.fail(
new GitCommandError({
operation: "GitVcsDriver.createWorktree",
command: "git",
cwd: project.workspaceRoot,
detail: "git worktree add failed",
reason: "branch_already_exists",
exitCode: 128,
}),
)
: Effect.succeed({
worktree: { path: "/repo-worktrees/recovered", refName: input.newRefName! },
}),
});
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const launched = yield* launches.launch(
launchInput({
command: "command:launch:taken-retry",
thread: "thread:launch:taken-retry",
message: "Retry me",
workspace: { type: "worktree", baseRef: "main", branch: "my-feature" },
}),
);
yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe(
Stream.filter(
(stored) =>
stored.commandId === CommandId.make("command:launch:taken-retry:fail") &&
stored.event.type === "run.updated",
),
Stream.runHead,
);
const failed = yield* threads.getThreadProjection(launched.threadId);
fails = false;
yield* launches.retryPreparation({
commandId: CommandId.make("command:launch:taken-retry:retry"),
threadId: launched.threadId,
runId: failed.runs[0]!.id,
});
yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe(
Stream.filter(
(stored) =>
stored.commandId === CommandId.make("command:launch:taken-retry:retry:release") &&
stored.event.type === "run.updated",
),
Stream.runHead,
);
const recovered = yield* threads.getThreadProjection(launched.threadId);
assert.equal(recovered.runs.length, 1);
assert.equal(recovered.runs[0]!.id, failed.runs[0]!.id);
assert.equal(recovered.messages.length, 1);
assert.equal(recovered.runs[0]?.status, "starting");
assert.equal(recovered.thread.branch, "my-feature-1");
assert.equal(recovered.thread.worktreePath, "/repo-worktrees/recovered");
assert.equal(recovered.turnItems.find((item) => item.type === "error")?.status, "cancelled");
}).pipe(Effect.provide(harness.layer));
});

it.effect("keeps the temporary branch when branch generation fails", () =>
Effect.gen(function* () {
const harness = makeHarness({
Expand Down Expand Up @@ -1484,39 +1769,56 @@ it.effect("a retry reuses a recorded worktree without undoing its branch rename"
}).pipe(Effect.provide(harness.layer));
});

it.effect("removes a worktree that failed before the thread recorded it", () => {
const harness = makeHarness({
// A checkout that dies after claiming its directory.
createWorktree: (_input, options) =>
(options?.progress?.onWorktreeClaimed?.("/repo-worktrees/partial") ?? Effect.void).pipe(
Effect.andThen(Effect.fail(new Error("checkout failed") as never)),
),
});
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const launched = yield* launches.launch(
launchInput({
command: "command:launch:partial-worktree",
thread: "thread:launch:partial-worktree",
message: "Partial checkout",
workspace: { type: "worktree", baseRef: "main" },
}),
);
yield* waitUntil(() =>
threads
.getThreadProjection(launched.threadId)
.pipe(Effect.map((projection) => projection.runs[0]?.status === "failed")),
);
const projection = yield* threads.getThreadProjection(launched.threadId);
// Unrecorded, so a retry would create a second checkout beside it.
assert.equal(projection.thread.worktreePath, null);
assert.deepEqual(
harness.removeWorktree.mock.calls.map(([input]) => input.path),
["/repo-worktrees/partial"],
);
}).pipe(Effect.provide(harness.layer));
});
it.effect.each([undefined, "branch_already_exists"] as const)(
"removes a claimed worktree after failure %s without creating another",
(reason) => {
const harness = makeHarness({
// A checkout that dies after claiming its directory.
createWorktree: (_input, options) =>
(options?.progress?.onWorktreeClaimed?.("/repo-worktrees/partial") ?? Effect.void).pipe(
Effect.andThen(
Effect.fail(
new GitCommandError({
operation: "GitVcsDriver.createWorktree.configureBaseRef",
command: "git",
cwd: project.workspaceRoot,
detail: "checkout failed",
...(reason === undefined ? {} : { reason }),
}),
),
),
),
});
return Effect.gen(function* () {
const launches = yield* ThreadLaunch.ThreadLaunchService;
const threads = yield* ThreadManagement.ThreadManagementService;
const launched = yield* launches.launch(
launchInput({
command: "command:launch:partial-worktree",
thread: "thread:launch:partial-worktree",
message: "Partial checkout",
workspace: { type: "worktree", baseRef: "main" },
}),
);
yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe(
Stream.filter(
(stored) =>
stored.commandId === CommandId.make("command:launch:partial-worktree:fail") &&
stored.event.type === "run.updated",
),
Stream.runHead,
);
const projection = yield* threads.getThreadProjection(launched.threadId);
// Unrecorded, so a retry would create a second checkout beside it.
assert.equal(projection.thread.worktreePath, null);
assert.equal(harness.createWorktree.mock.calls.length, 1);
assert.deepEqual(
harness.removeWorktree.mock.calls.map(([input]) => input.path),
["/repo-worktrees/partial"],
);
}).pipe(Effect.provide(harness.layer));
},
);

it.effect.each(["worktree", "setup"] as const)(
"%s failure keeps the thread and message visible and emits failure items",
Expand Down
Loading
Loading