From 046e2883475f347b85a4e77d3b100aac94392507 Mon Sep 17 00:00:00 2001 From: Adam Firestone Date: Thu, 13 Aug 2026 13:05:31 -0500 Subject: [PATCH] fix(server): recover pruned thread worktrees --- .../OrchestrationEngineHarness.integration.ts | 6 + .../Layers/ProviderCommandReactor.test.ts | 224 ++++++++++++++++++ .../Layers/ProviderCommandReactor.ts | 123 +++++++++- apps/server/src/server.ts | 2 +- 4 files changed, 352 insertions(+), 3 deletions(-) diff --git a/apps/server/integration/OrchestrationEngineHarness.integration.ts b/apps/server/integration/OrchestrationEngineHarness.integration.ts index 7f00a05c1989..faaf3e23c77c 100644 --- a/apps/server/integration/OrchestrationEngineHarness.integration.ts +++ b/apps/server/integration/OrchestrationEngineHarness.integration.ts @@ -47,6 +47,7 @@ import { ProviderService } from "../src/provider/Services/ProviderService.ts"; import { AnalyticsService } from "../src/telemetry/Services/AnalyticsService.ts"; import { CheckpointReactorLive } from "../src/orchestration/Layers/CheckpointReactor.ts"; import * as RepositoryIdentityResolver from "../src/project/RepositoryIdentityResolver.ts"; +import { ProjectSetupScriptRunner } from "../src/project/ProjectSetupScriptRunner.ts"; import { OrchestrationEngineLive } from "../src/orchestration/Layers/OrchestrationEngine.ts"; import { OrchestrationProjectionPipelineLive } from "../src/orchestration/Layers/ProjectionPipeline.ts"; import { OrchestrationProjectionSnapshotQueryLive } from "../src/orchestration/Layers/ProjectionSnapshotQuery.ts"; @@ -335,6 +336,11 @@ export const makeOrchestrationIntegrationHarness = ( const providerCommandReactorLayer = ProviderCommandReactorLive.pipe( Layer.provideMerge(runtimeServicesLayer), Layer.provideMerge(gitWorkflowLayer), + Layer.provideMerge( + Layer.succeed(ProjectSetupScriptRunner, { + runForThread: () => Effect.succeed({ status: "no-script" }), + }), + ), Layer.provideMerge(textGenerationLayer), Layer.provideMerge(serverSettingsLayer), ); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 0b0bfff5ffed..305da9274ca7 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -9,6 +9,7 @@ import { ProviderSession, ProviderDriverKind, ProviderInstanceId, + type VcsRef, } from "@t3tools/contracts"; import { createModelSelection } from "@t3tools/shared/model"; import { @@ -46,12 +47,14 @@ import { import { makeProviderRegistryLayer } from "../../provider/testUtils/providerRegistryMock.ts"; import { TextGeneration, type TextGenerationShape } from "../../textGeneration/TextGeneration.ts"; import * as RepositoryIdentityResolver from "../../project/RepositoryIdentityResolver.ts"; +import { ProjectSetupScriptRunner } from "../../project/ProjectSetupScriptRunner.ts"; import { OrchestrationEngineLive } from "./OrchestrationEngine.ts"; import { OrchestrationProjectionPipelineLive } from "./ProjectionPipeline.ts"; import { OrchestrationProjectionSnapshotQueryLive } from "./ProjectionSnapshotQuery.ts"; import * as ThreadBackgroundLiveness from "../ThreadBackgroundLiveness.ts"; import * as ThreadPlanProgress from "../ThreadPlanProgress.ts"; import { + managedWorktreeRecoveryPlan, providerErrorLabel, providerErrorLabelFromInstanceHint, ProviderCommandReactorLive, @@ -145,6 +148,100 @@ describe("ProviderCommandReactor", () => { }); }); + describe("managed worktree recovery planning", () => { + it("reattaches an existing local branch", () => { + expect( + managedWorktreeRecoveryPlan({ + branch: "t3code/fix", + path: "/repo/worktrees/fix", + refs: [ + { + name: "t3code/fix", + current: false, + isDefault: false, + isRemote: false, + worktreePath: null, + }, + ], + checkpoints: [], + }), + ).toEqual({ + refName: "t3code/fix", + path: "/repo/worktrees/fix", + }); + }); + + it("recreates a deleted branch from its latest ready checkpoint", () => { + expect( + managedWorktreeRecoveryPlan({ + branch: "t3code/merged-fix", + path: "/repo/worktrees/merged-fix", + refs: [ + { + name: "main", + current: true, + isDefault: true, + isRemote: false, + worktreePath: "/repo", + }, + { + name: "origin/t3code/merged-fix", + remoteName: "origin", + current: false, + isDefault: false, + isRemote: true, + worktreePath: null, + }, + ], + checkpoints: [ + { + turnId: asTurnId("turn-1"), + checkpointTurnCount: 1, + checkpointRef: CheckpointRef.make("refs/t3/checkpoints/old"), + status: "ready", + files: [], + assistantMessageId: null, + completedAt: "2026-01-01T00:00:00.000Z", + }, + { + turnId: asTurnId("turn-2"), + checkpointTurnCount: 2, + checkpointRef: CheckpointRef.make("refs/t3/checkpoints/latest"), + status: "ready", + files: [], + assistantMessageId: null, + completedAt: "2026-01-01T00:01:00.000Z", + }, + ], + }), + ).toEqual({ + refName: "refs/t3/checkpoints/latest", + newRefName: "t3code/merged-fix", + baseRefName: "main", + path: "/repo/worktrees/merged-fix", + }); + }); + + it("does not steal a branch checked out by another worktree", () => { + expect( + managedWorktreeRecoveryPlan({ + branch: "t3code/fix", + path: "/repo/worktrees/fix", + refs: [ + { + name: "t3code/fix", + current: false, + isDefault: false, + isRemote: false, + worktreePath: "/repo/worktrees/other", + }, + ], + checkpoints: [], + }), + ).toBeNull(); + }); + }); + async function createHarness(input?: { readonly baseDir?: string; readonly threadModelSelection?: ModelSelection; @@ -157,6 +254,12 @@ describe("ProviderCommandReactor", () => { ) => Effect.Effect; readonly beforeThreadDetailRead?: (callIndex: number) => Effect.Effect; readonly ensurePreTurnBaselineEffect?: () => Effect.Effect; + readonly missingManagedWorktree?: { + readonly branch: string; + readonly path: string; + readonly refs: ReadonlyArray; + readonly checkpointRef?: CheckpointRef; + }; }) { const now = "2026-01-01T00:00:00.000Z"; const baseDir = @@ -271,6 +374,29 @@ describe("ProviderCommandReactor", () => { : "renamed-branch", }), ); + const listRefs = vi.fn(() => + Effect.succeed({ + refs: [...(input?.missingManagedWorktree?.refs ?? [])], + isRepo: true, + hasPrimaryRemote: true, + nextCursor: null, + totalCount: input?.missingManagedWorktree?.refs.length ?? 0, + }), + ); + const createWorktree = vi.fn( + ( + worktreeInput: Parameters< + GitWorkflowService.GitWorkflowService["Service"]["createWorktree"] + >[0], + ) => + Effect.succeed({ + worktree: { + path: worktreeInput.path ?? "/tmp/provider-created-worktree", + refName: worktreeInput.newRefName ?? worktreeInput.refName, + }, + }), + ); + const runSetupScript = vi.fn(() => Effect.succeed({ status: "no-script" as const })); const refreshStatus = vi.fn((_: string) => Effect.succeed({ isRepo: true, @@ -427,9 +553,12 @@ describe("ProviderCommandReactor", () => { Layer.provideMerge(makeProviderRegistryLayer(providerSnapshots as never)), Layer.provideMerge( Layer.mock(GitWorkflowService.GitWorkflowService)({ + createWorktree, + listRefs, renameBranch, } satisfies Partial), ), + Layer.provideMerge(Layer.succeed(ProjectSetupScriptRunner, { runForThread: runSetupScript })), Layer.provideMerge( Layer.succeed(VcsStatusBroadcaster, { getStatus: () => Effect.die("getStatus should not be called in this test"), @@ -483,6 +612,37 @@ describe("ProviderCommandReactor", () => { createdAt: now, }), ); + if (input?.missingManagedWorktree) { + await Effect.runPromise( + engine.dispatch({ + type: "thread.managed-worktree.record", + commandId: CommandId.make("cmd-managed-worktree-record"), + threadId: ThreadId.make("thread-1"), + branch: input.missingManagedWorktree.branch, + managedWorktree: { + projectCwd: "/tmp/provider-project", + path: input.missingManagedWorktree.path, + createdForCommandId: CommandId.make("cmd-thread-create"), + }, + }), + ); + if (input.missingManagedWorktree.checkpointRef !== undefined) { + await Effect.runPromise( + engine.dispatch({ + type: "thread.turn.diff.complete", + commandId: CommandId.make("cmd-managed-worktree-checkpoint"), + threadId: ThreadId.make("thread-1"), + turnId: asTurnId("turn-before-restart"), + checkpointTurnCount: 1, + checkpointRef: input.missingManagedWorktree.checkpointRef, + status: "ready", + files: [], + completedAt: now, + createdAt: now, + }), + ); + } + } if (input?.titleRegenerationBeforeStart === "two") { await Effect.runPromise( engine.dispatch({ @@ -533,6 +693,9 @@ describe("ProviderCommandReactor", () => { respondToUserInput, stopSession, renameBranch, + listRefs, + createWorktree, + runSetupScript, refreshStatus, generateBranchName, generateThreadTitle, @@ -587,6 +750,67 @@ describe("ProviderCommandReactor", () => { expect(thread?.session?.runtimeMode).toBe("approval-required"); }); + it("recreates a pruned managed worktree before resuming a provider turn", async () => { + const baseDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3code-reactor-recovery-")); + const worktreePath = NodePath.join(baseDir, "missing-worktree"); + const checkpointRef = CheckpointRef.make("refs/t3/checkpoints/thread-1/turn/1"); + const harness = await createHarness({ + baseDir, + missingManagedWorktree: { + branch: "t3code/merged-fix", + path: worktreePath, + checkpointRef, + refs: [ + { + name: "main", + current: true, + isDefault: true, + isRemote: false, + worktreePath: "/tmp/provider-project", + }, + ], + }, + }); + const now = "2026-01-01T00:02:00.000Z"; + + await Effect.runPromise( + harness.engine.dispatch({ + type: "thread.turn.start", + commandId: CommandId.make("cmd-turn-start-after-prune"), + threadId: ThreadId.make("thread-1"), + message: { + messageId: asMessageId("user-message-after-prune"), + role: "user", + text: "continue after restart", + attachments: [], + }, + interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE, + runtimeMode: "approval-required", + createdAt: now, + }), + ); + + await waitFor(() => harness.createWorktree.mock.calls.length === 1); + await waitFor(() => harness.startSession.mock.calls.length === 1); + await waitFor(() => harness.sendTurn.mock.calls.length === 1); + expect(harness.createWorktree.mock.calls[0]?.[0]).toEqual({ + cwd: "/tmp/provider-project", + refName: checkpointRef, + newRefName: "t3code/merged-fix", + baseRefName: "main", + path: worktreePath, + }); + expect(harness.startSession.mock.calls[0]?.[1]).toMatchObject({ + cwd: worktreePath, + }); + expect(harness.runSetupScript).toHaveBeenCalledWith({ + threadId: ThreadId.make("thread-1"), + projectId: asProjectId("project-1"), + projectCwd: "/tmp/provider-project", + worktreePath, + }); + }); + effectIt.effect("cancels an unclaimed retraction before provider session creation", () => Effect.gen(function* () { const readEntered = yield* Deferred.make(); diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts index 927adb8a85d0..e39d094ec774 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.ts @@ -3,16 +3,22 @@ import { CommandId, EventId, type ModelSelection, + type OrchestrationCheckpointSummary, type OrchestrationEvent, + type OrchestrationProjectShell, ProviderDriverKind, type ProjectId, type OrchestrationSession, + type OrchestrationThread, ThreadId, type ProviderSession, type RuntimeMode, type TurnId, + type VcsCreateWorktreeInput, + type VcsRef, } from "@t3tools/contracts"; import { isTemporaryWorktreeBranch, WORKTREE_BRANCH_PREFIX } from "@t3tools/shared/git"; +import * as FileSystem from "effect/FileSystem"; import * as Cache from "effect/Cache"; import * as Cause from "effect/Cause"; import * as Crypto from "effect/Crypto"; @@ -32,6 +38,7 @@ import type { ProviderServiceError } from "../../provider/Errors.ts"; import { TextGeneration } from "../../textGeneration/TextGeneration.ts"; import { ProviderService } from "../../provider/Services/ProviderService.ts"; import { ProviderRegistry } from "../../provider/Services/ProviderRegistry.ts"; +import { ProjectSetupScriptRunner } from "../../project/ProjectSetupScriptRunner.ts"; import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; import { CheckpointReactor } from "../Services/CheckpointReactor.ts"; @@ -101,6 +108,55 @@ const MAX_FIRST_USER_TITLE_CONTEXT_CHARS = 2_000; const THREAD_TITLE_CONTEXT_TRUNCATION_MARKER = "[Earlier content truncated]\n\n"; const FIRST_USER_CONTEXT_TRUNCATION_MARKER = "\n[First user message truncated]"; +type ManagedWorktreeRecoveryPlan = Pick< + VcsCreateWorktreeInput, + "refName" | "newRefName" | "baseRefName" | "path" +>; + +export function managedWorktreeRecoveryPlan(input: { + readonly branch: string; + readonly path: string; + readonly refs: ReadonlyArray; + readonly checkpoints: ReadonlyArray; +}): ManagedWorktreeRecoveryPlan | null { + const localBranch = input.refs.find((ref) => ref.isRemote !== true && ref.name === input.branch); + if (localBranch?.worktreePath === null) { + return { + refName: input.branch, + path: input.path, + }; + } + + if (localBranch !== undefined) { + return null; + } + + const defaultRef = + input.refs.find((ref) => ref.isRemote !== true && ref.current) ?? + input.refs.find((ref) => ref.isRemote !== true && ref.isDefault) ?? + input.refs.find((ref) => ref.isDefault); + const latestCheckpoint = input.checkpoints + .toReversed() + .find((checkpoint) => checkpoint.status === "ready"); + const remoteBranch = input.refs.find( + (ref) => + ref.isRemote === true && + (ref.name === input.branch || + (ref.remoteName !== undefined && ref.name === `${ref.remoteName}/${input.branch}`)), + ); + const recoveryRef = latestCheckpoint?.checkpointRef ?? remoteBranch?.name ?? defaultRef?.name; + if (recoveryRef === undefined) { + return null; + } + + return { + refName: recoveryRef, + newRefName: input.branch, + ...(defaultRef !== undefined ? { baseRefName: defaultRef.name } : {}), + path: input.path, + }; +} + type ThreadTitleMessage = { readonly role: "user" | "assistant" | "system"; readonly text: string; @@ -322,6 +378,8 @@ const make = Effect.gen(function* () { const providerService = yield* ProviderService; const providerRegistry = yield* ProviderRegistry; const gitWorkflow = yield* GitWorkflowService; + const fileSystem = yield* FileSystem.FileSystem; + const projectSetupScriptRunner = yield* ProjectSetupScriptRunner; const vcsStatusBroadcaster = yield* VcsStatusBroadcaster; const textGeneration = yield* TextGeneration; const serverSettingsService = yield* ServerSettingsService; @@ -483,6 +541,66 @@ const make = Effect.gen(function* () { }); }); + const ensureThreadWorkspaceCwd = Effect.fn("ensureThreadWorkspaceCwd")(function* (input: { + readonly thread: OrchestrationThread; + readonly project: OrchestrationProjectShell | undefined; + readonly provider: ProviderDriverKind; + }) { + const cwd = resolveThreadWorkspaceCwd({ + thread: input.thread, + projects: input.project ? [input.project] : [], + }); + const managedWorktree = input.thread.managedWorktree; + if ( + input.thread.worktreePath === null || + managedWorktree == null || + managedWorktree.path !== input.thread.worktreePath || + input.thread.branch === null || + (yield* fileSystem.exists(input.thread.worktreePath)) + ) { + return cwd; + } + + const refs = yield* gitWorkflow.listRefs({ + cwd: managedWorktree.projectCwd, + includeMatchingRemoteRefs: true, + refresh: true, + limit: 200, + }); + const recoveryPlan = managedWorktreeRecoveryPlan({ + branch: input.thread.branch, + path: managedWorktree.path, + refs: refs.refs, + checkpoints: input.thread.checkpoints, + }); + if (!refs.isRepo || recoveryPlan === null) { + return yield* new ProviderAdapterRequestError({ + provider: input.provider, + method: "thread.turn.start", + detail: `Thread '${input.thread.id}' cannot resume because its managed worktree '${managedWorktree.path}' no longer exists and no safe Git ref is available to recreate it.`, + }); + } + + yield* gitWorkflow.createWorktree({ + cwd: managedWorktree.projectCwd, + ...recoveryPlan, + }); + yield* projectSetupScriptRunner.runForThread({ + threadId: input.thread.id, + projectId: input.thread.projectId, + projectCwd: managedWorktree.projectCwd, + worktreePath: managedWorktree.path, + }); + yield* Effect.logInfo("recreated missing managed thread worktree", { + threadId: input.thread.id, + branch: input.thread.branch, + worktreePath: managedWorktree.path, + recoveryRef: recoveryPlan.refName, + createdBranch: recoveryPlan.newRefName !== undefined, + }); + return managedWorktree.path; + }); + const ensureSessionForThread = Effect.fn("ensureSessionForThread")(function* ( threadId: ThreadId, createdAt: string, @@ -617,9 +735,10 @@ const make = Effect.gen(function* () { } } const project = yield* resolveProject(thread.projectId); - const effectiveCwd = resolveThreadWorkspaceCwd({ + const effectiveCwd = yield* ensureThreadWorkspaceCwd({ thread, - projects: project ? [project] : [], + project, + provider: preferredProvider, }); const startProviderSession = (input?: { diff --git a/apps/server/src/server.ts b/apps/server/src/server.ts index a76e0edb75b2..1f949372f5bd 100644 --- a/apps/server/src/server.ts +++ b/apps/server/src/server.ts @@ -373,7 +373,7 @@ const RuntimeCoreDependenciesLive = ReactorLayerLive.pipe( Layer.provideMerge(CheckpointingLayerLive), Layer.provideMerge(SourceControlProviderRegistryLayerLive), Layer.provideMerge(GitLayerLive), - Layer.provideMerge(VcsLayerLive), + Layer.provideMerge(Layer.mergeAll(VcsLayerLive, ProjectSetupScriptRunner.layer)), Layer.provideMerge(ProviderRuntimeLayerLive), Layer.provideMerge(Layer.mergeAll(TerminalLayerLive, PreviewLayerLive)), Layer.provideMerge(PersistenceLayerLive),