diff --git a/apps/server/src/orchestration-v2/Orchestrator.ts b/apps/server/src/orchestration-v2/Orchestrator.ts index abdfc29579cd..0271ac875f78 100644 --- a/apps/server/src/orchestration-v2/Orchestrator.ts +++ b/apps/server/src/orchestration-v2/Orchestrator.ts @@ -2370,13 +2370,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) { @@ -5166,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), }; @@ -7564,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, @@ -7601,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 5e53250915f8..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,8 +97,11 @@ 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"]; + readonly closeTerminal?: TerminalManager.TerminalManager["Service"]["close"]; readonly fetchRemote?: GitWorkflow.GitWorkflowService["Service"]["fetchRemote"]; readonly renameBranch?: GitWorkflow.GitWorkflowService["Service"]["renameBranch"]; readonly runSetup?: ProjectSetupScriptRunner.ProjectSetupScriptRunner["Service"]["runForThread"]; @@ -105,6 +109,12 @@ interface HarnessOptions { readonly generateBranchName?: TextGeneration.TextGeneration["Service"]["generateBranchName"]; readonly serverSettings?: Parameters[0]; readonly providers?: ReadonlyArray; + readonly beforeLaunchDispatch?: ( + command: Parameters[0], + ) => Effect.Effect; + readonly afterLaunchDispatch?: ( + command: Parameters[0], + ) => Effect.Effect; } function makeHarness(options: HarnessOptions = {}) { @@ -128,10 +138,8 @@ function makeHarness(options: HarnessOptions = {}) { const renameBranch = vi.fn( options.renameBranch ?? ((input) => Effect.succeed({ branch: input.newBranch })), ); - const removeWorktree = vi.fn( - (_input: Parameters[0]) => - Effect.void, - ); + 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 })), ); @@ -142,9 +150,10 @@ 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: () => Effect.void }), + Layer.mock(TerminalManager.TerminalManager)({ close: closeTerminal }), Layer.succeed(ProjectService.ProjectService, { create: () => Effect.die("unused"), bootstrap: () => Effect.die("unused"), @@ -188,8 +197,29 @@ function makeHarness(options: HarnessOptions = {}) { folderForThread: () => Effect.succeed(Option.none()), }), ); + const beforeLaunchDispatch = options.beforeLaunchDispatch; + 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, threadManagement, receipts, IdAllocator.layer)), + Layer.provide( + Layer.mergeAll(externalServices, launchThreadManagement, receipts, IdAllocator.layer), + ), ); const projectedProjects = Layer.mock(ProjectStore.ProjectStoreV2)({ get: (requestedProjectId) => @@ -223,9 +253,12 @@ function makeHarness(options: HarnessOptions = {}) { outbox, database, externalServices, + IdAllocator.layer, + receipts, ), createWorktree, removeWorktree, + closeTerminal, renameBranch, generateBranchName, generateThreadTitle, @@ -537,9 +570,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(); @@ -547,7 +580,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* () { @@ -561,6 +598,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({ @@ -584,22 +623,27 @@ 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, "/repo-worktrees/feature"); assert.equal( projection.checkpointScopes.find((scope) => scope.id === rootNode?.checkpointScopeId) ?.cwd, @@ -1431,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. @@ -1485,14 +1846,22 @@ 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.equal( + projection.thread.worktreePath, + failurePoint === "setup" ? "/repo-worktrees/feature" : null, + ); + if (failurePoint === "worktree") assert.isNull(projection.thread.branch); + assert.isEmpty(harness.removeWorktree.mock.calls); assert.equal( projection.turnItems.find((item) => item.type === "command_execution")?.status, "failed", @@ -1554,6 +1923,361 @@ it.effect("replays a server-allocated launch", () => }), ); +it.effect.each(["worktree", "existing_worktree"] as const)( + "keeps a failed setup terminal and its recorded %s workspace available for retry", + (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.isEmpty(harness.closeTerminal.mock.calls); + assert.isEmpty(harness.removeWorktree.mock.calls); + assert.equal( + projection.thread.worktreePath, + workspaceType === "worktree" ? "/repo-worktrees/feature" : "/existing", + ); + assert.equal(projection.thread.branch, workspaceType === "worktree" ? "feature" : "existing"); + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + if (workspaceType === "worktree") + assert.equal((yield* tracker.get(launched.threadId))?.setupScript?.terminalId, "setup"); + }).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 newer workspace binding after removing a cancelled launch worktree", () => + Effect.gen(function* () { + const setupEntered = yield* Deferred.make(); + const harness = makeHarness({ + runSetup: () => Deferred.succeed(setupEntered, undefined).pipe(Effect.andThen(Effect.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", + }); + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + assert.isTrue(yield* tracker.cancel(launched.threadId)); + 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 recorded failed launch worktree without attempting cleanup", () => { + 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"); + assert.isEmpty(harness.removeWorktree.mock.calls); + }).pipe(Effect.provide(harness.layer)); +}); + +it.effect("stops branch renaming before removing a cancelled 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); + const tracker = yield* WorktreeSetupTracker.WorktreeSetupTracker; + assert.isTrue(yield* tracker.cancel(launched.threadId)); + 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 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.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", + }); + } + assert.isTrue(yield* tracker.cancel(launched.threadId)); + 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* () { @@ -1983,8 +2707,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( @@ -2134,7 +2860,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; @@ -2152,14 +2888,72 @@ 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.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* () { 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 555d3266eba2..af15f4ebfb2f 100644 --- a/apps/server/src/orchestration-v2/ThreadLaunchService.ts +++ b/apps/server/src/orchestration-v2/ThreadLaunchService.ts @@ -23,6 +23,8 @@ 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 FileSystem from "effect/FileSystem"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; import * as Ref from "effect/Ref"; @@ -172,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; @@ -243,7 +246,11 @@ const make = Effect.gen(function* () { const reused = input.reusedWorktree; const tracked = input.workspaceStrategy.type === "worktree" || reused !== undefined; 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; let workspaceRecorded = false; if (input.workspaceStrategy.type === "worktree") { yield* setupTracker.begin({ @@ -301,19 +308,24 @@ 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; + 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 : 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({ @@ -369,7 +381,7 @@ const make = Effect.gen(function* () { { cwd: project.workspaceRoot, refName: startRef, - newRefName: branch!, + newRefName: worktreeBranch, baseRefName: input.workspaceStrategy.baseRef, path: null, }, @@ -378,6 +390,7 @@ const make = Effect.gen(function* () { onWorktreeClaimed: (path) => Effect.sync(() => { createdWorktreePath = path; + createdWorktreeBranch = worktreeBranch; }), onCheckoutProgress: (progress) => setupTracker.stage(threadId, "checkout", { percent: progress.percent }), @@ -388,6 +401,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"); } @@ -396,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; @@ -420,20 +455,32 @@ 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, - 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({ type: "thread.metadata.update", commandId: CommandId.make(`${input.commandId}:branch-rename`), threadId, + expectedWorktreePath: worktreeCwd, + expectedBranch: oldBranch, branch: renamed.branch, worktreePath: worktreeCwd, }), @@ -496,9 +543,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"}`, @@ -538,7 +586,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; @@ -547,47 +605,122 @@ 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)), - ); - // A cancelled setup leaves nothing behind. A failed one keeps a worktree - // the thread recorded, so a retry reuses it, and removes one it never - // recorded, which a retry would otherwise duplicate. - if (tracked && createdWorktreePath && (cancelled || !workspaceRecorded)) { - if (setupTerminalId) + // Failed recorded checkouts and their setup output remain available for + // retry. Released runs and prepared idle launches also keep their checkout. + const cleanup = !preparationComplete && (cancelled || !workspaceRecorded); + if (cleanup) { + if (branchRenameFiber !== null) yield* Fiber.interrupt(branchRenameFiber); + if (setupTerminalId !== null) yield* terminals .close({ threadId, terminalId: setupTerminalId, deleteHistory: true }) - .pipe(Effect.ignore); - const removedPath = createdWorktreePath; - // The thread forgets the worktree only once it is gone; a failed - // removal leaves the directory for the user to clean up rather than - // reusing a checkout that may be half written. - yield* git - .removeWorktree({ cwd: project.workspaceRoot, path: removedPath, force: true }) + .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 (cleanup && tracked && createdWorktreePath !== null) { + const removed = yield* git + .removeWorktree({ + cwd: project.workspaceRoot, + path: createdWorktreePath, + force: true, + }) .pipe( - Effect.andThen( - threads + Effect.as(true), + Effect.catchCause((cleanupCause) => + Effect.logWarning("Failed to remove thread launch worktree", { + threadId, + worktreePath: createdWorktreePath, + cause: cleanupCause, + }).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) { + 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 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}:cancel-workspace`), + commandId: CommandId.make(`${input.commandId}:cleanup-workspace`), threadId, - worktreePath: null, - branch: null, + expectedWorktreePath: shell.worktreePath, + expectedBranch: shell.branch, + worktreePath: createdWorktreePath, + branch: survivingBranch, }) - .pipe(Effect.ignore), - ), - Effect.catchCause((removeCause) => - Effect.logWarning("Failed to remove an abandoned thread worktree", { - commandId: input.commandId, - threadId, - path: removedPath, - cause: removeCause, - }), - ), - ); + .pipe(Effect.ignore); + } + yield* setupTracker.update(threadId, (snapshot) => ({ + ...snapshot, + worktreePath: createdWorktreePath, + branch: survivingBranch, + })); + } + } } + yield* setupTracker.finish( + threadId, + cancelled ? "cancelled" : "failed", + cancelled ? null : failureDetail(Cause.squash(cause)), + ); }), ), ); @@ -923,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/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts b/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts index ef12204541b6..f931fd7eb421 100644 --- a/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts +++ b/apps/server/src/orchestration-v2/ThreadSettlementService.test.ts @@ -270,8 +270,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 7c64ce02804a..4d78ca173504 100644 --- a/apps/server/src/orchestration-v2/ThreadSettlementService.ts +++ b/apps/server/src/orchestration-v2/ThreadSettlementService.ts @@ -193,16 +193,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/orchestration-v2/runtimeLayer.test.ts b/apps/server/src/orchestration-v2/runtimeLayer.test.ts index c2f0024b0542..0b586a071743 100644 --- a/apps/server/src/orchestration-v2/runtimeLayer.test.ts +++ b/apps/server/src/orchestration-v2/runtimeLayer.test.ts @@ -1677,6 +1677,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/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/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 eb7c4e6a9ffd..88081f3f20ba 100644 --- a/docs/user/thread-sidebar.md +++ b/docs/user/thread-sidebar.md @@ -13,6 +13,18 @@ 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 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. 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. 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. + ### Start without a project A thread does not need a project. To start one without a project, click **or @@ -141,7 +153,8 @@ Press `Escape` while dragging to cancel. 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 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 resumed after it closed. Only your own messages count as resuming. A turn that diff --git a/packages/contracts/src/orchestrationV2.ts b/packages/contracts/src/orchestrationV2.ts index 220fec2b15fb..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; @@ -2612,6 +2614,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)), @@ -2759,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"),