From 1ec8dac1d1bd123686a8d243d9b2eb426d6d267f Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sat, 3 Oct 2026 06:38:53 +0200 Subject: [PATCH 1/6] fix(server): clean up failed launches and settle empty threads --- .../ThreadLaunchService.test.ts | 311 ++++++++++++++++-- .../orchestration-v2/ThreadLaunchService.ts | 135 ++++++-- .../ThreadSettlementService.test.ts | 34 +- .../ThreadSettlementService.ts | 17 +- .../project/ProjectSetupScriptRunner.test.ts | 58 ++++ .../src/project/ProjectSetupScriptRunner.ts | 12 +- docs/user/thread-sidebar.md | 10 +- 7 files changed, 513 insertions(+), 64 deletions(-) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 10eb4ac9485e..8f9338c13e42 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -97,6 +97,8 @@ const adapter = { interface HarnessOptions { readonly managedFolders?: Layer.Layer; readonly createWorktree?: GitWorkflow.GitWorkflowService["Service"]["createWorktree"]; + readonly removeWorktree?: GitWorkflow.GitWorkflowService["Service"]["removeWorktree"]; + readonly closeTerminal?: TerminalManager.TerminalManager["Service"]["close"]; readonly fetchRemote?: GitWorkflow.GitWorkflowService["Service"]["fetchRemote"]; readonly renameBranch?: GitWorkflow.GitWorkflowService["Service"]["renameBranch"]; readonly runSetup?: ProjectSetupScriptRunner.ProjectSetupScriptRunner["Service"]["runForThread"]; @@ -127,6 +129,8 @@ function makeHarness(options: HarnessOptions = {}) { const renameBranch = vi.fn( options.renameBranch ?? ((input) => Effect.succeed({ branch: input.newBranch })), ); + const removeWorktree = vi.fn(options.removeWorktree ?? (() => Effect.void)); + const closeTerminal = vi.fn(options.closeTerminal ?? (() => Effect.void)); const runSetup = vi.fn( options.runSetup ?? (() => Effect.succeed({ status: "no-script" as const })), ); @@ -139,7 +143,7 @@ function makeHarness(options: HarnessOptions = {}) { const externalServices = Layer.mergeAll( WorktreeSetupTracker.layer, Layer.mock(ProjectCloneTracker.ProjectCloneTracker)({ get: () => Effect.succeed(null) }), - Layer.mock(TerminalManager.TerminalManager)({ close: () => Effect.void }), + Layer.mock(TerminalManager.TerminalManager)({ close: closeTerminal }), Layer.succeed(ProjectService.ProjectService, { create: () => Effect.die("unused"), bootstrap: () => Effect.die("unused"), @@ -164,7 +168,7 @@ function makeHarness(options: HarnessOptions = {}) { fetchRemote: options.fetchRemote ?? (() => Effect.void), remoteExists: () => Effect.succeed(true), remoteBranchExists: () => Effect.succeed(true), - removeWorktree: () => Effect.void, + removeWorktree, resolveRemoteTrackingCommit: () => Effect.succeed({ commitSha: "remote-main-sha", remoteRefName: "origin/main" }), }), @@ -220,6 +224,8 @@ function makeHarness(options: HarnessOptions = {}) { externalServices, ), createWorktree, + removeWorktree, + closeTerminal, renameBranch, generateBranchName, generateThreadTitle, @@ -531,9 +537,9 @@ it.effect("enqueues provider work only after setup has been initiated", () => }), ); -it.effect( - "queues follow-up messages behind preparation and checkpoints them in the final workspace", - () => +it.effect.each(["success", "failure"] as const)( + "queues follow-ups behind preparation and checkpoints the bound workspace after %s", + (setupResult) => Effect.gen(function* () { const setupEntered = yield* Deferred.make(); const failSetup = yield* Deferred.make(); @@ -541,7 +547,11 @@ it.effect( runSetup: () => Deferred.succeed(setupEntered, undefined).pipe( Effect.andThen(Deferred.await(failSetup)), - Effect.andThen(Effect.fail(new Error("setup failed") as never)), + Effect.andThen( + setupResult === "failure" + ? Effect.fail(new Error("setup failed") as never) + : Effect.succeed({ status: "no-script" as const }), + ), ), }); yield* Effect.gen(function* () { @@ -555,6 +565,8 @@ it.effect( workspace: { type: "worktree", baseRef: "main" }, }), ); + const initialRun = launched.projection.runs[0]; + assert.ok(initialRun); yield* Deferred.await(setupEntered); const followUp = yield* threads.sendToThread({ @@ -578,26 +590,34 @@ it.effect( ); yield* Deferred.succeed(failSetup, undefined); - yield* waitUntil(() => - threads - .getThreadProjection(launched.threadId) - .pipe( - Effect.map( - (projection) => - projection.runs.find((run) => run.id === followUp.run.id)?.status === "starting", - ), - ), + const startedRunId = setupResult === "failure" ? followUp.run.id : initialRun.id; + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && + stored.event.payload.id === startedRunId && + stored.event.payload.status === "starting", + ), + Stream.runHead, ); const projection = yield* threads.getThreadProjection(launched.threadId); const rootNode = projection.nodes.find( - (node) => node.runId === followUp.run.id && node.kind === "root_turn", + (node) => node.runId === startedRunId && node.kind === "root_turn", + ); + assert.equal( + projection.runs.find((run) => run.id === followUp.run.id)?.status, + setupResult === "failure" ? "starting" : "queued", ); assert.isNotNull(rootNode?.checkpointScopeId); + assert.equal( + projection.thread.worktreePath, + setupResult === "failure" ? null : "/repo-worktrees/feature", + ); assert.equal( projection.checkpointScopes.find((scope) => scope.id === rootNode?.checkpointScopeId) ?.cwd, - "/repo-worktrees/feature", + setupResult === "failure" ? process.cwd() : "/repo-worktrees/feature", ); }).pipe(Effect.provide(harness.layer)); }), @@ -1313,14 +1333,19 @@ it.effect.each(["worktree", "setup"] as const)( workspace: { type: "worktree", baseRef: "main" }, }); const launched = yield* launches.launch(input); - yield* waitUntil(() => - threads - .getThreadProjection(launched.threadId) - .pipe(Effect.map((projection) => projection.runs[0]?.status === "failed")), + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, ); const projection = yield* threads.getThreadProjection(launched.threadId); assert.equal(projection.messages[0]?.text, `Fail during ${failurePoint}`); assert.equal(projection.runs[0]?.status, "failed"); + assert.isNull(projection.thread.worktreePath); + assert.isNull(projection.thread.branch); + assert.equal(harness.removeWorktree.mock.calls.length, failurePoint === "setup" ? 1 : 0); assert.equal( projection.turnItems.find((item) => item.type === "command_execution")?.status, "failed", @@ -1382,6 +1407,185 @@ it.effect("replays a server-allocated launch", () => }), ); +it.effect.each(["worktree", "existing_worktree"] as const)( + "closes a failed setup terminal and only removes the owned %s workspace", + (workspaceType) => { + const harness = makeHarness({ + runSetup: () => + Effect.succeed({ + status: "started" as const, + async: false, + scriptId: "setup", + scriptName: "Setup", + scriptCommand: "vp install", + terminalId: "setup", + cwd: "/repo-worktrees/feature", + completion: Effect.succeed({ exitCode: 1, durationMs: 1 }), + }), + }); + return Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const input = launchInput({ + command: `launch:failed-terminal:${workspaceType}`, + thread: `thread:failed-terminal:${workspaceType}`, + message: "Start", + workspace: + workspaceType === "worktree" + ? { type: "worktree", baseRef: "main", branch: "feature" } + : { type: "existing_worktree", worktreePath: "/existing", branch: "existing" }, + }); + const launched = yield* launches.launch(input); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.equal(projection.runs[0]?.status, "failed"); + assert.deepEqual( + harness.closeTerminal.mock.calls.map(([input]) => input), + [{ threadId: launched.threadId, terminalId: "setup", deleteHistory: true }], + ); + assert.equal(harness.removeWorktree.mock.calls.length, workspaceType === "worktree" ? 1 : 0); + assert.equal( + projection.thread.worktreePath, + workspaceType === "worktree" ? null : "/existing", + ); + assert.equal(projection.thread.branch, workspaceType === "worktree" ? null : "existing"); + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + assert.isNull((yield* tracker.get(launched.threadId))?.setupScript ?? null); + }).pipe(Effect.provide(harness.layer)); + }, +); + +it.effect.each([true, false])( + "cleans a partially claimed worktree with removal success %s", + (removed) => { + const harness = makeHarness({ + createWorktree: (_input, options) => + (options?.progress?.onWorktreeClaimed?.("/claimed") ?? Effect.void).pipe( + Effect.andThen(Effect.fail(new Error("checkout failed") as never)), + ), + removeWorktree: () => + removed ? Effect.void : Effect.fail(new Error("remove failed") as never), + }); + return Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const launched = yield* launches.launch( + launchInput({ + command: "launch:claimed-worktree", + thread: "thread:claimed-worktree", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }), + ); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + assert.deepEqual( + harness.removeWorktree.mock.calls.map(([input]) => input), + [{ cwd: "/repo", path: "/claimed", force: true }], + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.equal(projection.thread.worktreePath, removed ? null : "/claimed"); + assert.equal(projection.thread.branch, removed ? null : "feature"); + }).pipe(Effect.provide(harness.layer)); + }, +); + +it.effect("keeps a failed launch worktree bound when cleanup cannot remove it", () => { + const harness = makeHarness({ + runSetup: () => Effect.fail(new Error("setup failed") as never), + removeWorktree: () => Effect.fail(new Error("remove failed") as never), + }); + return Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + const launched = yield* launches.launch( + launchInput({ + command: "launch:cleanup-failed", + thread: "thread:cleanup-failed", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }), + ); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.equal(projection.thread.worktreePath, "/repo-worktrees/feature"); + assert.equal(projection.thread.branch, "feature"); + assert.equal((yield* tracker.get(launched.threadId))?.worktreePath, "/repo-worktrees/feature"); + assert.equal(projection.runs[0]?.status, "failed"); + }).pipe(Effect.provide(harness.layer)); +}); + +it.effect("stops branch renaming before removing a failed launch worktree", () => + Effect.gen(function* () { + const renameEntered = yield* Deferred.make(); + const setupCompletion = yield* Deferred.make<{ exitCode: number; durationMs: number }>(); + const renameInterrupted = yield* Ref.make(false); + const harness = makeHarness({ + renameBranch: () => + Deferred.succeed(renameEntered, undefined).pipe( + Effect.andThen(Effect.never), + Effect.onInterrupt(() => Ref.set(renameInterrupted, true)), + ), + runSetup: () => + Effect.succeed({ + status: "started" as const, + async: false, + scriptId: "setup", + scriptName: "Setup", + scriptCommand: "vp install", + terminalId: "setup", + cwd: "/repo-worktrees/feature", + completion: Deferred.await(setupCompletion), + }), + }); + yield* Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const launched = yield* launches.launch( + launchInput({ + command: "launch:rename-failure", + thread: "thread:rename-failure", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "t3code/abcd1234" }, + }), + ); + yield* Deferred.await(renameEntered); + yield* Deferred.succeed(setupCompletion, { exitCode: 1, durationMs: 1 }); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + assert.isTrue(yield* Ref.get(renameInterrupted)); + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.isNull(projection.thread.worktreePath); + assert.isNull(projection.thread.branch); + assert.equal(harness.removeWorktree.mock.calls.length, 1); + assert.isNull((yield* tracker.get(launched.threadId))?.setupScript); + }).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* () { @@ -1961,7 +2165,17 @@ it.effect("cancels tracked setup before provider work is released", () => Effect.gen(function* () { const entered = yield* Deferred.make(); const harness = makeHarness({ - runSetup: () => Deferred.succeed(entered, undefined).pipe(Effect.andThen(Effect.never)), + runSetup: () => + Effect.succeed({ + status: "started" as const, + async: false, + scriptId: "setup", + scriptName: "Setup", + scriptCommand: "vp install", + terminalId: "setup", + cwd: "/repo-worktrees/feature", + completion: Deferred.succeed(entered, undefined).pipe(Effect.andThen(Effect.never)), + }), }); yield* Effect.gen(function* () { const launches = yield* ThreadLaunch.ThreadLaunchService; @@ -1979,14 +2193,67 @@ it.effect("cancels tracked setup before provider work is released", () => assert.equal((yield* tracker.get(launched.threadId))?.phase, "running"); assert.isTrue(yield* tracker.cancel(launched.threadId)); assert.equal((yield* tracker.get(launched.threadId))?.phase, "cancelled"); + assert.isNull((yield* tracker.get(launched.threadId))?.setupScript); const projection = yield* threads.getThreadProjection(launched.threadId); assert.equal(projection.runs[0]?.status, "failed"); assert.isNull(projection.thread.worktreePath); + assert.isNull(projection.thread.branch); + assert.equal(harness.removeWorktree.mock.calls.length, 1); + assert.deepEqual( + harness.closeTerminal.mock.calls.map(([input]) => input), + [{ threadId: launched.threadId, terminalId: "setup", deleteHistory: true }], + ); assert.isEmpty(yield* outbox.listByCommandId(CommandId.make(`${input.commandId}:release`))); }).pipe(Effect.provide(harness.layer)); }), ); +it.effect("keeps a released worktree when shutdown interrupts its async setup", () => { + const harness = makeHarness({ + runSetup: () => + Effect.succeed({ + status: "started" as const, + async: true, + scriptId: "setup", + scriptName: "Setup", + scriptCommand: "vp install", + terminalId: "setup", + cwd: "/repo-worktrees/feature", + completion: Effect.never, + }), + }); + return Effect.gen(function* () { + yield* Effect.scoped( + Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + const threads = yield* ThreadManagement.ThreadManagementService; + const launched = yield* launches.launch( + launchInput({ + command: "launch:released-shutdown", + thread: "thread:released-shutdown", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }), + ); + yield* tracker.stream(launched.threadId).pipe( + Stream.filter( + (snapshot) => + snapshot?.stages.some((stage) => stage.id === "agent" && stage.status === "done") === + true, + ), + Stream.runHead, + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.equal(projection.runs[0]?.status, "starting"); + assert.equal(projection.thread.worktreePath, "/repo-worktrees/feature"); + }).pipe(Effect.provide(harness.layer)), + ); + assert.isEmpty(harness.removeWorktree.mock.calls); + assert.isEmpty(harness.closeTerminal.mock.calls); + }); +}); + it.effect.each([0, 1])("releases an async setup before its completion with exit %s", (exitCode) => Effect.gen(function* () { const completion = yield* Deferred.make<{ exitCode: number | null; durationMs: number }>(); diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index faf613376920..9e942abaf788 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -23,6 +23,7 @@ import * as Cause from "effect/Cause"; import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Ref from "effect/Ref"; @@ -218,7 +219,10 @@ const make = Effect.gen(function* () { const tracked = input.workspaceStrategy.type === "worktree"; let createdWorktreePath: string | null = null; + let createdWorktreeBranch: string | null = null; let setupTerminalId: string | null = null; + let branchRenameFiber: Fiber.Fiber | null = null; + let preparationComplete = false; if (tracked) { yield* setupTracker.begin({ threadId, @@ -267,19 +271,18 @@ const make = Effect.gen(function* () { // The server owns worktree naming: without an explicit branch, provision // under a temporary `t3code/` name so the worktree never waits on // name generation, then rename in the background below. - const requestedBranch = input.workspaceStrategy.branch; - let branch: string | null; - if (input.workspaceStrategy.type === "worktree" && requestedBranch === undefined) { - const uuid = yield* randomUuidV4; - branch = buildTemporaryWorktreeBranchName(() => uuid.replaceAll("-", "")); - } else { - branch = requestedBranch ?? null; - } + let branch = input.workspaceStrategy.branch ?? null; let worktreePath = input.workspaceStrategy.type === "existing_worktree" ? input.workspaceStrategy.worktreePath : null; if (input.workspaceStrategy.type === "worktree") { + const worktreeBranch = + branch ?? + (yield* randomUuidV4.pipe( + Effect.map((uuid) => buildTemporaryWorktreeBranchName(() => uuid.replaceAll("-", ""))), + )); + branch = worktreeBranch; if (runId !== null) { yield* threads .dispatch({ @@ -335,7 +338,7 @@ const make = Effect.gen(function* () { { cwd: project.workspaceRoot, refName: startRef, - newRefName: branch!, + newRefName: worktreeBranch, baseRefName: input.workspaceStrategy.baseRef, path: null, }, @@ -344,6 +347,7 @@ const make = Effect.gen(function* () { onWorktreeClaimed: (path) => Effect.sync(() => { createdWorktreePath = path; + createdWorktreeBranch = worktreeBranch; }), onCheckoutProgress: (progress) => setupTracker.stage(threadId, "checkout", { percent: progress.percent }), @@ -354,6 +358,7 @@ const make = Effect.gen(function* () { worktreePath = worktree.worktree.path; branch = worktree.worktree.refName; createdWorktreePath = worktreePath; + createdWorktreeBranch = branch; yield* setupTracker.update(threadId, (snapshot) => ({ ...snapshot, worktreePath, branch })); yield* setupTracker.stageStatus(threadId, "checkout", "done"); } @@ -380,7 +385,7 @@ const make = Effect.gen(function* () { ) { const oldBranch = branch; const worktreeCwd = worktreePath; - yield* generateBranchNameFor(worktreeCwd, initialMessage).pipe( + branchRenameFiber = yield* generateBranchNameFor(worktreeCwd, initialMessage).pipe( Effect.flatMap(({ branch: newBranch, exactName }) => git.renameBranch({ cwd: worktreeCwd, @@ -456,9 +461,10 @@ const make = Effect.gen(function* () { terminalId: setup.terminalId, }, })); - if (setup.completion) { + const setupCompletion = setup.completion; + if (setupCompletion) { const awaitCompletion = Effect.gen(function* () { - const completion = yield* setup.completion!; + const completion = yield* setupCompletion; yield* setupTracker.stage(threadId, "setup-script", { status: completion.exitCode === 0 ? "done" : "failed", detail: `exited with ${completion.exitCode ?? "no exit code"}`, @@ -498,7 +504,17 @@ const make = Effect.gen(function* () { threadId, runId, }) - .pipe(Effect.mapError(mapError(input, "release-run", threadId))); + .pipe( + Effect.mapError(mapError(input, "release-run", threadId)), + Effect.tap(() => + Effect.sync(() => { + preparationComplete = true; + }), + ), + Effect.uninterruptible, + ); + } else { + preparationComplete = true; } yield* setupTracker.stageStatus(threadId, "agent", "done"); yield* awaitAsyncSetup; @@ -507,33 +523,92 @@ const make = Effect.gen(function* () { Effect.onError((cause) => Effect.gen(function* () { const cancelled = Cause.hasInterruptsOnly(cause); - yield* setupTracker.finish( - threadId, - cancelled ? "cancelled" : "failed", - cancelled ? null : failureDetail(Cause.squash(cause)), - ); - if (cancelled && tracked && createdWorktreePath) { - if (setupTerminalId) + // Once the prepared run is released, the provider owns its workspace. + // An async setup interrupted during shutdown must not delete it. + if (!preparationComplete) { + if (branchRenameFiber !== null) yield* Fiber.interrupt(branchRenameFiber); + if (setupTerminalId !== null) yield* terminals .close({ threadId, terminalId: setupTerminalId, deleteHistory: true }) - .pipe(Effect.ignore); - yield* git + .pipe( + Effect.tap(() => + setupTracker.update(threadId, (snapshot) => ({ + ...snapshot, + setupScript: null, + })), + ), + Effect.catchCause((cleanupCause) => + Effect.logWarning("Failed to close thread launch setup terminal", { + threadId, + terminalId: setupTerminalId, + cause: cleanupCause, + }), + ), + ); + } + if (!preparationComplete && tracked && createdWorktreePath !== null) { + const removed = yield* git .removeWorktree({ cwd: project.workspaceRoot, path: createdWorktreePath, force: true, }) - .pipe(Effect.ignore); - yield* threads - .dispatch({ - type: "thread.metadata.update", - commandId: CommandId.make(`${input.commandId}:cancel-workspace`), - threadId, + .pipe( + Effect.as(true), + Effect.catchCause((cleanupCause) => + Effect.logWarning("Failed to remove thread launch worktree", { + threadId, + worktreePath: createdWorktreePath, + cause: cleanupCause, + }).pipe(Effect.as(false)), + ), + ); + // Keep the binding when removal fails so the surviving worktree + // remains discoverable instead of becoming an orphan. + if (removed) { + yield* threads + .dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), + threadId, + worktreePath: null, + branch: null, + }) + .pipe(Effect.ignore); + yield* setupTracker.update(threadId, (snapshot) => ({ + ...snapshot, worktreePath: null, branch: null, - }) - .pipe(Effect.ignore); + })); + } else { + const shell = yield* threads + .getThreadShell(threadId) + .pipe(Effect.catchCause(() => Effect.succeed(null))); + // Checkout can fail after claiming a path but before the normal + // workspace update. Bind that survivor so it can still be found. + if (shell !== null && shell.worktreePath !== createdWorktreePath) { + yield* threads + .dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), + threadId, + worktreePath: createdWorktreePath, + branch: createdWorktreeBranch, + }) + .pipe(Effect.ignore); + yield* setupTracker.update(threadId, (snapshot) => ({ + ...snapshot, + worktreePath: createdWorktreePath, + branch: createdWorktreeBranch, + })); + } + } } + yield* setupTracker.finish( + threadId, + cancelled ? "cancelled" : "failed", + cancelled ? null : failureDetail(Cause.squash(cause)), + ); }), ), ); diff --git a/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts b/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts index 03933d9d0048..6752e199a0fa 100644 --- a/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts @@ -251,8 +251,40 @@ describe("resolveAutoSettlementAt", () => { expect( ThreadSettlementService.resolveAutoSettlementAt({ ...input, autoSettleAfterDays: null }), ).toBeNull(); + expect(ThreadSettlementService.resolveAutoSettlementAt({ ...input, thread: shell() })).toEqual( + shell().createdAt, + ); + }); + + it("uses creation time for empty threads only after the inactivity window", () => { + const input = { + thread: shell({ createdAt: at(-3 * DAY_MS) }), + pullRequest: null, + nowMs: NOW_MS, + autoSettleAfterDays: 2, + autoSettleOnMerge: true, + }; + expect(ThreadSettlementService.resolveAutoSettlementAt(input)).toEqual(input.thread.createdAt); expect( - ThreadSettlementService.resolveAutoSettlementAt({ ...input, thread: shell() }), + ThreadSettlementService.resolveAutoSettlementAt({ + ...input, + thread: shell({ createdAt: at(-DAY_MS) }), + }), + ).toBeNull(); + expect( + ThreadSettlementService.resolveAutoSettlementAt({ + ...input, + thread: shell({ createdAt: at(-2 * DAY_MS) }), + }), + ).toBeNull(); + expect( + ThreadSettlementService.resolveAutoSettlementAt({ ...input, autoSettleAfterDays: null }), + ).toBeNull(); + expect( + ThreadSettlementService.resolveAutoSettlementAt({ + ...input, + thread: shell({ createdAt: at(-3 * DAY_MS), pinnedAt: at(-DAY_MS) }), + }), ).toBeNull(); }); diff --git a/apps/server/src/orchestration-v2/ThreadSettlementService.ts b/apps/server/src/orchestration-v2/ThreadSettlementService.ts index 74c9de8d6b37..f828249894c4 100644 --- a/apps/server/src/orchestration-v2/ThreadSettlementService.ts +++ b/apps/server/src/orchestration-v2/ThreadSettlementService.ts @@ -187,16 +187,17 @@ export function resolveAutoSettlementAt(input: { }; } if (!isAutoSettlementCandidate(thread, input.nowMs)) return null; - const activityAtMs = latestMillis([ - toMillis(thread.latestUserMessageAt), - toMillis(thread.latestRunRequestedAt), - toMillis(thread.latestRunStartedAt), - toMillis(thread.latestRunCompletedAt), - ]); + const activityAtMs = + latestMillis([ + toMillis(thread.latestUserMessageAt), + toMillis(thread.latestRunRequestedAt), + toMillis(thread.latestRunStartedAt), + toMillis(thread.latestRunCompletedAt), + ]) ?? DateTime.toEpochMillis(thread.createdAt); if (pullRequest !== null && pullRequestSettles(thread, pullRequest, input.autoSettleOnMerge)) { - return activityAtMs === null ? thread.createdAt : DateTime.makeUnsafe(activityAtMs); + return DateTime.makeUnsafe(activityAtMs); } - if (input.autoSettleAfterDays === null || activityAtMs === null) return null; + if (input.autoSettleAfterDays === null) return null; return activityAtMs < input.nowMs - input.autoSettleAfterDays * DAY_MS ? DateTime.makeUnsafe(activityAtMs) : null; diff --git a/apps/server/src/project/ProjectSetupScriptRunner.test.ts b/apps/server/src/project/ProjectSetupScriptRunner.test.ts index 85af1b3e3438..783371b7c04d 100644 --- a/apps/server/src/project/ProjectSetupScriptRunner.test.ts +++ b/apps/server/src/project/ProjectSetupScriptRunner.test.ts @@ -1,6 +1,7 @@ import { assert, it, vi } from "@effect/vitest"; import { ProjectId } from "@t3tools/contracts"; import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -117,3 +118,60 @@ it.effect("resolves setup scripts through the standalone project service", () => yield* listener({ type: "closed", threadId: "thread-1", terminalId: "setup-setup" }); }).pipe(Effect.provide(layer)); }); + +it.effect.each(["failure", "defect"] as const)( + "closes the setup terminal and unsubscribes when writing ends with a %s", + (failureType) => { + const writeFailure = new TerminalManager.TerminalWriteError({ + threadId: "thread-failed", + terminalId: "setup-setup", + terminalPid: 123, + cause: new Error("write failed"), + }); + const close = vi.fn(() => Effect.void); + const unsubscribe = vi.fn(); + const layer = ProjectSetupScriptRunner.layer.pipe( + Layer.provide( + Layer.mergeAll( + Layer.mock(ProjectService.ProjectService)({}), + Layer.mock(TerminalManager.TerminalManager)({ + open: () => Effect.succeed({} as never), + write: () => + failureType === "failure" ? Effect.fail(writeFailure) : Effect.die(writeFailure), + close, + subscribe: () => Effect.succeed(unsubscribe), + }), + ServerSettings.layerTest(), + ), + ), + ); + return Effect.gen(function* () { + const runner = yield* ProjectSetupScriptRunner.ProjectSetupScriptRunner; + const result = yield* runner + .runForThread({ + threadId: "thread-failed", + worktreePath: "/repo-worktree", + project: { + id: ProjectId.make("project:setup-write-failure"), + workspaceRoot: "/repo", + scripts: [ + { + id: "setup", + name: "Setup", + command: "vp install", + icon: "configure", + runOnWorktreeCreate: true, + }, + ], + }, + observeCompletion: {}, + }) + .pipe(Effect.exit); + assert.isTrue(Exit.isFailure(result)); + assert.equal(unsubscribe.mock.calls.length, 1); + assert.deepEqual(close.mock.calls[0], [ + { threadId: "thread-failed", terminalId: "setup-setup", deleteHistory: true }, + ]); + }).pipe(Effect.provide(layer)); + }, +); diff --git a/apps/server/src/project/ProjectSetupScriptRunner.ts b/apps/server/src/project/ProjectSetupScriptRunner.ts index f97c0f4df9c0..6d8dc0c1ebd5 100644 --- a/apps/server/src/project/ProjectSetupScriptRunner.ts +++ b/apps/server/src/project/ProjectSetupScriptRunner.ts @@ -424,8 +424,16 @@ export const make = Effect.gen(function* () { cause, }), ), - // Nothing will ever settle the completion if the command never ran. - Effect.tapError(() => Effect.sync(() => observed?.unsubscribe())), + // No caller receives the terminal id when writing fails, so unwind + // the shell here before launch cleanup removes its workspace. + Effect.onError(() => + Effect.sync(() => observed?.unsubscribe()).pipe( + Effect.andThen( + terminalManager.close({ threadId: input.threadId, terminalId, deleteHistory: true }), + ), + Effect.ignoreCause, + ), + ), ); // A clean run leaves only an idle prompt behind; its output stays in the diff --git a/docs/user/thread-sidebar.md b/docs/user/thread-sidebar.md index 1172678440c7..9f2eb47aa392 100644 --- a/docs/user/thread-sidebar.md +++ b/docs/user/thread-sidebar.md @@ -13,6 +13,13 @@ an existing worktree, use **New thread in this worktree** from the branch toolba When you change a new thread's project, T3 Code stays in the current environment if that project exists there. Otherwise it selects an environment that has it. +If preparing a new worktree fails before the agent starts, T3 Code closes its setup +terminal and removes the new worktree. The thread and its messages stay visible. +Messages queued during preparation then run from the project's root directory. +To try again in a separate worktree, start a new thread with **New worktree**. +An existing worktree you selected stays in place. If removing a failed worktree +fails, it stays linked to the thread so you can inspect or remove it. + ### Start without a project A thread does not need a project. To start one without a project, click **or @@ -131,7 +138,8 @@ runs a command, such as a dev server, stays open. By default, environments settle inactive threads after three days and settle threads whose pull request merged. A closed pull request can also settle an idle -thread. Work in progress, pending questions or approvals, and live background work +thread. A thread with no messages or runs becomes inactive from its creation time. +Work in progress, pending questions or approvals, and live background work prevent automatic settlement. An open pull request does not prevent inactivity settlement, but an old closed or merged pull request does not settle work you resumed after it closed. From d227e96761c4685b15f4b6c5376be10f0a7e223d Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sat, 3 Oct 2026 06:57:47 +0200 Subject: [PATCH 2/6] fix(server): report launch cleanup binding failures --- .../src/orchestration-v2/ThreadLaunchService.ts | 15 ++++++++++++--- 1 file changed, 12 insertions(+), 3 deletions(-) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index 9e942abaf788..bc3a7fdde2de 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -581,9 +581,18 @@ const make = Effect.gen(function* () { branch: null, })); } else { - const shell = yield* threads - .getThreadShell(threadId) - .pipe(Effect.catchCause(() => Effect.succeed(null))); + const shell = yield* threads.getThreadShell(threadId).pipe( + Effect.catchCause((cleanupCause) => + Effect.logWarning( + "Failed to read thread launch worktree binding during cleanup", + { + threadId, + worktreePath: createdWorktreePath, + cause: cleanupCause, + }, + ).pipe(Effect.as(null)), + ), + ); // Checkout can fail after claiming a path but before the normal // workspace update. Bind that survivor so it can still be found. if (shell !== null && shell.worktreePath !== createdWorktreePath) { From 3eef9eeaa66f212435f62ceaf47ac2109743c7e2 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sat, 3 Oct 2026 07:17:49 +0200 Subject: [PATCH 3/6] fix(server): retain renamed branches after launch cleanup fails --- .../ThreadLaunchService.test.ts | 98 ++++++++++++++++++- .../orchestration-v2/ThreadLaunchService.ts | 61 ++++++++---- 2 files changed, 139 insertions(+), 20 deletions(-) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 8f9338c13e42..6d8a58b8cc2b 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -106,6 +106,9 @@ interface HarnessOptions { readonly generateBranchName?: TextGeneration.TextGeneration["Service"]["generateBranchName"]; readonly serverSettings?: Parameters[0]; readonly providers?: ReadonlyArray; + readonly beforeLaunchDispatch?: ( + command: Parameters[0], + ) => Effect.Effect; } function makeHarness(options: HarnessOptions = {}) { @@ -187,8 +190,24 @@ function makeHarness(options: HarnessOptions = {}) { folderForThread: () => Effect.succeed(Option.none()), }), ); + const beforeLaunchDispatch = options.beforeLaunchDispatch; + const launchThreadManagement = beforeLaunchDispatch + ? Layer.effect( + ThreadManagement.ThreadManagementService, + Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + return ThreadManagement.ThreadManagementService.of({ + ...threads, + dispatch: (command) => + beforeLaunchDispatch(command).pipe(Effect.andThen(threads.dispatch(command))), + }); + }), + ).pipe(Layer.provide(threadManagement)) + : threadManagement; const launch = ThreadLaunch.layer.pipe( - Layer.provide(Layer.mergeAll(externalServices, threadManagement, receipts, IdAllocator.layer)), + Layer.provide( + Layer.mergeAll(externalServices, launchThreadManagement, receipts, IdAllocator.layer), + ), ); const projectedProjects = Layer.mock(ProjectStore.ProjectStoreV2)({ get: (requestedProjectId) => @@ -1586,6 +1605,83 @@ it.effect("stops branch renaming before removing a failed launch worktree", () = }), ); +it.effect.each(["original", "newer-branch", "newer-workspace"] as const)( + "keeps survivor metadata accurate after interrupting a completed rename with %s binding", + (binding) => + Effect.gen(function* () { + const metadataEntered = yield* Deferred.make(); + const setupCompletion = yield* Deferred.make<{ exitCode: number; durationMs: number }>(); + const actualBranch = yield* Ref.make("t3code/abcd1234"); + const harness = makeHarness({ + renameBranch: () => + Ref.set(actualBranch, "renamed-branch").pipe(Effect.as({ branch: "renamed-branch" })), + beforeLaunchDispatch: (command) => + String(command.commandId).endsWith(":branch-rename") + ? Deferred.succeed(metadataEntered, undefined).pipe(Effect.andThen(Effect.never)) + : Effect.void, + removeWorktree: () => Effect.fail(new Error("remove failed") as never), + runSetup: () => + Effect.succeed({ + status: "started" as const, + async: false, + scriptId: "setup", + scriptName: "Setup", + scriptCommand: "vp install", + terminalId: "setup", + cwd: "/repo-worktrees/feature", + completion: Deferred.await(setupCompletion), + }), + }); + yield* Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + const input = launchInput({ + command: `launch:renamed-survivor:${binding}`, + thread: `thread:renamed-survivor:${binding}`, + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "t3code/abcd1234" }, + }); + const launched = yield* launches.launch(input); + yield* Deferred.await(metadataEntered); + assert.equal(yield* Ref.get(actualBranch), "renamed-branch"); + assert.equal( + (yield* threads.getThreadProjection(launched.threadId)).thread.branch, + "t3code/abcd1234", + ); + const newerPath = + binding === "newer-workspace" ? "/replacement-worktree" : "/repo-worktrees/feature"; + if (binding !== "original") { + yield* threads.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:newer-binding`), + threadId: launched.threadId, + worktreePath: newerPath, + branch: "newer-branch", + }); + } + yield* Deferred.succeed(setupCompletion, { exitCode: 1, durationMs: 1 }); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.equal( + projection.thread.branch, + binding === "original" ? "renamed-branch" : "newer-branch", + ); + assert.equal(projection.thread.worktreePath, newerPath); + if (binding === "original") { + assert.equal((yield* tracker.get(launched.threadId))?.branch, "renamed-branch"); + } + assert.equal(projection.runs[0]?.status, "failed"); + }).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* () { diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index bc3a7fdde2de..da2d774877b8 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -220,6 +220,7 @@ const make = Effect.gen(function* () { const tracked = input.workspaceStrategy.type === "worktree"; let createdWorktreePath: string | null = null; let createdWorktreeBranch: string | null = null; + let renamedWorktreeBranch: string | null = null; let setupTerminalId: string | null = null; let branchRenameFiber: Fiber.Fiber | null = null; let preparationComplete = false; @@ -387,12 +388,22 @@ const make = Effect.gen(function* () { const worktreeCwd = worktreePath; branchRenameFiber = yield* generateBranchNameFor(worktreeCwd, initialMessage).pipe( Effect.flatMap(({ branch: newBranch, exactName }) => - git.renameBranch({ - cwd: worktreeCwd, - oldBranch, - newBranch, - ...(exactName ? { exactName: true } : {}), - }), + git + .renameBranch({ + cwd: worktreeCwd, + oldBranch, + newBranch, + ...(exactName ? { exactName: true } : {}), + }) + .pipe( + // Record a completed rename before cancellation can stop its + // metadata write. The external Git operation stays interruptible. + Effect.onExit((exit) => + Effect.sync(() => { + if (Exit.isSuccess(exit)) renamedWorktreeBranch = exit.value.branch; + }), + ), + ), ), Effect.flatMap((renamed) => threads.dispatch({ @@ -593,22 +604,34 @@ const make = Effect.gen(function* () { ).pipe(Effect.as(null)), ), ); - // Checkout can fail after claiming a path but before the normal - // workspace update. Bind that survivor so it can still be found. - if (shell !== null && shell.worktreePath !== createdWorktreePath) { - yield* threads - .dispatch({ - type: "thread.metadata.update", - commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), - threadId, - worktreePath: createdWorktreePath, - branch: createdWorktreeBranch, - }) - .pipe(Effect.ignore); + const survivingBranch = renamedWorktreeBranch ?? createdWorktreeBranch; + // Repair only the launch's binding, including a checkout that + // failed before publishing its path. Keep newer bindings intact. + const ownsBinding = + shell !== null && + ((shell.worktreePath === createdWorktreePath && + (shell.branch === createdWorktreeBranch || shell.branch === survivingBranch)) || + (shell.worktreePath === null && + shell.branch === (input.workspaceStrategy.branch ?? null))); + if (ownsBinding) { + if ( + shell.worktreePath !== createdWorktreePath || + shell.branch !== survivingBranch + ) { + yield* threads + .dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), + threadId, + worktreePath: createdWorktreePath, + branch: survivingBranch, + }) + .pipe(Effect.ignore); + } yield* setupTracker.update(threadId, (snapshot) => ({ ...snapshot, worktreePath: createdWorktreePath, - branch: createdWorktreeBranch, + branch: survivingBranch, })); } } From 84e6660e3b4516562663934ea3a935490b593502 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sat, 3 Oct 2026 13:03:23 +0200 Subject: [PATCH 4/6] fix(server): preserve newer workspace bindings during launch cleanup --- .../src/orchestration-v2/Orchestrator.ts | 7 +- .../ThreadLaunchService.test.ts | 108 +++++++++++++++++- .../orchestration-v2/ThreadLaunchService.ts | 50 ++++---- .../src/orchestration-v2/runtimeLayer.test.ts | 12 ++ packages/contracts/src/orchestrationV2.ts | 1 + 5 files changed, 153 insertions(+), 25 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 2707fe5b8344..6d749418db5a 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -2340,13 +2340,14 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio } if ( command.type === "thread.metadata.update" && - command.expectedWorktreePath !== undefined && - command.expectedWorktreePath !== thread.worktreePath + ((command.expectedWorktreePath !== undefined && + command.expectedWorktreePath !== thread.worktreePath) || + (command.expectedBranch !== undefined && command.expectedBranch !== thread.branch)) ) { return yield* new OrchestratorDispatchError({ commandId: command.commandId, commandType: command.type, - cause: `Thread ${command.threadId} worktree changed before the metadata update could be applied.`, + cause: `Thread ${command.threadId} workspace binding changed before the metadata update could be applied.`, }); } if (command.type === "thread.metadata.update" && command.expectedEmpty === true) { diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 6d8a58b8cc2b..dcbd4422078e 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -1520,6 +1520,110 @@ it.effect.each([true, false])( }, ); +it.effect("keeps a newer workspace binding after removing a failed launch worktree", () => + Effect.gen(function* () { + const setupEntered = yield* Deferred.make(); + const setupFailed = yield* Deferred.make(); + const harness = makeHarness({ + runSetup: () => + Deferred.succeed(setupEntered, undefined).pipe( + Effect.andThen(Deferred.await(setupFailed)), + Effect.andThen(Effect.fail(new Error("setup failed") as never)), + ), + }); + yield* Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const input = launchInput({ + command: "launch:removed-newer-workspace", + thread: "thread:removed-newer-workspace", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }); + const launched = yield* launches.launch(input); + yield* Deferred.await(setupEntered); + yield* threads.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:newer-binding`), + threadId: launched.threadId, + worktreePath: "/replacement-worktree", + branch: "newer-branch", + }); + yield* Deferred.succeed(setupFailed, undefined); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.equal(projection.thread.worktreePath, "/replacement-worktree"); + assert.equal(projection.thread.branch, "newer-branch"); + assert.deepEqual( + harness.removeWorktree.mock.calls.map(([input]) => input.path), + ["/repo-worktrees/feature"], + ); + }).pipe(Effect.provide(harness.layer)); + }), +); + +it.effect.each([ + { removed: true, newPath: "/replacement-worktree" }, + { removed: false, newPath: "/replacement-worktree" }, + { removed: true, newPath: null }, + { removed: false, newPath: null }, +])("keeps a workspace rebound during partial checkout cleanup: %o", ({ removed, newPath }) => + Effect.gen(function* () { + const cleanupEntered = yield* Deferred.make(); + const cleanupResume = yield* Deferred.make(); + const harness = makeHarness({ + createWorktree: (_input, options) => + (options?.progress?.onWorktreeClaimed?.("/claimed") ?? Effect.void).pipe( + Effect.andThen(Effect.fail(new Error("checkout failed") as never)), + ), + removeWorktree: () => + removed ? Effect.void : Effect.fail(new Error("remove failed") as never), + beforeLaunchDispatch: (command) => + String(command.commandId).endsWith(":cleanup-workspace") + ? Deferred.succeed(cleanupEntered, undefined).pipe( + Effect.andThen(Deferred.await(cleanupResume)), + ) + : Effect.void, + }); + yield* Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const input = launchInput({ + command: `launch:cleanup-rebound:${removed}`, + thread: `thread:cleanup-rebound:${removed}`, + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }); + const launched = yield* launches.launch(input); + yield* Deferred.await(cleanupEntered); + yield* threads.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:newer-binding`), + threadId: launched.threadId, + worktreePath: newPath, + branch: "newer-branch", + }); + yield* Deferred.succeed(cleanupResume, undefined); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + assert.equal(projection.thread.worktreePath, newPath); + assert.equal(projection.thread.branch, "newer-branch"); + }).pipe(Effect.provide(harness.layer)); + }), +); + it.effect("keeps a failed launch worktree bound when cleanup cannot remove it", () => { const harness = makeHarness({ runSetup: () => Effect.fail(new Error("setup failed") as never), @@ -2110,8 +2214,10 @@ it.effect("shared intake preserves durable attachment bytes after a lost launch ); assert.isDefined(stored); assert.notEqual(stored.attachments[0]?.id, attachment.id); + const contextRecord = stored.context?.records[0]; + assert.isDefined(contextRecord); assert.equal( - (stored.context?.records[0] as { attachmentId: string }).attachmentId, + (contextRecord as { attachmentId: string }).attachmentId, stored.attachments[0]?.id, ); const userItem = accepted.turnItems.find( diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index da2d774877b8..219d902d0ff9 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -574,36 +574,42 @@ const make = Effect.gen(function* () { }).pipe(Effect.as(false)), ), ); + const shell = yield* threads.getThreadShell(threadId).pipe( + Effect.catchCause((cleanupCause) => + Effect.logWarning("Failed to read thread launch worktree binding during cleanup", { + threadId, + worktreePath: createdWorktreePath, + cause: cleanupCause, + }).pipe(Effect.as(null)), + ), + ); // Keep the binding when removal fails so the surviving worktree // remains discoverable instead of becoming an orphan. if (removed) { - yield* threads - .dispatch({ - type: "thread.metadata.update", - commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), - threadId, - worktreePath: null, - branch: null, - }) - .pipe(Effect.ignore); + if ( + shell !== null && + (shell.worktreePath === createdWorktreePath || + (shell.worktreePath === null && + shell.branch === (input.workspaceStrategy.branch ?? null))) + ) { + yield* threads + .dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), + threadId, + expectedWorktreePath: shell.worktreePath, + expectedBranch: shell.branch, + worktreePath: null, + branch: null, + }) + .pipe(Effect.ignore); + } yield* setupTracker.update(threadId, (snapshot) => ({ ...snapshot, worktreePath: null, branch: null, })); } else { - const shell = yield* threads.getThreadShell(threadId).pipe( - Effect.catchCause((cleanupCause) => - Effect.logWarning( - "Failed to read thread launch worktree binding during cleanup", - { - threadId, - worktreePath: createdWorktreePath, - cause: cleanupCause, - }, - ).pipe(Effect.as(null)), - ), - ); const survivingBranch = renamedWorktreeBranch ?? createdWorktreeBranch; // Repair only the launch's binding, including a checkout that // failed before publishing its path. Keep newer bindings intact. @@ -623,6 +629,8 @@ const make = Effect.gen(function* () { type: "thread.metadata.update", commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), threadId, + expectedWorktreePath: shell.worktreePath, + expectedBranch: shell.branch, worktreePath: createdWorktreePath, branch: survivingBranch, }) diff --git a/apps/server/src/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index 2c5235c767d7..dc8cb6390e62 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -1675,6 +1675,18 @@ it.layer(TestLayer)("OrchestrationV2LayerLive lifecycle", (it) => { }) .pipe(Effect.flip); assert.instanceOf(staleWorkspaceUpdate, Orchestrator.OrchestratorDispatchError); + const staleBranchUpdate = yield* orchestrator + .dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make("runtime-layer-lifecycle-stale-branch"), + threadId, + branch: null, + worktreePath: null, + expectedWorktreePath: "/tmp/t3-v2-worktree", + expectedBranch: null, + }) + .pipe(Effect.flip); + assert.instanceOf(staleBranchUpdate, Orchestrator.OrchestratorDispatchError); const projectionAfterStaleWorkspaceUpdate = yield* orchestrator.getThreadProjection(threadId); assert.equal(projectionAfterStaleWorkspaceUpdate.thread.branch, "feature/v2"); assert.equal(projectionAfterStaleWorkspaceUpdate.thread.worktreePath, "/tmp/t3-v2-worktree"); diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 6004ab119917..8aedf894d67a 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -2573,6 +2573,7 @@ export const OrchestrationV2Command = Schema.Union([ branch: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), worktreePath: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), expectedWorktreePath: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), + expectedBranch: Schema.optional(Schema.NullOr(TrimmedNonEmptyString)), /** Reject unless no message or run has landed on this thread. */ expectedEmpty: Schema.optional(Schema.Boolean), limitRecovery: Schema.optional(Schema.NullOr(OrchestrationV2LimitRecoveryUpdate)), From 844ff8c07be08fc24e0049905fe9194a6054dd26 Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sun, 4 Oct 2026 13:39:41 +0200 Subject: [PATCH 5/6] test(server): document idle launch retention during shutdown --- .../ThreadLaunchService.test.ts | 95 ++++++++++--------- .../orchestration-v2/ThreadLaunchService.ts | 4 +- docs/internals/overview.md | 2 +- docs/user/thread-sidebar.md | 8 +- 4 files changed, 58 insertions(+), 51 deletions(-) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index dcbd4422078e..b5a719bcb76f 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -2410,51 +2410,56 @@ it.effect("cancels tracked setup before provider work is released", () => }), ); -it.effect("keeps a released worktree when shutdown interrupts its async setup", () => { - const harness = makeHarness({ - runSetup: () => - Effect.succeed({ - status: "started" as const, - async: true, - scriptId: "setup", - scriptName: "Setup", - scriptCommand: "vp install", - terminalId: "setup", - cwd: "/repo-worktrees/feature", - completion: Effect.never, - }), - }); - return Effect.gen(function* () { - yield* Effect.scoped( - Effect.gen(function* () { - const launches = yield* ThreadLaunch.ThreadLaunchService; - const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; - const threads = yield* ThreadManagement.ThreadManagementService; - const launched = yield* launches.launch( - launchInput({ - command: "launch:released-shutdown", - thread: "thread:released-shutdown", - message: "Start", - workspace: { type: "worktree", baseRef: "main", branch: "feature" }, - }), - ); - yield* tracker.stream(launched.threadId).pipe( - Stream.filter( - (snapshot) => - snapshot?.stages.some((stage) => stage.id === "agent" && stage.status === "done") === - true, - ), - Stream.runHead, - ); - const projection = yield* threads.getThreadProjection(launched.threadId); - assert.equal(projection.runs[0]?.status, "starting"); - assert.equal(projection.thread.worktreePath, "/repo-worktrees/feature"); - }).pipe(Effect.provide(harness.layer)), - ); - assert.isEmpty(harness.removeWorktree.mock.calls); - assert.isEmpty(harness.closeTerminal.mock.calls); - }); -}); +it.effect.each([undefined, "Start"])( + "keeps a prepared worktree when shutdown interrupts its async setup with message %s", + (message) => { + const harness = makeHarness({ + runSetup: () => + Effect.succeed({ + status: "started" as const, + async: true, + scriptId: "setup", + scriptName: "Setup", + scriptCommand: "vp install", + terminalId: "setup", + cwd: "/repo-worktrees/feature", + completion: Effect.never, + }), + }); + return Effect.gen(function* () { + yield* Effect.scoped( + Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + const threads = yield* ThreadManagement.ThreadManagementService; + const launched = yield* launches.launch( + launchInput({ + command: "launch:released-shutdown", + thread: "thread:released-shutdown", + ...(message === undefined ? {} : { message }), + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }), + ); + yield* tracker.stream(launched.threadId).pipe( + Stream.filter( + (snapshot) => + snapshot?.stages.some( + (stage) => stage.id === "agent" && stage.status === "done", + ) === true, + ), + Stream.runHead, + ); + const projection = yield* threads.getThreadProjection(launched.threadId); + if (message === undefined) assert.isEmpty(projection.runs); + else assert.equal(projection.runs[0]?.status, "starting"); + assert.equal(projection.thread.worktreePath, "/repo-worktrees/feature"); + }).pipe(Effect.provide(harness.layer)), + ); + assert.isEmpty(harness.removeWorktree.mock.calls); + assert.isEmpty(harness.closeTerminal.mock.calls); + }); + }, +); it.effect.each([0, 1])("releases an async setup before its completion with exit %s", (exitCode) => Effect.gen(function* () { diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index 219d902d0ff9..9248a6eb96de 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -534,8 +534,8 @@ const make = Effect.gen(function* () { Effect.onError((cause) => Effect.gen(function* () { const cancelled = Cause.hasInterruptsOnly(cause); - // Once the prepared run is released, the provider owns its workspace. - // An async setup interrupted during shutdown must not delete it. + // A released run owns its workspace. An idle launch also keeps its + // prepared checkout when shutdown interrupts the remaining async setup. if (!preparationComplete) { if (branchRenameFiber !== null) yield* Fiber.interrupt(branchRenameFiber); if (setupTerminalId !== null) diff --git a/docs/internals/overview.md b/docs/internals/overview.md index 149bcf57dff0..6d2733b96861 100644 --- a/docs/internals/overview.md +++ b/docs/internals/overview.md @@ -91,7 +91,7 @@ and trigger a check. A merge outside T3, such as an agent running `gh pr merge`, notification, so the [PR sync reactor](../../apps/server/src/orchestration-v2/PullRequestSyncReactor.ts) re-reads a thread's open links when a run that ran a merge or close command ends. The guarded `thread.auto-settle` command rejects newer activity, explicit settlement overrides, and live or -blocked work. It records the activity timestamp for stable +blocked work. It records the activity timestamp, or creation time for a thread with no activity, for stable sorting and detaches idle provider sessions. Clients render the persisted result; they do not derive settlement from their own clocks or PR caches. diff --git a/docs/user/thread-sidebar.md b/docs/user/thread-sidebar.md index 9f2eb47aa392..7f0199fd869a 100644 --- a/docs/user/thread-sidebar.md +++ b/docs/user/thread-sidebar.md @@ -13,12 +13,14 @@ an existing worktree, use **New thread in this worktree** from the branch toolba When you change a new thread's project, T3 Code stays in the current environment if that project exists there. Otherwise it selects an environment that has it. -If preparing a new worktree fails before the agent starts, T3 Code closes its setup -terminal and removes the new worktree. The thread and its messages stay visible. +If preparing a new worktree fails before a queued agent can start, T3 Code closes +its setup terminal and removes the new worktree. The thread and its messages stay visible. Messages queued during preparation then run from the project's root directory. To try again in a separate worktree, start a new thread with **New worktree**. An existing worktree you selected stays in place. If removing a failed worktree fails, it stays linked to the thread so you can inspect or remove it. +An idle thread created without a message keeps its prepared worktree if shutdown +interrupts an asynchronous setup script. ### Start without a project @@ -138,7 +140,7 @@ runs a command, such as a dev server, stays open. By default, environments settle inactive threads after three days and settle threads whose pull request merged. A closed pull request can also settle an idle -thread. A thread with no messages or runs becomes inactive from its creation time. +thread. A thread with no messages or runs uses its creation time as the settlement timestamp. Work in progress, pending questions or approvals, and live background work prevent automatic settlement. An open pull request does not prevent inactivity settlement, but an old closed or merged pull request does not settle work you From 43e0632517ef76b38890a39944e07b32e77041cc Mon Sep 17 00:00:00 2001 From: Adamulek123 Date: Sun, 4 Oct 2026 23:40:40 +0200 Subject: [PATCH 6/6] fix(server): distinguish completed checkouts on preparation retry --- .../src/orchestration-v2/Orchestrator.ts | 65 +++- .../ThreadLaunchService.test.ts | 356 +++++++++++++++++- .../orchestration-v2/ThreadLaunchService.ts | 139 ++++--- docs/user/thread-sidebar.md | 9 +- packages/contracts/src/orchestrationV2.ts | 10 + 5 files changed, 520 insertions(+), 59 deletions(-) diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index 8c131ac9e641..0271ac875f78 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -5167,7 +5167,12 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio ? {} : { restartContinuationOfRunId: command.restartContinuationOfRunId }), ...(dispatchMode.type === "defer_start" && dispatchMode.workspaceStrategy !== undefined - ? { workspacePreparation: dispatchMode.workspaceStrategy } + ? { + workspacePreparation: dispatchMode.workspaceStrategy, + ...(dispatchMode.workspaceStrategy.type === "worktree" + ? { completedWorktreePath: null } + : {}), + } : {}), ...wakeWorkStartedAt(projection.runs, command), }; @@ -7565,6 +7570,52 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio }); } const now = yield* DateTime.now; + const workspace = command.completedWorkspace; + if ( + workspace !== undefined && + (command.phase !== "setup" || + state.run.workspacePreparation?.type !== "worktree" || + projection.thread.worktreePath !== workspace.expectedWorktreePath || + projection.thread.branch !== workspace.expectedBranch) + ) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: "The prepared run no longer owns the workspace binding.", + }); + } + if (workspace !== undefined) { + yield* emit( + events, + command, + )({ + type: "thread.metadata-updated", + threadId: command.threadId, + occurredAt: now, + payload: { + ...projection.thread, + worktreePath: workspace.worktreePath, + branch: workspace.branch, + updatedAt: now, + }, + }); + } + if (workspace !== undefined || command.phase === "worktree") { + yield* emit( + events, + command, + )({ + type: "run.updated", + threadId: command.threadId, + runId: state.run.id, + providerInstanceId: state.run.providerInstanceId, + occurredAt: now, + payload: { + ...state.run, + completedWorktreePath: workspace?.worktreePath ?? null, + }, + }); + } yield* emit( events, command, @@ -7602,6 +7653,18 @@ const makeOrchestrator = Effect.fn("orchestrationV2.Orchestrator.layer")(functio cause: `Run ${command.runId} is not awaiting workspace preparation.`, }); } + if ( + state.run.workspacePreparation?.type === "worktree" && + state.run.completedWorktreePath !== undefined && + (state.run.completedWorktreePath === null || + state.run.completedWorktreePath !== projection.thread.worktreePath) + ) { + return yield* new OrchestratorDispatchError({ + commandId: command.commandId, + commandType: command.type, + cause: "The prepared run has no completed checkout at its current workspace binding.", + }); + } const now = yield* DateTime.now; const resolvedRuntimePolicy = yield* runtimePolicy .resolve({ thread: projection.thread, modelSelection: state.run.modelSelection }) diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts index 21e6577d55fe..6f89cebcfc9c 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.test.ts @@ -38,6 +38,7 @@ 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 SqlClient from "effect/unstable/sql/SqlClient"; import * as GitWorkflow from "../git/GitWorkflowService.ts"; import { SqlitePersistenceMemory } from "../persistence/Layers/Sqlite.ts"; @@ -96,6 +97,7 @@ const adapter = { } as ProviderAdapterV2Shape; interface HarnessOptions { + readonly pathExists?: FileSystem.FileSystem["exists"]; readonly managedFolders?: Layer.Layer; readonly createWorktree?: GitWorkflow.GitWorkflowService["Service"]["createWorktree"]; readonly removeWorktree?: GitWorkflow.GitWorkflowService["Service"]["removeWorktree"]; @@ -110,6 +112,9 @@ interface HarnessOptions { readonly beforeLaunchDispatch?: ( command: Parameters[0], ) => Effect.Effect; + readonly afterLaunchDispatch?: ( + command: Parameters[0], + ) => Effect.Effect; } function makeHarness(options: HarnessOptions = {}) { @@ -145,6 +150,7 @@ function makeHarness(options: HarnessOptions = {}) { options.generateTitle ?? (() => Effect.succeed({ title: "Generated title" })), ); const externalServices = Layer.mergeAll( + FileSystem.layerNoop({ exists: options.pathExists ?? (() => Effect.succeed(false)) }), WorktreeSetupTracker.layer, Layer.mock(ProjectCloneTracker.ProjectCloneTracker)({ get: () => Effect.succeed(null) }), Layer.mock(TerminalManager.TerminalManager)({ close: closeTerminal }), @@ -192,19 +198,24 @@ function makeHarness(options: HarnessOptions = {}) { }), ); const beforeLaunchDispatch = options.beforeLaunchDispatch; - const launchThreadManagement = beforeLaunchDispatch - ? Layer.effect( - ThreadManagement.ThreadManagementService, - Effect.gen(function* () { - const threads = yield* ThreadManagement.ThreadManagementService; - return ThreadManagement.ThreadManagementService.of({ - ...threads, - dispatch: (command) => - beforeLaunchDispatch(command).pipe(Effect.andThen(threads.dispatch(command))), - }); - }), - ).pipe(Layer.provide(threadManagement)) - : threadManagement; + const afterLaunchDispatch = options.afterLaunchDispatch; + const launchThreadManagement = + beforeLaunchDispatch || afterLaunchDispatch + ? Layer.effect( + ThreadManagement.ThreadManagementService, + Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + return ThreadManagement.ThreadManagementService.of({ + ...threads, + dispatch: (command) => + (beforeLaunchDispatch?.(command) ?? Effect.void).pipe( + Effect.andThen(threads.dispatch(command)), + Effect.ensuring(afterLaunchDispatch?.(command) ?? Effect.void), + ), + }); + }), + ).pipe(Layer.provide(threadManagement)) + : threadManagement; const launch = ThreadLaunch.layer.pipe( Layer.provide( Layer.mergeAll(externalServices, launchThreadManagement, receipts, IdAllocator.layer), @@ -242,6 +253,8 @@ function makeHarness(options: HarnessOptions = {}) { outbox, database, externalServices, + IdAllocator.layer, + receipts, ), createWorktree, removeWorktree, @@ -1462,6 +1475,323 @@ it.effect("a retry reuses a recorded worktree without undoing its branch rename" }).pipe(Effect.provide(harness.layer)); }); +it.effect.each([false, true])( + "preserves incomplete checkout files until explicit removal before Retry (restart: %s)", + (restart) => { + let survivorExists = true; + let provisions = 0; + const harness = makeHarness({ + pathExists: () => Effect.succeed(survivorExists), + createWorktree: (input, options) => + ++provisions === 1 + ? (options?.progress?.onWorktreeClaimed?.("/partial-survivor") ?? Effect.void).pipe( + Effect.andThen( + Effect.fail(new Error("provisioning failed after claiming checkout") as never), + ), + ) + : Effect.succeed({ worktree: { path: "/fresh-checkout", refName: input.newRefName! } }), + removeWorktree: () => Effect.fail(new Error("cleanup failed") as never), + }); + return Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + const withLaunch = (effect: Effect.Effect) => + restart + ? Effect.scoped( + effect.pipe( + Effect.provide(ThreadLaunch.layer.pipe(Layer.provide(WorktreeSetupTracker.layer))), + ), + ) + : effect; + const input = launchInput({ + command: "launch:incomplete", + thread: "thread:incomplete", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }); + const failed = yield* withLaunch( + Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const launched = yield* launches.launch(input); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + return yield* threads.getThreadProjection(launched.threadId); + }), + ); + assert.equal(failed.thread.worktreePath, "/partial-survivor"); + assert.equal(failed.runs[0]?.completedWorktreePath, null); + const runId = failed.runs[0]!.id; + const threadId = failed.thread.id; + yield* withLaunch( + Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + yield* launches.retryPreparation({ + commandId: CommandId.make("launch:incomplete:retry"), + threadId, + runId, + }); + yield* threads.streamStoredEventsFrom({ threadId }).pipe( + Stream.filter( + (stored) => + stored.commandId === CommandId.make("launch:incomplete:retry:fail") && + stored.event.type === "run.updated", + ), + Stream.runHead, + ); + }), + ); + const blocked = yield* threads.getThreadProjection(threadId); + assert.equal(blocked.runs[0]?.status, "failed"); + assert.equal(blocked.thread.worktreePath, "/partial-survivor"); + assert.include( + blocked.turnItems.find((item) => item.type === "error")?.failure.message, + "Its files have been preserved", + ); + assert.equal(harness.createWorktree.mock.calls.length, 1); + assert.equal(harness.removeWorktree.mock.calls.length, 1); + assert.isEmpty(harness.runSetup.mock.calls); + assert.isEmpty(blocked.checkpointScopes); + + // The user explicitly removes the preserved checkout and its branch. + survivorExists = false; + yield* withLaunch( + Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + yield* launches.retryPreparation({ + commandId: CommandId.make("launch:incomplete:recovered"), + threadId, + runId, + }); + yield* threads.streamStoredEventsFrom({ threadId }).pipe( + Stream.filter( + (stored) => + stored.commandId === CommandId.make("launch:incomplete:recovered:release") && + stored.event.type === "run.updated", + ), + Stream.runHead, + ); + }), + ); + const recovered = yield* threads.getThreadProjection(threadId); + assert.equal(recovered.runs[0]?.status, "starting"); + assert.equal(recovered.runs[0]?.completedWorktreePath, "/fresh-checkout"); + assert.equal(recovered.thread.worktreePath, "/fresh-checkout"); + assert.equal(recovered.checkpointScopes[0]?.cwd, "/fresh-checkout"); + assert.equal(harness.createWorktree.mock.calls.length, 2); + assert.equal(harness.runSetup.mock.calls.length, 1); + }).pipe(Effect.provide(harness.layer)); + }, +); + +it.effect.each([false, true])( + "reuses a completed checkout after setup failure and service restart (legacy: %s)", + (legacy) => { + let setupFailures = 1; + const harness = makeHarness({ + runSetup: () => + setupFailures-- > 0 + ? Effect.fail(new Error("setup failed") as never) + : Effect.succeed({ status: "no-script" as const }), + }); + return Effect.gen(function* () { + const threads = yield* ThreadManagement.ThreadManagementService; + const freshLaunch = ThreadLaunch.layer.pipe(Layer.provide(WorktreeSetupTracker.layer)); + const failed = yield* Effect.scoped( + Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const launched = yield* launches.launch( + launchInput({ + command: "launch:completed", + thread: "thread:completed", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }), + ); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + return yield* threads.getThreadProjection(launched.threadId); + }).pipe(Effect.provide(freshLaunch)), + ); + assert.equal(failed.runs[0]?.completedWorktreePath, "/repo-worktrees/feature"); + if (legacy) { + // Emulate a persisted run from released versions, before this optional fact existed. + const sql = yield* SqlClient.SqlClient; + yield* sql`UPDATE orchestration_v2_projection_runs SET payload_json = json_remove(payload_json, '$.completedWorktreePath') WHERE run_id = ${failed.runs[0]!.id}`; + const decoded = yield* threads.getThreadProjection(failed.thread.id); + assert.equal(decoded.runs[0]?.completedWorktreePath, undefined); + } + yield* Effect.scoped( + Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + yield* launches.retryPreparation({ + commandId: CommandId.make("launch:completed:retry"), + threadId: failed.thread.id, + runId: failed.runs[0]!.id, + }); + yield* threads.streamStoredEventsFrom({ threadId: failed.thread.id }).pipe( + Stream.filter( + (stored) => + stored.commandId === CommandId.make("launch:completed:retry:release") && + stored.event.type === "run.updated", + ), + Stream.runHead, + ); + }).pipe(Effect.provide(freshLaunch)), + ); + const retried = yield* threads.getThreadProjection(failed.thread.id); + assert.equal(retried.runs[0]?.status, "starting"); + assert.equal(retried.thread.worktreePath, "/repo-worktrees/feature"); + assert.equal(harness.createWorktree.mock.calls.length, 1); + assert.equal(harness.runSetup.mock.calls.length, 2); + assert.isEmpty(harness.removeWorktree.mock.calls); + }).pipe(Effect.provide(harness.layer)); + }, +); + +it.effect("a completed checkout does not authorize Retry at a rebound path", () => { + const harness = makeHarness({ + pathExists: () => Effect.succeed(true), + runSetup: () => Effect.fail(new Error("setup failed") as never), + }); + return Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const launched = yield* launches.launch( + launchInput({ + command: "launch:rebound", + thread: "thread:rebound", + message: "Start", + workspace: { type: "worktree", baseRef: "main", branch: "feature" }, + }), + ); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => stored.event.type === "run.updated" && stored.event.payload.status === "failed", + ), + Stream.runHead, + ); + const failed = yield* threads.getThreadProjection(launched.threadId); + yield* threads.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make("launch:rebound:change"), + threadId: launched.threadId, + worktreePath: "/new-binding", + branch: "other", + }); + yield* launches.retryPreparation({ + commandId: CommandId.make("launch:rebound:retry"), + threadId: launched.threadId, + runId: failed.runs[0]!.id, + }); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.commandId === CommandId.make("launch:rebound:retry:fail") && + stored.event.type === "run.updated", + ), + Stream.runHead, + ); + const blocked = yield* threads.getThreadProjection(launched.threadId); + assert.equal(blocked.runs[0]?.status, "failed"); + assert.equal(blocked.runs[0]?.completedWorktreePath, "/repo-worktrees/feature"); + assert.equal(blocked.thread.worktreePath, "/new-binding"); + assert.equal(harness.runSetup.mock.calls.length, 1); + assert.equal(harness.createWorktree.mock.calls.length, 1); + assert.isEmpty(harness.removeWorktree.mock.calls); + assert.isEmpty(blocked.checkpointScopes); + }).pipe(Effect.provide(harness.layer)); +}); + +it.effect("does not release agent work when the workspace changes during setup", () => { + return Effect.gen(function* () { + const setupStarted = yield* Deferred.make(); + const finishSetup = yield* Deferred.make(); + const branchNameStarted = yield* Deferred.make(); + const finishBranchName = yield* Deferred.make(); + const branchRenameWritten = yield* Deferred.make(); + const harness = makeHarness({ + generateBranchName: () => + Deferred.succeed(branchNameStarted, undefined).pipe( + Effect.andThen(Deferred.await(finishBranchName)), + Effect.as({ branch: "generated-branch" }), + ), + afterLaunchDispatch: (command) => + command.commandId === CommandId.make("launch:setup-rebound:branch-rename") + ? Deferred.succeed(branchRenameWritten, undefined).pipe(Effect.asVoid) + : Effect.void, + runSetup: () => + Deferred.succeed(setupStarted, undefined).pipe( + Effect.andThen(Deferred.await(finishSetup)), + Effect.as({ status: "no-script" as const }), + ), + }); + yield* Effect.gen(function* () { + const launches = yield* ThreadLaunch.ThreadLaunchService; + const threads = yield* ThreadManagement.ThreadManagementService; + const launched = yield* launches.launch( + launchInput({ + command: "launch:setup-rebound", + thread: "thread:setup-rebound", + message: "Start", + workspace: { type: "worktree", baseRef: "main" }, + }), + ); + yield* Deferred.await(setupStarted); + yield* Deferred.await(branchNameStarted); + const ready = yield* threads.getThreadProjection(launched.threadId); + assert.equal(ready.thread.worktreePath, ready.runs[0]?.completedWorktreePath); + const workspaceEvents = yield* threads + .streamStoredEventsFrom({ threadId: launched.threadId }) + .pipe( + Stream.filter( + (stored) => stored.commandId === CommandId.make("launch:setup-rebound:workspace"), + ), + Stream.take(3), + Stream.runCollect, + ); + assert.deepEqual( + workspaceEvents.map((stored) => stored.event.type), + ["thread.metadata-updated", "run.updated", "turn-item.updated"], + ); + yield* threads.dispatch({ + type: "thread.metadata.update", + commandId: CommandId.make("launch:setup-rebound:change"), + threadId: launched.threadId, + worktreePath: "/new-binding", + branch: "other", + }); + yield* Deferred.succeed(finishBranchName, undefined); + yield* Deferred.await(branchRenameWritten); + const afterRename = yield* threads.getThreadProjection(launched.threadId); + assert.equal(afterRename.thread.worktreePath, "/new-binding"); + assert.equal(afterRename.thread.branch, "other"); + yield* Deferred.succeed(finishSetup, undefined); + yield* threads.streamStoredEventsFrom({ threadId: launched.threadId }).pipe( + Stream.filter( + (stored) => + stored.commandId === CommandId.make("launch:setup-rebound:fail") && + stored.event.type === "run.updated", + ), + Stream.runHead, + ); + const failed = yield* threads.getThreadProjection(launched.threadId); + assert.equal(failed.runs[0]?.status, "failed"); + assert.equal(failed.thread.worktreePath, "/new-binding"); + assert.isEmpty(failed.checkpointScopes); + assert.isEmpty(harness.removeWorktree.mock.calls); + }).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. diff --git a/apps/server/src/orchestration-v2/ThreadLaunchService.ts b/apps/server/src/orchestration-v2/ThreadLaunchService.ts index df75125936ab..af15f4ebfb2f 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -24,6 +24,7 @@ import * as Context from "effect/Context"; import * as Effect from "effect/Effect"; import * as Exit from "effect/Exit"; import * as Fiber from "effect/Fiber"; +import * as FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Ref from "effect/Ref"; @@ -173,6 +174,7 @@ const make = Effect.gen(function* () { const cloneTracker = yield* ProjectCloneTracker.ProjectCloneTracker; const terminals = yield* TerminalManager.TerminalManager; const git = yield* GitWorkflow.GitWorkflowService; + const fileSystem = yield* FileSystem.FileSystem; const setupScripts = yield* ProjectSetupScriptRunner.ProjectSetupScriptRunner; const providerRegistry = yield* ProviderRegistry.ProviderRegistry; const serverSettings = yield* ServerSettings.ServerSettingsService; @@ -307,6 +309,12 @@ const make = Effect.gen(function* () { // under a temporary `t3code/` name so the worktree never waits on // name generation, then rename in the background below. let branch = input.workspaceStrategy.branch ?? null; + const binding = yield* threads + .getThreadShell(threadId) + .pipe(Effect.mapError(mapError(input, "update-thread", threadId))); + if (binding === null) { + return yield* mapError(input, "update-thread", threadId)("Thread no longer exists."); + } let worktreePath = input.workspaceStrategy.type === "existing_worktree" ? input.workspaceStrategy.worktreePath @@ -402,13 +410,34 @@ const make = Effect.gen(function* () { // the first attempt's branch rename. if (reused === undefined) { yield* threads - .dispatch({ - type: "thread.metadata.update", - commandId: CommandId.make(`${input.commandId}:workspace`), - threadId, - branch, - worktreePath, - }) + .dispatch( + input.workspaceStrategy.type === "worktree" && + runId !== null && + worktreePath !== null && + branch !== null + ? { + type: "prepared-run.progress", + commandId: CommandId.make(`${input.commandId}:workspace`), + threadId, + runId, + phase: "setup", + completedWorkspace: { + worktreePath, + branch, + expectedWorktreePath: binding.worktreePath, + expectedBranch: binding.branch, + }, + } + : { + type: "thread.metadata.update", + commandId: CommandId.make(`${input.commandId}:workspace`), + threadId, + branch, + worktreePath, + expectedWorktreePath: binding.worktreePath, + expectedBranch: binding.branch, + }, + ) .pipe(Effect.mapError(mapError(input, "update-thread", threadId))); } workspaceRecorded = true; @@ -450,6 +479,8 @@ const make = Effect.gen(function* () { type: "thread.metadata.update", commandId: CommandId.make(`${input.commandId}:branch-rename`), threadId, + expectedWorktreePath: worktreeCwd, + expectedBranch: oldBranch, branch: renamed.branch, worktreePath: worktreeCwd, }), @@ -1025,42 +1056,66 @@ const make = Effect.gen(function* () { projection: OrchestrationV2ThreadProjection, run: OrchestrationV2ThreadProjection["runs"][number], workspacePreparation: ThreadLaunchWorkspaceStrategy, - ) => { - const message = projection.messages.find((candidate) => candidate.id === run.userMessageId); - // A worktree the failed attempt already created is reused, not created again. - const reuse = - workspacePreparation.type === "worktree" && - projection.thread.worktreePath !== null && - projection.thread.branch !== null - ? { - strategy: { - type: "existing_worktree" as const, - worktreePath: projection.thread.worktreePath, - branch: projection.thread.branch, - }, - reusedWorktree: { baseRef: workspacePreparation.baseRef }, - } - : null; - return schedulePreparation( - { - commandId: input.commandId, - projectId: projection.thread.projectId, - workspaceStrategy: reuse?.strategy ?? workspacePreparation, - ...(reuse === null ? {} : { reusedWorktree: reuse.reusedWorktree }), - ...(message === undefined - ? {} - : { - initialMessage: { - text: message.text, - attachments: message.attachments, - ...(message.context ? { context: message.context } : {}), + ) => + Effect.gen(function* () { + const message = projection.messages.find((candidate) => candidate.id === run.userMessageId); + // Released versions only published a worktree after successful provisioning. + // Preserve that legacy invariant; new runs explicitly start with null proof. + const legacyCompleted = run.completedWorktreePath === undefined; + if ( + workspacePreparation.type === "worktree" && + projection.thread.worktreePath !== null && + !legacyCompleted && + run.completedWorktreePath !== projection.thread.worktreePath && + (yield* fileSystem.exists(projection.thread.worktreePath)) + ) { + return yield* mapError( + { + commandId: input.commandId, + projectId: projection.thread.projectId, + workspaceStrategy: workspacePreparation, + }, + "provision-worktree", + input.threadId, + )( + `Checkout completion is unknown for ${projection.thread.worktreePath}. Its files have been preserved. Back up any changes, remove this worktree and its branch with Git, then select Retry to create a fresh checkout.`, + ); + } + // Only the exact checkout whose successful provisioning was committed is reused. + const reuse = + workspacePreparation.type === "worktree" && + projection.thread.worktreePath !== null && + (legacyCompleted || run.completedWorktreePath === projection.thread.worktreePath) && + projection.thread.branch !== null + ? { + strategy: { + type: "existing_worktree" as const, + worktreePath: projection.thread.worktreePath, + branch: projection.thread.branch, }, - }), - }, - input.threadId, - run.id, - ); - }; + reusedWorktree: { baseRef: workspacePreparation.baseRef }, + } + : null; + return yield* schedulePreparation( + { + commandId: input.commandId, + projectId: projection.thread.projectId, + workspaceStrategy: reuse?.strategy ?? workspacePreparation, + ...(reuse === null ? {} : { reusedWorktree: reuse.reusedWorktree }), + ...(message === undefined + ? {} + : { + initialMessage: { + text: message.text, + attachments: message.attachments, + ...(message.context ? { context: message.context } : {}), + }, + }), + }, + input.threadId, + run.id, + ); + }); return ThreadLaunchService.of({ launch, retryPreparation }); }); diff --git a/docs/user/thread-sidebar.md b/docs/user/thread-sidebar.md index 874be30b861d..88081f3f20ba 100644 --- a/docs/user/thread-sidebar.md +++ b/docs/user/thread-sidebar.md @@ -14,11 +14,14 @@ When you change a new thread's project, T3 Code stays in the current environment if that project exists there. Otherwise it selects an environment that has it. If preparing a new worktree fails before a queued agent can start, the thread, -messages, recorded worktree, and setup output stay visible. Use **Retry** on the preparation failure to try again on the same thread and -reuse its worktree. +messages, recorded worktree, and setup output stay visible. Use **Retry** on the +preparation failure to try again on the same thread. Retry reuses a worktree that +finished provisioning, including one whose setup script failed. Cancelling preparation removes a newly created worktree. A checkout that fails before it is recorded is also removed. If removal fails, the surviving worktree -stays linked to the thread so you can inspect or remove it. +stays linked to the thread so you can inspect or remove it. Retry preserves an +incomplete checkout and asks you to back up any changes and remove the worktree +and its branch with Git before creating a fresh checkout. An idle thread created without a message keeps its prepared worktree if shutdown interrupts an asynchronous setup script. diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 73497f8f47b0..591186d4610a 100644 --- a/packages/contracts/src/orchestrationV2.ts +++ b/packages/contracts/src/orchestrationV2.ts @@ -574,6 +574,8 @@ export const OrchestrationV2Run = Schema.Struct({ delegatedCompletion: Schema.optional(OrchestrationV2DelegatedCompletionCohort), /** How a launch prepares this run's workspace; prepared-run.retry repeats it. */ workspacePreparation: Schema.optional(OrchestrationV2ThreadLaunchWorkspaceStrategy), + /** Successful provisioning of this exact checkout, before its setup script runs. */ + completedWorktreePath: Schema.optional(Schema.NullOr(Schema.String)), }); export type OrchestrationV2Run = typeof OrchestrationV2Run.Type; @@ -2760,6 +2762,14 @@ export const OrchestrationV2Command = Schema.Union([ threadId: ThreadId, runId: RunId, phase: Schema.Literals(["worktree", "setup"]), + completedWorkspace: Schema.optional( + Schema.Struct({ + worktreePath: Schema.String, + branch: Schema.String, + expectedWorktreePath: Schema.NullOr(Schema.String), + expectedBranch: Schema.NullOr(Schema.String), + }), + ), }), Schema.Struct({ type: Schema.Literal("prepared-run.fail"),