diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index d58c014d86e8..19959e9fb4d6 100644 --- a/apps/server/src/environment/ServerEnvironment.ts +++ b/apps/server/src/environment/ServerEnvironment.ts @@ -223,6 +223,7 @@ export const make = Effect.gen(function* () { inlineMessageContext: true, requiredWorktreeBootstrap: true, threadSettlement: true, + recoverableThreadDeletion: true, threadAutoSettlement: true, storageCleanup: true, projectWorktreeCleanup: true, diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 9ccf78ca1744..48d204fe85f4 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts @@ -20,6 +20,8 @@ import { import * as NodeServices from "@effect/platform-node/NodeServices"; import { it as effectIt } from "@effect/vitest"; import * as Effect from "effect/Effect"; +import * as Deferred from "effect/Deferred"; +import * as Fiber from "effect/Fiber"; import * as Layer from "effect/Layer"; import * as ManagedRuntime from "effect/ManagedRuntime"; import * as Metric from "effect/Metric"; @@ -54,6 +56,7 @@ import { } from "../Services/ProjectionPipeline.ts"; import { ProjectionSnapshotQuery } from "../Services/ProjectionSnapshotQuery.ts"; import { ServerConfig } from "../../config.ts"; +import { removeUnusedWorktree } from "../threadWorktreeDeletion.ts"; const asProjectId = (value: string): ProjectId => ProjectId.make(value); const asMessageId = (value: string): MessageId => MessageId.make(value); @@ -130,6 +133,84 @@ const hasMetricSnapshot = ( ); describe("OrchestrationEngine", () => { + effectIt.effect.each(["thread.create", "thread.meta.update"] as const)( + "prevents concurrent %s adoption during cleanup without blocking unrelated commands", + (type) => + Effect.gen(function* () { + const engine = yield* OrchestrationEngineService; + const snapshots = yield* ProjectionSnapshotQuery; + const projectId = ProjectId.make("cleanup-project"); + const threadId = ThreadId.make("cleanup-adopter"); + yield* engine.dispatch({ + type: "project.create", + commandId: CommandId.make("cleanup-project"), + projectId, + title: "Test", + workspaceRoot: "/repo", + createdAt: now(), + }); + const create = { + type: "thread.create", + commandId: CommandId.make("cleanup-create"), + projectId, + threadId, + title: "Test", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "test" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: null, + worktreePath: "/repo/staged", + createdAt: now(), + } as const; + if (type === "thread.meta.update") + yield* engine.dispatch({ ...create, worktreePath: null }); + const started = yield* Deferred.make(); + const release = yield* Deferred.make(); + const cleanup = yield* removeUnusedWorktree( + { cwd: "/repo", path: create.worktreePath, force: true }, + Deferred.succeed(started, undefined).pipe( + Effect.andThen(Deferred.await(release)), + Effect.andThen( + type === "thread.meta.update" ? Effect.die("cleanup failed") : Effect.void, + ), + ), + ).pipe(Effect.exit, Effect.forkChild({ startImmediately: true })); + yield* Deferred.await(started); + const adoption: OrchestrationCommand = + type === "thread.create" + ? create + : { + type, + commandId: CommandId.make("cleanup-adopt"), + threadId, + worktreePath: create.worktreePath, + }; + expect((yield* engine.dispatch(adoption).pipe(Effect.exit))._tag).toBe("Failure"); + expect( + (yield* snapshots.getCommandReadModel()).threads.some( + (thread) => thread.worktreePath === create.worktreePath, + ), + ).toBe(false); + // This receipt must arrive while filesystem removal is still blocked. + yield* engine.dispatch({ + type: "project.create", + commandId: CommandId.make("unrelated-project"), + projectId: ProjectId.make("unrelated-project"), + title: "Other work", + workspaceRoot: "/other", + createdAt: now(), + }); + yield* Deferred.succeed(release, undefined); + expect((yield* Fiber.join(cleanup))._tag).toBe( + type === "thread.meta.update" ? "Failure" : "Success", + ); + yield* engine.dispatch({ ...adoption, commandId: CommandId.make("retry-adoption") }); + expect((yield* snapshots.getCommandReadModel()).threads[0]?.worktreePath).toBe( + create.worktreePath, + ); + }).pipe(Effect.provide(makeOrchestrationLayer())), + ); + it.each(["running", "stopped"] as const)( "sends async answers with a %s session and rejects old duplicate replies", async (status) => { diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts index 9136d080c1c3..233d9bfa74bb 100644 --- a/apps/server/src/orchestration/Layers/OrchestrationEngine.ts +++ b/apps/server/src/orchestration/Layers/OrchestrationEngine.ts @@ -17,6 +17,7 @@ import * as Exit from "effect/Exit"; import * as Layer from "effect/Layer"; import * as Metric from "effect/Metric"; import * as Option from "effect/Option"; +import * as Path from "effect/Path"; import * as PubSub from "effect/PubSub"; import * as Queue from "effect/Queue"; import * as Schema from "effect/Schema"; @@ -89,11 +90,13 @@ const makeOrchestrationEngine = Effect.gen(function* () { const projectionSnapshotQuery = yield* ProjectionSnapshotQuery; const threadBackgroundLiveness = yield* ThreadBackgroundLivenessService; const crypto = yield* Crypto.Crypto; + const path = yield* Path.Path; + const worktreeCleanups = new Map(); const nowIso = Effect.map(DateTime.now, DateTime.formatIso); let commandReadModel = createEmptyReadModel(yield* nowIso); - const commandQueue = yield* Queue.unbounded(); + const commandQueue = yield* Queue.unbounded>(); const eventPubSub = yield* PubSub.unbounded(); const projectEventsOntoReadModel = ( @@ -242,6 +245,17 @@ const makeOrchestrationEngine = Effect.gen(function* () { envelope.command.type === "thread.user-input.dismiss" ? yield* projectionSnapshotQuery.getUserInputActivity(envelope.command) : Option.none(); + if ( + (envelope.command.type === "thread.create" || + envelope.command.type === "thread.meta.update") && + envelope.command.worktreePath && + worktreeCleanups.has(path.resolve(envelope.command.worktreePath)) + ) { + return yield* new OrchestrationCommandInvariantError({ + commandType: envelope.command.type, + detail: "Worktree cleanup is in progress. Try again after it finishes.", + }); + } const eventBase = yield* decideOrchestrationCommand({ command: envelope.command, readModel: commandReadModel, @@ -413,7 +427,7 @@ const makeOrchestrationEngine = Effect.gen(function* () { yield* projectionPipeline.bootstrap; commandReadModel = yield* projectionSnapshotQuery.getCommandReadModel(); - const worker = Effect.forever(Queue.take(commandQueue).pipe(Effect.flatMap(processEnvelope))); + const worker = Effect.forever(Queue.take(commandQueue).pipe(Effect.flatten)); yield* Effect.forkScoped(worker); yield* Effect.logDebug("orchestration engine started").pipe( Effect.annotateLogs({ sequence: commandReadModel.snapshotSequence }), @@ -438,20 +452,61 @@ const makeOrchestrationEngine = Effect.gen(function* () { const dispatch: OrchestrationEngineShape["dispatch"] = (command, options) => Effect.gen(function* () { const result = yield* Deferred.make<{ sequence: number }, OrchestrationDispatchError>(); - yield* Queue.offer(commandQueue, { + const envelope: CommandEnvelope = { command, origin: options?.origin, result, startedAtMs: yield* Clock.currentTimeMillis, - }); + }; + yield* Queue.offer( + commandQueue, + Effect.suspend(() => processEnvelope(envelope)), + ); return yield* Deferred.await(result); }); + const runSerialized = (effect: Effect.Effect) => + Effect.gen(function* () { + const result = yield* Deferred.make(); + yield* Queue.offer( + commandQueue, + effect.pipe( + Effect.exit, + Effect.flatMap((exit) => Deferred.done(result, exit)), + Effect.asVoid, + ), + ); + return yield* Deferred.await(result); + }); + + const withWorktreeCleanup: OrchestrationEngineShape["withWorktreeCleanup"] = (paths, cleanup) => { + const keys = [...new Set(paths.map((entry) => path.resolve(entry)))]; + return Effect.acquireUseRelease( + runSerialized( + Effect.sync(() => { + for (const key of keys) worktreeCleanups.set(key, (worktreeCleanups.get(key) ?? 0) + 1); + }), + ), + () => cleanup, + () => + runSerialized( + Effect.sync(() => { + for (const key of keys) { + const remaining = worktreeCleanups.get(key)! - 1; + if (remaining === 0) worktreeCleanups.delete(key); + else worktreeCleanups.set(key, remaining); + } + }), + ), + ); + }; + return { readEvents, readThreadEvents, getThreadReplayStats, dispatch, + withWorktreeCleanup, subscribeDomainEvents: PubSub.subscribe(eventPubSub).pipe(Effect.map(Stream.fromSubscription)), // Each access creates a fresh PubSub subscription so that multiple // consumers (wsServer, ProviderRuntimeIngestion, CheckpointReactor, etc.) diff --git a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts index 87f9bd03d46c..00db50513b2b 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -425,6 +425,7 @@ describe("ProviderCommandReactor", () => { Effect.gen(function* () { const engine = yield* OrchestrationEngineService; return { + withWorktreeCleanup: engine.withWorktreeCleanup, readEvents: engine.readEvents, readThreadEvents: engine.readThreadEvents, getThreadReplayStats: engine.getThreadReplayStats, diff --git a/apps/server/src/orchestration/Services/OrchestrationEngine.ts b/apps/server/src/orchestration/Services/OrchestrationEngine.ts index cf2ef2a40354..3dcf44012286 100644 --- a/apps/server/src/orchestration/Services/OrchestrationEngine.ts +++ b/apps/server/src/orchestration/Services/OrchestrationEngine.ts @@ -75,6 +75,12 @@ export interface OrchestrationEngineShape { options?: { readonly origin?: OrchestrationClientOrigin }, ) => Effect.Effect<{ sequence: number }, OrchestrationDispatchError, never>; + /** Reserve paths on the command queue, preventing adoption while filesystem cleanup runs outside it. */ + readonly withWorktreeCleanup: ( + paths: ReadonlyArray, + cleanup: Effect.Effect, + ) => Effect.Effect; + /** * Stream persisted domain events in dispatch order. * diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 21df39a162e6..1153fbf2b1ba 100644 --- a/apps/server/src/orchestration/decider.ts +++ b/apps/server/src/orchestration/decider.ts @@ -21,6 +21,7 @@ import { threadPullRequestKeysEqual, } from "@t3tools/shared/threadPullRequests"; import { compareDateTimeStrings } from "@t3tools/shared/dateTime"; +import { normalizeProjectPathForComparison } from "@t3tools/shared/path"; import * as DateTime from "effect/DateTime"; import * as Crypto from "effect/Crypto"; import * as Effect from "effect/Effect"; @@ -411,11 +412,30 @@ export const decideOrchestrationCommand = Effect.fn("decideOrchestrationCommand" } case "thread.delete": { - yield* requireThread({ + const thread = yield* requireThread({ readModel, command, threadId: command.threadId, }); + if (command.deleteWorktreePath) { + const worktreePath = normalizeProjectPathForComparison(command.deleteWorktreePath); + if ( + thread.worktreePath === null || + normalizeProjectPathForComparison(thread.worktreePath) !== worktreePath || + readModel.threads.some( + (entry) => + entry.id !== thread.id && + entry.deletedAt === null && + entry.worktreePath !== null && + normalizeProjectPathForComparison(entry.worktreePath) === worktreePath, + ) + ) { + return yield* new OrchestrationCommandInvariantError({ + commandType: command.type, + detail: "Worktree references changed during deletion. Try deleting the thread again.", + }); + } + } const occurredAt = yield* nowIso; return { ...(yield* withEventBase({ diff --git a/apps/server/src/orchestration/threadWorktreeDeletion.test.ts b/apps/server/src/orchestration/threadWorktreeDeletion.test.ts new file mode 100644 index 000000000000..c1b5b49d0d5c --- /dev/null +++ b/apps/server/src/orchestration/threadWorktreeDeletion.test.ts @@ -0,0 +1,438 @@ +import * as NodeServices from "@effect/platform-node/NodeServices"; +import { expect, it } from "@effect/vitest"; +import { + CommandId, + OrchestrationDispatchCommandError, + ProjectId, + ProviderInstanceId, + ThreadId, + type OrchestrationCommand, +} from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as Crypto from "effect/Crypto"; +import * as Schema from "effect/Schema"; +import * as FileSystem from "effect/FileSystem"; +import * as Layer from "effect/Layer"; +import * as Path from "effect/Path"; +import { ServerConfig } from "../config.ts"; +import * as Git from "../vcs/GitVcsDriver.ts"; +import { decideOrchestrationCommand } from "./decider.ts"; +import { createEmptyReadModel, projectEvent } from "./projector.ts"; +import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.ts"; +import { removeUnusedWorktree, withThreadWorktreeDeletion } from "./threadWorktreeDeletion.ts"; +import { PersistenceSqlError } from "../persistence/Errors.ts"; +import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts"; + +const TestLayer = Layer.merge( + Git.layer, + Layer.mock(OrchestrationEngineService)({ withWorktreeCleanup: (_paths, effect) => effect }), +).pipe( + Layer.provide(ServerConfig.layerTest(process.cwd(), { prefix: "t3-delete-test-" })), + Layer.provideMerge(NodeServices.layer), +); +const fixture = Effect.gen(function* () { + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const git = yield* Git.GitVcsDriver; + const root = yield* fs.makeTempDirectoryScoped({ prefix: "t3-delete-worktree-" }); + const cwd = path.join(root, "repo"); + const worktree = path.join(root, "worktree"); + yield* fs.makeDirectory(cwd); + const run = (args: string[], directory = cwd) => + git.execute({ operation: "test", cwd: directory, args }); + yield* run(["init", "--initial-branch=main"]); + yield* run(["config", "user.name", "Test"]); + yield* run(["config", "user.email", "test@example.com"]); + yield* fs.writeFileString(path.join(cwd, "tracked"), "original"); + yield* fs.writeFileString(path.join(cwd, ".gitignore"), "ignored\n"); + yield* run(["add", "."]); + yield* run(["commit", "-m", "initial"]); + yield* run(["worktree", "add", "-b", "test", worktree]); + yield* fs.writeFileString(path.join(worktree, "tracked"), "staged"); + yield* run(["add", "tracked"], worktree); + yield* fs.writeFileString(path.join(worktree, "tracked"), "unstaged"); + yield* fs.writeFileString(path.join(worktree, "untracked"), "untracked contents"); + yield* fs.writeFileString(path.join(worktree, "ignored"), "ignored contents"); + let snapshot = createEmptyReadModel("2026-09-17T00:00:00.000Z"); + const threadId = ThreadId.make("thread"); + const projectId = ProjectId.make("project"); + const commands: OrchestrationCommand[] = [ + { + type: "project.create", + commandId: CommandId.make("project"), + projectId, + title: "Test", + workspaceRoot: cwd, + createdAt: snapshot.updatedAt, + }, + { + type: "thread.create", + commandId: CommandId.make("thread"), + projectId, + threadId, + title: "Test", + modelSelection: { instanceId: ProviderInstanceId.make("codex"), model: "test" }, + runtimeMode: "full-access", + interactionMode: "default", + branch: "test", + worktreePath: worktree, + createdAt: snapshot.updatedAt, + }, + ]; + for (const command of commands) { + const decided = yield* decideOrchestrationCommand({ command, readModel: snapshot }); + for (const event of Array.isArray(decided) ? decided : [decided]) + snapshot = yield* projectEvent(snapshot, { + ...event, + sequence: snapshot.snapshotSequence + 1, + }); + } + const command = { + type: "thread.delete", + commandId: CommandId.make("delete"), + threadId, + deleteWorktreePath: worktree, + } as const; + const crypto = yield* Crypto.Crypto; + let snapshotError: PersistenceSqlError | null = null; + const snapshots = Layer.mock(ProjectionSnapshotQuery)({ + getCommandReadModel: () => + snapshotError ? Effect.fail(snapshotError) : Effect.succeed(snapshot), + }); + const execute = ( + commit: Effect.Effect<{ sequence: number }, E>, + afterCommit: Effect.Effect = Effect.void, + ) => + withThreadWorktreeDeletion(command, (staged) => + Effect.gen(function* () { + const result = yield* commit; + const decided = yield* decideOrchestrationCommand({ + command: staged ? command : { ...command, deleteWorktreePath: undefined }, + readModel: snapshot, + }); + for (const event of Array.isArray(decided) ? decided : [decided]) { + snapshot = yield* projectEvent(snapshot, { + ...event, + sequence: snapshot.snapshotSequence + 1, + }); + } + yield* afterCommit; + return result; + }).pipe( + Effect.provideService(Crypto.Crypto, crypto), + Effect.mapError((error) => + Schema.is(OrchestrationDispatchCommandError)(error) + ? error + : new OrchestrationDispatchCommandError({ + message: "test commit failed", + cause: error, + }), + ), + ), + ).pipe(Effect.provide(snapshots)); + + return { + fs, + path, + git, + run, + cwd, + root, + worktree, + command, + execute, + retryCleanup: (target: string) => { + const input = { cwd, path: target, force: true }; + return removeUnusedWorktree(input, git.removeWorktree(input)).pipe(Effect.provide(snapshots)); + }, + snapshot, + getSnapshot: () => snapshot, + failSnapshotRead: () => { + snapshotError = new PersistenceSqlError({ operation: "test snapshot read" }); + }, + setSnapshot: (next: typeof snapshot) => { + snapshot = next; + }, + }; +}); + +it.layer(TestLayer)("recoverable worktree deletion", (it) => { + it.effect( + "restores staged, unstaged, untracked and ignored contents when thread deletion fails", + () => + Effect.gen(function* () { + const f = yield* fixture; + let committed = false; + const failure = new OrchestrationDispatchCommandError({ + message: "database rejected deletion", + }); + const result = yield* f + .execute( + Effect.gen(function* () { + expect(yield* f.fs.exists(f.worktree)).toBe(false); + committed = true; + return yield* failure; + }), + ) + .pipe(Effect.flip); + expect(result).toBe(failure); + expect(committed).toBe(true); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "tracked"))).toBe("unstaged"); + expect((yield* f.run(["show", ":tracked"], f.worktree)).stdout).toBe("staged"); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "untracked"))).toBe( + "untracked contents", + ); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "ignored"))).toBe( + "ignored contents", + ); + expect( + (yield* f.run(["rev-parse", "--abbrev-ref", "HEAD"], f.worktree)).stdout.trim(), + ).toBe("test"); + expect((yield* f.execute(Effect.succeed({ sequence: 3 }))).sequence).toBe(3); + expect(yield* f.fs.exists(f.worktree)).toBe(false); + expect(yield* f.fs.readDirectory(f.root)).toEqual(["repo"]); + }), + ); + + it.effect("does not delete the thread or change files when the worktree is locked", () => + Effect.gen(function* () { + const f = yield* fixture; + yield* f.run(["worktree", "lock", f.worktree]); + let committed = false; + yield* f + .execute( + Effect.sync(() => { + committed = true; + return { sequence: 3 }; + }), + ) + .pipe(Effect.flip); + expect(committed).toBe(false); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "tracked"))).toBe("unstaged"); + expect(yield* f.fs.readDirectory(f.root)).toEqual(["repo", "worktree"]); + }), + ); + + it.effect("keeps recovery files when rollback is blocked, then restores them on retry", () => + Effect.gen(function* () { + const f = yield* fixture; + yield* f + .execute( + Effect.gen(function* () { + yield* f.fs.makeDirectory(f.worktree); + yield* f.fs.writeFileString(f.path.join(f.worktree, "replacement"), "do not overwrite"); + return yield* new OrchestrationDispatchCommandError({ message: "commit failed" }); + }), + ) + .pipe(Effect.flip); + const staged = (yield* f.fs.readDirectory(f.root)).find((name) => + name.startsWith(".t3-delete-"), + )!; + expect(yield* f.fs.readFileString(f.path.join(f.root, staged, "ignored"))).toBe( + "ignored contents", + ); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "replacement"))).toBe( + "do not overwrite", + ); + yield* f.fs.remove(f.worktree, { recursive: true }); + // Same on-disk state as a server exiting after staging and before committing. + yield* f + .execute(Effect.fail(new OrchestrationDispatchCommandError({ message: "still offline" }))) + .pipe(Effect.flip); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "ignored"))).toBe( + "ignored contents", + ); + expect((yield* f.run(["show", ":tracked"], f.worktree)).stdout).toBe("staged"); + }), + ); + + it.effect("keeps a worktree used by an archived thread", () => + Effect.gen(function* () { + const f = yield* fixture; + const thread = f.snapshot.threads[0]!; + f.setSnapshot({ + ...f.snapshot, + threads: [ + ...f.snapshot.threads, + { ...thread, id: ThreadId.make("archived"), archivedAt: thread.createdAt }, + ], + }); + yield* f.execute(Effect.succeed({ sequence: 3 })); + expect(yield* f.fs.exists(f.worktree)).toBe(true); + }), + ); + + it.effect("rolls back when another thread starts sharing the worktree during staging", () => + Effect.gen(function* () { + const f = yield* fixture; + const thread = f.snapshot.threads[0]!; + const result = yield* f + .execute( + Effect.sync(() => { + f.setSnapshot({ + ...f.snapshot, + threads: [...f.snapshot.threads, { ...thread, id: ThreadId.make("new-sharer") }], + }); + return { sequence: 3 }; + }), + ) + .pipe(Effect.flip); + expect(result.message).toBe("test commit failed"); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "ignored"))).toBe( + "ignored contents", + ); + expect((yield* f.run(["show", ":tracked"], f.worktree)).stdout).toBe("staged"); + }), + ); + + it.effect("rolls back when the target changes worktrees during staging", () => + Effect.gen(function* () { + const f = yield* fixture; + yield* f + .execute( + Effect.sync(() => { + f.setSnapshot({ + ...f.snapshot, + threads: f.snapshot.threads.map((thread) => ({ + ...thread, + worktreePath: f.path.join(f.root, "other"), + })), + }); + return { sequence: 3 }; + }), + ) + .pipe(Effect.flip); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "ignored"))).toBe( + "ignored contents", + ); + }), + ); + + it.effect("returns a cleanup retry path when final removal fails after the thread commits", () => + Effect.gen(function* () { + const f = yield* fixture; + const result = yield* f.execute( + Effect.gen(function* () { + const staged = (yield* f.fs.readDirectory(f.root)).find((name) => + name.startsWith(".t3-delete-"), + )!; + yield* f.run(["worktree", "lock", f.path.join(f.root, staged)]); + return { sequence: 3 }; + }), + ); + expect(result.sequence).toBe(3); + expect(result.worktreeCleanupPending?.cwd).toBe(f.cwd); + const staged = result.worktreeCleanupPending!.path; + expect(yield* f.fs.readFileString(f.path.join(staged, "ignored"))).toBe("ignored contents"); + yield* f.run(["worktree", "unlock", staged]); + const snapshot = f.getSnapshot(); + for (const archivedAt of [null, "2026-09-17T00:00:00.000Z"]) { + f.setSnapshot({ + ...snapshot, + threads: [ + ...snapshot.threads, + { + ...f.snapshot.threads[0]!, + id: ThreadId.make("new-user"), + worktreePath: staged, + archivedAt, + }, + ], + }); + const failure = yield* f.retryCleanup(staged).pipe(Effect.flip); + expect(failure.message).toContain("still used by a thread"); + expect(yield* f.fs.readFileString(f.path.join(staged, "ignored"))).toBe("ignored contents"); + expect((yield* f.run(["show", ":tracked"], staged)).stdout).toBe("staged"); + } + f.setSnapshot(snapshot); + yield* f.retryCleanup(staged); + expect(yield* f.fs.exists(staged)).toBe(false); + }), + ); + + it.effect("keeps files if cleanup retry cannot read current references", () => + Effect.gen(function* () { + const f = yield* fixture; + f.failSnapshotRead(); + const failure = yield* f.retryCleanup(f.worktree).pipe(Effect.flip); + expect(failure.message).toContain("Could not verify worktree references"); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "ignored"))).toBe( + "ignored contents", + ); + }), + ); + + it.effect("restores files if a deduplicated receipt did not delete the current thread", () => + Effect.gen(function* () { + const f = yield* fixture; + yield* withThreadWorktreeDeletion(f.command, () => Effect.succeed({ sequence: 1 })).pipe( + Effect.provide( + Layer.mock(ProjectionSnapshotQuery)({ + getCommandReadModel: () => Effect.succeed(f.snapshot), + }), + ), + Effect.flip, + ); + expect(yield* f.fs.readFileString(f.path.join(f.worktree, "ignored"))).toBe( + "ignored contents", + ); + expect((yield* f.run(["show", ":tracked"], f.worktree)).stdout).toBe("staged"); + expect(yield* f.fs.readDirectory(f.root)).toEqual(["repo", "worktree"]); + }), + ); + + it.effect("reports cleanup uncertainty as success after a committed deletion", () => + Effect.gen(function* () { + const f = yield* fixture; + const result = yield* f.execute( + Effect.succeed({ sequence: 3 }), + Effect.sync(f.failSnapshotRead), + ); + expect(result.sequence).toBe(3); + expect( + f + .getSnapshot() + .threads.some((thread) => thread.id === f.command.threadId && thread.deletedAt === null), + ).toBe(false); + expect(result.worktreeCleanupPending?.retryable).toBe(false); + expect( + yield* f.fs.readFileString(f.path.join(result.worktreeCleanupPending!.path, "ignored")), + ).toBe("ignored contents"); + }), + ); + + it.effect("keeps a staged worktree newly used by another thread without offering cleanup", () => + Effect.gen(function* () { + const f = yield* fixture; + const originalThread = f.snapshot.threads[0]!; + let survivorPath = ""; + const result = yield* f.execute( + Effect.succeed({ sequence: 3 }), + Effect.gen(function* () { + const directory = (yield* f.fs.readDirectory(f.root)).find((name) => + name.startsWith(".t3-delete-"), + )!; + survivorPath = f.path.join(f.root, directory); + const snapshot = f.getSnapshot(); + f.setSnapshot({ + ...snapshot, + threads: [ + ...snapshot.threads, + { + ...originalThread, + id: ThreadId.make("staged-survivor"), + worktreePath: survivorPath, + }, + ], + }); + }), + ); + expect(result.worktreeCleanupPending).toBeUndefined(); + expect(yield* f.fs.exists(f.worktree)).toBe(false); + expect(yield* f.fs.readFileString(f.path.join(survivorPath, "ignored"))).toBe( + "ignored contents", + ); + expect((yield* f.run(["show", ":tracked"], survivorPath)).stdout).toBe("staged"); + }), + ); +}); diff --git a/apps/server/src/orchestration/threadWorktreeDeletion.ts b/apps/server/src/orchestration/threadWorktreeDeletion.ts new file mode 100644 index 000000000000..509ac4f63baa --- /dev/null +++ b/apps/server/src/orchestration/threadWorktreeDeletion.ts @@ -0,0 +1,203 @@ +import * as NodeCrypto from "node:crypto"; +import { + GitCommandError, + OrchestrationDispatchCommandError, + type DispatchResult, + type OrchestrationCommand, + type VcsRemoveWorktreeInput, +} from "@t3tools/contracts"; +import * as Effect from "effect/Effect"; +import * as Exit from "effect/Exit"; +import * as FileSystem from "effect/FileSystem"; +import * as Path from "effect/Path"; +import * as Semaphore from "effect/Semaphore"; + +import { GitVcsDriver } from "../vcs/GitVcsDriver.ts"; +import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.ts"; +import { OrchestrationEngineService } from "./Services/OrchestrationEngine.ts"; + +const locks = new Map(); + +/** Cleanup retries may arrive long after deletion, when another thread uses the files. */ +export const removeUnusedWorktree = Effect.fn("removeUnusedWorktree")(function* ( + input: VcsRemoveWorktreeInput, + remove: Effect.Effect, +) { + const path = yield* Path.Path; + const snapshots = yield* ProjectionSnapshotQuery; + const error = (detail: string) => + new GitCommandError({ + operation: "removeWorktree", + command: "git worktree remove", + cwd: input.cwd, + detail, + }); + const engine = yield* OrchestrationEngineService; + return yield* engine.withWorktreeCleanup( + [input.path], + Effect.gen(function* () { + const snapshot = yield* snapshots + .getCommandReadModel() + .pipe( + Effect.mapError(() => error("Could not verify worktree references. Files were kept.")), + ); + const target = path.resolve(input.path); + if ( + snapshot.threads.some( + (thread) => + thread.deletedAt === null && + thread.worktreePath !== null && + path.resolve(thread.worktreePath) === target, + ) + ) { + return yield* error("This worktree is still used by a thread. Files were kept."); + } + yield* remove; + }), + ); +}); + +/** Keep all files (including ignored files and the Git index) until deletion commits. + * A stable staging path also lets a retry restore a worktree after a server crash. + */ +export const withThreadWorktreeDeletion = Effect.fn("withThreadWorktreeDeletion")(function* ( + command: Extract, + commit: ( + worktreeStaged: boolean, + ) => Effect.Effect, +) { + if (!command.deleteWorktreePath) return yield* commit(false); + const fs = yield* FileSystem.FileSystem; + const path = yield* Path.Path; + const git = yield* GitVcsDriver; + const snapshots = yield* ProjectionSnapshotQuery; + const original = path.resolve(command.deleteWorktreePath); + const key = NodeCrypto.createHash("sha256").update(original).digest("hex").slice(0, 24); + const staged = path.join(path.dirname(original), `.t3-delete-${key}`); + const lock = locks.get(key) ?? { semaphore: Semaphore.makeUnsafe(1), users: 0 }; + lock.users++; + locks.set(key, lock); + return yield* Effect.gen(function* () { + const snapshot = yield* snapshots.getCommandReadModel(); + const thread = snapshot.threads.find( + (entry) => entry.id === command.threadId && entry.deletedAt === null, + ); + if (!thread || thread.worktreePath === null || path.resolve(thread.worktreePath) !== original) { + return yield* new OrchestrationDispatchCommandError({ + message: "The thread's worktree changed. Refresh and try deleting it again.", + }); + } + const project = snapshot.projects.find((entry) => entry.id === thread.projectId); + if (!project) + return yield* new OrchestrationDispatchCommandError({ + message: "The thread's project could not be found.", + }); + const cwd = project.workspaceRoot; + const move = (from: string, to: string) => + git.execute({ + operation: "thread.delete.move-worktree", + cwd, + args: ["worktree", "move", "--", from, to], + timeoutMs: 300_000, + }); + const restore = Effect.gen(function* () { + if (!(yield* fs.exists(staged))) return; + if (yield* fs.exists(original)) + return yield* new OrchestrationDispatchCommandError({ + message: `Worktree recovery could not replace ${original}. Your files are preserved at ${staged}.`, + }); + // A process can stop between the directory move and Git updating its pointers. + yield* git.execute({ + operation: "thread.delete.repair-worktree", + cwd, + args: ["worktree", "repair", "--", staged], + }); + yield* move(staged, original); + }); + yield* restore; + if ( + snapshot.threads.some( + (entry) => + entry.id !== thread.id && + entry.deletedAt === null && + entry.worktreePath !== null && + path.resolve(entry.worktreePath) === original, + ) + ) { + return yield* commit(false); + } + // Already removed externally: the ordinary delete remains safe and retryable. + if (!(yield* fs.exists(original))) return yield* commit(false); + const result = yield* Effect.gen(function* () { + yield* move(original, staged); + return yield* commit(true); + }).pipe(Effect.exit); + if (Exit.isFailure(result)) { + const recovery = yield* restore.pipe(Effect.exit); + if (Exit.isFailure(recovery)) + return yield* new OrchestrationDispatchCommandError({ + message: `Thread deletion failed and automatic worktree recovery failed. Your files are preserved at ${staged}; restore them before retrying.`, + cause: recovery.cause, + }); + return yield* Effect.failCause(result.cause); + } + const engine = yield* OrchestrationEngineService; + return yield* engine.withWorktreeCleanup( + [original, staged], + Effect.gen(function* () { + // Dispatch can replay a receipt for an earlier incarnation. Only remove files + // after an authoritative read confirms the current thread and references are gone. + const committed = yield* snapshots.getCommandReadModel().pipe(Effect.exit); + const pendingCleanup = (retryable: boolean): DispatchResult => ({ + ...result.value, + worktreeCleanupPending: { cwd, path: staged, retryable }, + }); + if (Exit.isFailure(committed)) return pendingCleanup(false); + const survivors = committed.value.threads.filter((entry) => entry.deletedAt === null); + const stagedInUse = survivors.some( + (entry) => entry.worktreePath !== null && path.resolve(entry.worktreePath) === staged, + ); + if (survivors.some((entry) => entry.id === command.threadId)) { + if (!stagedInUse) yield* restore; + return yield* new OrchestrationDispatchCommandError({ + message: + "The current thread was not deleted. Its files were kept; refresh before retrying.", + }); + } + // A new thread can have selected the staged checkout from Git's worktree list. + // Keep its path intact, without offering destructive cleanup for shared files. + if (stagedInUse) return result.value; + if ( + survivors.some( + (entry) => entry.worktreePath !== null && path.resolve(entry.worktreePath) === original, + ) + ) { + const recovery = yield* restore.pipe(Effect.exit); + return Exit.isFailure(recovery) ? pendingCleanup(false) : result.value; + } + // The thread is committed as deleted. A cleanup error must not masquerade as + // a failed deletion; return the preserved staging path for an explicit retry. + const cleanup = yield* git + .removeWorktree({ cwd, path: staged, force: true }) + .pipe(Effect.exit); + if (Exit.isFailure(cleanup)) { + yield* Effect.logWarning("Deleted thread has pending worktree cleanup", { + threadId: command.threadId, + cwd, + path: staged, + }); + return pendingCleanup(true); + } + return result.value; + }), + ); + }).pipe( + Effect.uninterruptible, + lock.semaphore.withPermits(1), + Effect.ensuring( + Effect.sync(() => { + if (--lock.users === 0) locks.delete(key); + }), + ), + ); +}); diff --git a/apps/server/src/project/AgentSessionImporter.test.ts b/apps/server/src/project/AgentSessionImporter.test.ts index 2438edca8b1b..2309416e4192 100644 --- a/apps/server/src/project/AgentSessionImporter.test.ts +++ b/apps/server/src/project/AgentSessionImporter.test.ts @@ -227,6 +227,8 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { getThreadReplayStats: () => Effect.die("unused"), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }); const directory = ProviderSessionDirectory.ProviderSessionDirectory.of({ @@ -332,6 +334,8 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { getThreadReplayStats: () => Effect.die("unused"), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }); const directory = ProviderSessionDirectory.ProviderSessionDirectory.of({ @@ -397,6 +401,8 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { getThreadReplayStats: () => Effect.die("unused"), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }); const directory = ProviderSessionDirectory.ProviderSessionDirectory.of({ @@ -468,6 +474,8 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { getThreadReplayStats: () => Effect.die("unused"), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }); @@ -506,6 +514,8 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => { getThreadReplayStats: () => Effect.die("unused"), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }); const directory = ProviderSessionDirectory.ProviderSessionDirectory.of({ diff --git a/apps/server/src/relay/AgentAwarenessRelay.test.ts b/apps/server/src/relay/AgentAwarenessRelay.test.ts index 71ec0719325d..e4f640288b41 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -558,6 +558,8 @@ describe.sequential("signRelayAgentActivityPublishProof", () => { dispatch: () => Effect.succeed({ sequence: 1 }), streamDomainEvents: Stream.fromQueue(events), subscribeDomainEvents: Effect.succeed(Stream.fromQueue(events)), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), } satisfies OrchestrationEngineShape; @@ -784,6 +786,10 @@ describe.sequential("signRelayAgentActivityPublishProof", () => { dispatch: () => Effect.succeed({ sequence: 1 }), streamDomainEvents: Stream.fromQueue(events), subscribeDomainEvents: Effect.succeed(Stream.fromQueue(events)), + withWorktreeCleanup: ( + _paths: ReadonlyArray, + effect: Effect.Effect, + ) => effect, latestSequence: Effect.succeed(0), } satisfies OrchestrationEngineShape), Layer.succeed(ProjectionSnapshotQuery, { diff --git a/apps/server/src/server.test.ts b/apps/server/src/server.test.ts index cb1ed672312e..52fa0bbd2ed5 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -1006,6 +1006,7 @@ const buildAppUnderTest = (options?: { }), dispatch: () => Effect.succeed({ sequence: 0 }), streamDomainEvents: Stream.empty, + withWorktreeCleanup: (_paths, effect) => effect, latestSequence: Effect.succeed(0), ...options?.layers?.orchestrationEngine, }), diff --git a/apps/server/src/serverRuntimeStartup.reconcile.test.ts b/apps/server/src/serverRuntimeStartup.reconcile.test.ts index 050650b11515..6a29c09d7ed9 100644 --- a/apps/server/src/serverRuntimeStartup.reconcile.test.ts +++ b/apps/server/src/serverRuntimeStartup.reconcile.test.ts @@ -103,6 +103,8 @@ const runReconciliation = (input: { dispatch: input.dispatch, streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }), Effect.provide( @@ -715,6 +717,8 @@ it.effect("does not fail startup when the live provider session inventory cannot dispatch: () => Effect.die("unused"), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }), Effect.provide(Layer.mergeAll(NodeServices.layer, ServerSettings.layerTest())), diff --git a/apps/server/src/serverRuntimeStartup.test.ts b/apps/server/src/serverRuntimeStartup.test.ts index 4bbdb2e9c8b5..3000d27d7173 100644 --- a/apps/server/src/serverRuntimeStartup.test.ts +++ b/apps/server/src/serverRuntimeStartup.test.ts @@ -211,6 +211,8 @@ it.effect("resolveAutoBootstrapWelcomeTargets returns existing project and threa ), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), } satisfies OrchestrationEngine.OrchestrationEngineService["Service"]), Effect.provide(NodeServices.layer), @@ -342,6 +344,8 @@ it.effect.each([ ), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), } satisfies OrchestrationEngine.OrchestrationEngineService["Service"]), Effect.provide(NodeServices.layer), @@ -416,6 +420,8 @@ it.effect( ), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), } satisfies OrchestrationEngine.OrchestrationEngineService["Service"]), Effect.provide(NodeServices.layer), @@ -482,6 +488,8 @@ it.effect("resolveAutoBootstrapWelcomeTargets preserves typed UUID generation fa ), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), } satisfies OrchestrationEngine.OrchestrationEngineService["Service"]), Effect.provideService(Crypto.Crypto, { diff --git a/apps/server/src/serverRuntimeStartup.worktreeSetup.test.ts b/apps/server/src/serverRuntimeStartup.worktreeSetup.test.ts index fda60d889716..f4d1ea3169df 100644 --- a/apps/server/src/serverRuntimeStartup.worktreeSetup.test.ts +++ b/apps/server/src/serverRuntimeStartup.worktreeSetup.test.ts @@ -97,6 +97,8 @@ const run = (activities: ReadonlyArray>) => }), streamDomainEvents: Stream.empty, subscribeDomainEvents: Effect.succeed(Stream.empty), + withWorktreeCleanup: (_paths: ReadonlyArray, effect: Effect.Effect) => + effect, latestSequence: Effect.succeed(0), }), Effect.provide(NodeServices.layer), diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 47974366a671..1f4c9577aff1 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -105,6 +105,10 @@ import { import * as OrchestrationEngine from "./orchestration/Services/OrchestrationEngine.ts"; import * as ProjectionSnapshotQuery from "./orchestration/Services/ProjectionSnapshotQuery.ts"; import { ThreadDeletionReactor } from "./orchestration/Services/ThreadDeletionReactor.ts"; +import { + removeUnusedWorktree, + withThreadWorktreeDeletion, +} from "./orchestration/threadWorktreeDeletion.ts"; import { observeRpcEffect as instrumentRpcEffect, observeRpcStream as instrumentRpcStream, @@ -555,6 +559,9 @@ const makeWsRpcLayer = ( const externalLauncher = yield* ExternalLauncher.ExternalLauncher; const remoteOpenTargets = yield* RemoteOpenTargets.RemoteOpenTargets; const gitWorkflow = yield* GitWorkflowService.GitWorkflowService; + const worktreeDeletionContext = yield* Effect.context< + GitVcsDriver.GitVcsDriver | FileSystem.FileSystem | Path.Path + >(); const review = yield* ReviewService.ReviewService; const vcsProvisioning = yield* VcsProvisioningService.VcsProvisioningService; const vcsStatusBroadcaster = yield* VcsStatusBroadcaster.VcsStatusBroadcaster; @@ -1907,7 +1914,31 @@ const makeWsRpcLayer = ( ), ) : false; - const result = yield* dispatchNormalizedCommand(normalizedCommand).pipe( + const dispatch = dispatchNormalizedCommand(normalizedCommand); + const result = yield* ( + normalizedCommand.type === "thread.delete" && normalizedCommand.deleteWorktreePath + ? withThreadWorktreeDeletion(normalizedCommand, (staged) => + dispatchNormalizedCommand( + staged + ? normalizedCommand + : { ...normalizedCommand, deleteWorktreePath: undefined }, + ), + ).pipe( + Effect.provide(worktreeDeletionContext), + Effect.provideService( + OrchestrationEngine.OrchestrationEngineService, + orchestrationEngine, + ), + Effect.provideService( + ProjectionSnapshotQuery.ProjectionSnapshotQuery, + projectionSnapshotQuery, + ), + Effect.mapError((cause) => + toDispatchCommandError(cause, "Failed to delete thread and worktree"), + ), + ) + : dispatch + ).pipe( Effect.tapError(() => cleanupFailedUploadedAttachments(command, normalizedCommand)), ); yield* recordClientCommandAnalytics(normalizedCommand); @@ -3396,7 +3427,18 @@ const makeWsRpcLayer = ( [WS_METHODS.vcsRemoveWorktree]: (input) => observeRpcEffect( WS_METHODS.vcsRemoveWorktree, - gitWorkflow.removeWorktree(input).pipe(Effect.tap(() => refreshGitStatus(input.cwd))), + removeUnusedWorktree(input, gitWorkflow.removeWorktree(input)).pipe( + Effect.provide(worktreeDeletionContext), + Effect.provideService( + OrchestrationEngine.OrchestrationEngineService, + orchestrationEngine, + ), + Effect.provideService( + ProjectionSnapshotQuery.ProjectionSnapshotQuery, + projectionSnapshotQuery, + ), + Effect.tap(() => refreshGitStatus(input.cwd)), + ), { "rpc.aggregate": "vcs" }, ), [WS_METHODS.vcsCreateRef]: (input) => diff --git a/apps/web/src/components/LegacySidebar.tsx b/apps/web/src/components/LegacySidebar.tsx index 1991dc502283..3ea3525acfd0 100644 --- a/apps/web/src/components/LegacySidebar.tsx +++ b/apps/web/src/components/LegacySidebar.tsx @@ -1964,34 +1964,48 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec const confirmed = await api.dialogs.confirm( [ `Delete ${count} thread${count === 1 ? "" : "s"}?`, - "This permanently clears conversation history for these threads.", + "This permanently clears conversation history and deletes worktrees no other threads use.", ].join("\n"), { variant: "destructive" }, ); if (!confirmed) return; } - const { deletedThreadKeys, firstFailure } = await deleteSelectedThreadEntries({ - entries: selectedThreadEntries, - delete: ({ threadRef }, deletedThreadKeys) => - deleteThread(threadRef, { deletedThreadKeys }), + const deletionToast = toastManager.add({ + type: "loading", + title: `Deleting ${count} thread${count === 1 ? "" : "s"}…`, + timeout: 0, }); - if (firstFailure !== null) { - const firstError = squashAtomCommandFailure(firstFailure); - toastManager.add( - stackedThreadToast({ - type: "error", - title: "Failed to delete threads", - description: firstError instanceof Error ? firstError.message : "An error occurred.", + try { + const { deletedThreadKeys, firstFailure } = await deleteSelectedThreadEntries({ + entries: selectedThreadEntries, + delete: ({ threadRef }, deletedThreadKeys, deferDeletion) => + deleteThread(threadRef, { + deletedThreadKeys, + deferDeletion, + worktreeDeletionConfirmed: true, + }), + }); + if (firstFailure !== null) { + const firstError = squashAtomCommandFailure(firstFailure); + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Some threads could not be deleted", + description: `${firstError instanceof Error ? firstError.message : "An error occurred."} Remaining threads were kept; you can retry deleting them.`, + timeout: 0, + }), + ); + } + removeFromSelection( + getThreadKeysToDeselectAfterDelete(threadKeys, deletedThreadKeys, (threadKey) => { + const threadRef = parseScopedThreadKey(threadKey); + return threadRef !== null && readThreadShell(threadRef) !== null; }), ); + } finally { + toastManager.close(deletionToast); } - removeFromSelection( - getThreadKeysToDeselectAfterDelete(threadKeys, deletedThreadKeys, (threadKey) => { - const threadRef = parseScopedThreadKey(threadKey); - return threadRef !== null && readThreadShell(threadRef) !== null; - }), - ); }, [ appSettingsConfirmThreadArchive, diff --git a/apps/web/src/components/Sidebar.logic.test.ts b/apps/web/src/components/Sidebar.logic.test.ts index fede3c183448..e510aa4c7d36 100644 --- a/apps/web/src/components/Sidebar.logic.test.ts +++ b/apps/web/src/components/Sidebar.logic.test.ts @@ -180,6 +180,49 @@ describe("deleteSelectedThreadEntries", () => { }); }); + it.each([false, true])( + "reports rejected deletions and waits for other entries (deferred=%s)", + async (deferred) => { + const error = new Error("Connection lost"); + let finish!: () => void; + let started!: () => void; + const pending = new Promise((resolve) => { + finish = resolve; + }); + const running = new Promise((resolve) => { + started = resolve; + }); + const deletion = deleteSelectedThreadEntries({ + entries, + delete: async ({ threadKey }, _deletedThreadKeys, deferDeletion) => { + const run = async () => { + if (threadKey === "one") throw error; + if (threadKey === "two") { + started(); + await pending; + } + return success; + }; + if (!deferred) return run(); + deferDeletion(run); + return success; + }, + }); + await running; + let completed = false; + void deletion.then(() => { + completed = true; + }); + await Promise.resolve(); + expect(completed).toBe(false); + finish(); + const outcome = await deletion; + expect(outcome.deletedThreadKeys).toEqual(new Set(["two", "three"])); + expect(outcome.firstFailure?._tag).toBe("Failure"); + if (outcome.firstFailure) expect(Cause.squash(outcome.firstFailure.cause)).toBe(error); + }, + ); + it.each([ { firstResult: success, deletedThreadKeys: new Set(["one"]), firstFailure: null }, { firstResult: failure, deletedThreadKeys: new Set(), firstFailure: failure }, diff --git a/apps/web/src/components/Sidebar.logic.ts b/apps/web/src/components/Sidebar.logic.ts index 2796b2f7f885..2224b355a70a 100644 --- a/apps/web/src/components/Sidebar.logic.ts +++ b/apps/web/src/components/Sidebar.logic.ts @@ -3,6 +3,7 @@ import * as React from "react"; import { defaultAnimateLayoutChanges, type AnimateLayoutChanges } from "@dnd-kit/sortable"; import { isAtomCommandInterrupted, + settlePromise, type AtomCommandResult, } from "@t3tools/client-runtime/state/runtime"; import { threadSearchMatchKey } from "@t3tools/client-runtime/state/thread-search"; @@ -422,14 +423,29 @@ export async function deleteSelectedThreadEntries< delete: ( entry: TEntry, deletedThreadKeys: ReadonlySet, + deferDeletion: (deleteThread: () => Promise>) => void, ) => Promise | null>; }) { const deletedThreadKeys = new Set(); + const pendingDeletions: Array> = []; let firstFailure: AsyncResult.Failure | null = null; for (const entry of input.entries) { - const result = await input.delete(entry, deletedThreadKeys); - if (result === null) continue; + let deferred = false; + const attempt = await settlePromise(() => + input.delete(entry, deletedThreadKeys, (deleteThread) => { + deferred = true; + pendingDeletions.push( + settlePromise(deleteThread).then((attempt) => { + const result = attempt._tag === "Failure" ? attempt : attempt.value; + if (result._tag === "Success") deletedThreadKeys.add(entry.threadKey); + else if (!isAtomCommandInterrupted(result)) firstFailure ??= result; + }), + ); + }), + ); + const result = attempt._tag === "Failure" ? attempt : attempt.value; + if (result === null || deferred) continue; if (result._tag === "Failure") { if (isAtomCommandInterrupted(result)) break; firstFailure ??= result; @@ -438,6 +454,8 @@ export async function deleteSelectedThreadEntries< deletedThreadKeys.add(entry.threadKey); } + // Wait for every outcome, including rejected promises, before updating selection. + await Promise.all(pendingDeletions); return { deletedThreadKeys, firstFailure }; } diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index 7b9e0e545d56..b29fc64d0989 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -3986,39 +3986,51 @@ export default function Sidebar() { api.dialogs.confirm( [ `Delete ${count} thread${count === 1 ? "" : "s"}?`, - "This permanently clears conversation history for these threads.", + "This permanently clears conversation history and deletes worktrees no other threads use.", ].join("\n"), { variant: "destructive" }, ), ); if (confirmed._tag === "Failure" || !confirmed.value) return; } - const { deletedThreadKeys, firstFailure } = await deleteSelectedThreadEntries({ - entries: threadKeys.map((threadKey) => ({ threadKey })), - delete: async ({ threadKey }, deletedThreadKeys) => { - const thread = threadByKeyRef.current.get(threadKey); - if (!thread) return null; - return deleteThread(scopeThreadRef(thread.environmentId, thread.id), { - deletedThreadKeys, - }); - }, + const deletionToast = toastManager.add({ + type: "loading", + title: `Deleting ${count} thread${count === 1 ? "" : "s"}…`, + timeout: 0, }); - if (firstFailure !== null) { - const firstError = squashAtomCommandFailure(firstFailure); - toastManager.add( - stackedThreadToast({ - type: "error", - title: "Failed to delete threads", - description: firstError instanceof Error ? firstError.message : "An error occurred.", + try { + const { deletedThreadKeys, firstFailure } = await deleteSelectedThreadEntries({ + entries: threadKeys.map((threadKey) => ({ threadKey })), + delete: async ({ threadKey }, deletedThreadKeys, deferDeletion) => { + const thread = threadByKeyRef.current.get(threadKey); + if (!thread) return null; + return deleteThread(scopeThreadRef(thread.environmentId, thread.id), { + deletedThreadKeys, + deferDeletion, + worktreeDeletionConfirmed: true, + }); + }, + }); + if (firstFailure !== null) { + const firstError = squashAtomCommandFailure(firstFailure); + toastManager.add( + stackedThreadToast({ + type: "error", + title: "Some threads could not be deleted", + description: `${firstError instanceof Error ? firstError.message : "An error occurred."} Remaining threads were kept; you can retry deleting them.`, + timeout: 0, + }), + ); + } + removeFromSelection( + getThreadKeysToDeselectAfterDelete(selectedThreadKeys, deletedThreadKeys, (threadKey) => { + const threadRef = parseScopedThreadKey(threadKey); + return threadRef !== null && readThreadShell(threadRef) !== null; }), ); + } finally { + toastManager.close(deletionToast); } - removeFromSelection( - getThreadKeysToDeselectAfterDelete(selectedThreadKeys, deletedThreadKeys, (threadKey) => { - const threadRef = parseScopedThreadKey(threadKey); - return threadRef !== null && readThreadShell(threadRef) !== null; - }), - ); }, [ attemptSettle, diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index fc5082fcf1a4..ba185e8e04b7 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -3052,7 +3052,7 @@ export function GeneralSettingsPanel() { ({ + confirm: vi.fn(async () => true), + run: vi.fn(), + readThreadShell: vi.fn(), + readProject: vi.fn(), + readEnvironmentThreadRefs: vi.fn(), + confirmThreadDelete: false, + recoverableDeletion: true, + automaticCleanup: false, + archived: vi.fn(), + toastAdd: vi.fn( + (_toast: { + title: string; + description: string; + actionProps?: { onClick: () => Promise }; + }) => "cleanup-toast", + ), + toastClose: vi.fn(), + toastUpdate: vi.fn(), +})); +vi.mock("../state/server", async (importOriginal) => { + const { Atom } = await import("effect/unstable/reactivity"); + const { DEFAULT_SERVER_SETTINGS, EnvironmentId } = await import("@t3tools/contracts"); + return { + ...(await importOriginal()), + environmentServerConfigsAtom: Atom.make( + () => + new Map([ + [ + EnvironmentId.make("local"), + { + settings: { + ...DEFAULT_SERVER_SETTINGS, + storageCleanup: { + ...DEFAULT_SERVER_SETTINGS.storageCleanup, + worktreeOnDelete: mocks.automaticCleanup, + }, + }, + }, + ], + ]), + ), + }; +}); +vi.mock("@t3tools/client-runtime/state/runtime", async (importOriginal) => ({ + ...(await importOriginal()), + executeAtomQuery: mocks.archived, +})); +vi.mock("../state/use-atom-command", () => ({ + useAtomCommand: (command: { label: string }) => (input: unknown) => + mocks.run(command.label, input), +})); +vi.mock("../state/entities", () => ({ + readThreadShell: mocks.readThreadShell, + readEnvironmentThreadRefs: mocks.readEnvironmentThreadRefs, + readProject: mocks.readProject, + readEnvironmentSupportsRecoverableDeletion: () => mocks.recoverableDeletion, +})); +vi.mock("./useSettings", () => ({ + useClientSettings: (select: (settings: object) => unknown) => + select({ confirmThreadDelete: mocks.confirmThreadDelete, sidebarThreadSortOrder: "updatedAt" }), +})); +vi.mock("@tanstack/react-router", () => ({ + useRouter: () => ({ state: { matches: [] } }), +})); +vi.mock("./useHandleNewThread", () => ({ useNewThreadHandler: () => vi.fn() })); +vi.mock("../localApi", () => ({ readLocalApi: () => ({ dialogs: { confirm: mocks.confirm } }) })); +vi.mock("../composerDraftStore", () => ({ useComposerDraftStore: () => vi.fn() })); +vi.mock("../terminalUiStateStore", () => ({ useTerminalUiStateStore: () => vi.fn() })); +vi.mock("../uiStateStore", () => ({ useUiStateStore: () => vi.fn() })); +vi.mock("../lib/composerDraftUploads", () => ({ releaseComposerDraftUploads: vi.fn() })); +vi.mock("../lib/archivedThreadsState", () => ({ refreshArchivedThreadsForEnvironment: vi.fn() })); +vi.mock("../components/ui/toast", () => ({ + stackedThreadToast: (value: unknown) => value, + toastManager: { add: mocks.toastAdd, close: mocks.toastClose, update: mocks.toastUpdate }, +})); + +const environmentId = EnvironmentId.make("local"); +const threads = ["one", "two", "three"].map((id) => ({ + id: ThreadId.make(id), + environmentId, + projectId: ProjectId.make(`project-${id}`), + title: id, + session: null, + worktreePath: `/repo/${id}`, + createdAt: "2026-09-17T00:00:00.000Z", + updatedAt: "2026-09-17T00:00:00.000Z", +})); +const entries = threads.map((thread) => { + const threadRef = scopeThreadRef(environmentId, thread.id); + return { threadRef, threadKey: scopedThreadKey(threadRef) }; +}); +let actions: ReturnType; +let renderer: ReactTestRenderer; +function Probe() { + const value = useThreadActions(); + useLayoutEffect(() => { + actions = value; + }); + return null; +} + +beforeEach(() => { + vi.clearAllMocks(); + vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); + mocks.confirmThreadDelete = false; + mocks.recoverableDeletion = true; + mocks.automaticCleanup = false; + appAtomRegistry.refresh(environmentServerConfigsAtom); + mocks.readProject.mockReturnValue({ workspaceRoot: "/repo" }); + mocks.confirm.mockResolvedValue(true); + mocks.run.mockResolvedValue(AsyncResult.success({ sequence: 1 })); + mocks.archived.mockResolvedValue(AsyncResult.success({ threads: [] })); + // Keep the snapshot stale to exercise successful-deletion tracking. + mocks.readThreadShell.mockImplementation( + ({ threadId }) => threads.find((thread) => thread.id === threadId) ?? null, + ); + mocks.readEnvironmentThreadRefs.mockReturnValue(entries.map((entry) => entry.threadRef)); + act(() => { + renderer = create(); + }); +}); +afterEach(() => { + act(() => renderer.unmount()); + vi.unstubAllGlobals(); +}); + +it("keeps conversations until worktree cleanup succeeds without blocking the rest of the batch", async () => { + mocks.readProject.mockImplementation(({ projectId }) => ({ workspaceRoot: `/${projectId}` })); + let startCleanup!: () => void; + const cleanupStarted = new Promise((resolve) => { + startCleanup = resolve; + }); + let finishCleanup!: () => void; + const releaseCleanup = new Promise((resolve) => { + finishCleanup = resolve; + }); + const deleted: string[] = []; + const removals: string[] = []; + mocks.run.mockImplementation(async (label, { input }) => { + if (label.endsWith(":thread:delete") && input.deleteWorktreePath) { + removals.push(input.deleteWorktreePath); + if (removals.length === 3) startCleanup(); + await releaseCleanup; + } + if (label.endsWith(":thread:delete")) deleted.push(input.threadId); + return AsyncResult.success({ sequence: 1 }); + }); + const deletion = deleteSelectedThreadEntries({ + entries, + delete: ({ threadRef }, deletedThreadKeys, deferDeletion) => + actions.deleteThread(threadRef, { deletedThreadKeys, deferDeletion }), + }); + await cleanupStarted; + expect(deleted).toEqual([]); + expect(removals).toEqual(["/repo/one", "/repo/two", "/repo/three"]); + expect(mocks.confirm).not.toHaveBeenCalled(); + finishCleanup(); + expect((await deletion).deletedThreadKeys.size).toBe(3); +}); + +it.each([false, true])( + "removes a shared worktree only if all its threads were deleted (failure=%s)", + async (failFirst) => { + mocks.readThreadShell.mockImplementation(({ threadId }) => ({ + ...threads.find((thread) => thread.id === threadId), + worktreePath: "/repo/shared", + })); + mocks.run.mockImplementation(async (label, { input }) => + failFirst && label.endsWith(":thread:delete") && input.threadId === "one" + ? AsyncResult.failure(Cause.fail(new Error("delete failed"))) + : AsyncResult.success({ sequence: 1 }), + ); + await deleteSelectedThreadEntries({ + entries, + delete: ({ threadRef }, deletedThreadKeys, deferDeletion) => + actions.deleteThread(threadRef, { deletedThreadKeys, deferDeletion }), + }); + expect( + mocks.run.mock.calls.filter( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toHaveLength(failFirst ? 0 : 1); + }, +); + +it("does not repeat a worktree confirmation already included in the bulk confirmation", async () => { + mocks.confirmThreadDelete = true; + act(() => renderer.update()); + await actions.deleteThread(entries[0]!.threadRef, { worktreeDeletionConfirmed: true }); + expect(mocks.confirm).not.toHaveBeenCalled(); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(true); +}); + +it("still offers to keep a worktree for single deletion when confirmations are on", async () => { + mocks.confirmThreadDelete = true; + mocks.confirm.mockResolvedValue(false); + act(() => renderer.update()); + await actions.deleteThread(entries[0]!.threadRef); + expect(mocks.confirm).toHaveBeenCalledOnce(); + expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":thread:delete"))).toBe(true); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); +}); + +it.each([false, true])( + "respects automatic cleanup with confirmations enabled=%s", + async (confirmations) => { + mocks.automaticCleanup = true; + mocks.confirmThreadDelete = confirmations; + appAtomRegistry.refresh(environmentServerConfigsAtom); + act(() => renderer.update()); + expect((await actions.deleteThread(entries[0]!.threadRef))._tag).toBe("Success"); + expect(mocks.confirm).not.toHaveBeenCalled(); + expect(mocks.run.mock.calls.find(([label]) => label.endsWith(":thread:delete"))?.[1]).toEqual({ + environmentId, + input: { + threadId: threads[0]!.id, + ...(!confirmations ? { deleteWorktreePath: "/repo/one" } : {}), + }, + }); + }, +); + +it("keeps a worktree still used by an archived thread", async () => { + mocks.archived.mockResolvedValue( + AsyncResult.success({ + threads: [{ id: ThreadId.make("archived"), worktreePath: "/repo/one" }], + }), + ); + await actions.deleteThread(entries[0]!.threadRef); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); +}); + +it.each([false, true])( + "uses recoverable deletion for an archived target (failure=%s)", + async (fail) => { + mocks.readThreadShell.mockReturnValue(null); + mocks.archived.mockResolvedValue(AsyncResult.success({ threads: [threads[0]] })); + const deletion = fail + ? AsyncResult.failure(Cause.fail(new Error("worktree locked"))) + : AsyncResult.success({ sequence: 1 }); + mocks.run.mockImplementation(async (label) => + label.endsWith(":thread:delete") ? deletion : AsyncResult.success({ sequence: 1 }), + ); + + expect(await actions.deleteThread(entries[0]!.threadRef)).toBe(deletion); + expect(mocks.run.mock.calls.find(([label]) => label.endsWith(":thread:delete"))?.[1]).toEqual({ + environmentId, + input: { threadId: threads[0]!.id, deleteWorktreePath: "/repo/one" }, + }); + expect(mocks.confirm).not.toHaveBeenCalled(); + }, +); + +it("keeps an archived target's worktree when another archived thread shares it", async () => { + mocks.readThreadShell.mockReturnValue(null); + mocks.archived.mockResolvedValue( + AsyncResult.success({ + threads: [threads[0], { ...threads[1], worktreePath: "/repo/one" }], + }), + ); + await actions.deleteThread(entries[0]!.threadRef); + expect(mocks.run.mock.calls.find(([label]) => label.endsWith(":thread:delete"))?.[1]).toEqual({ + environmentId, + input: { threadId: threads[0]!.id }, + }); +}); + +it("does not delete an unresolved target when its archived shell cannot be loaded", async () => { + mocks.readThreadShell.mockReturnValue(null); + mocks.archived.mockResolvedValue(AsyncResult.failure(Cause.fail(new Error("offline")))); + expect((await actions.deleteThread(entries[0]!.threadRef))._tag).toBe("Failure"); + expect(mocks.run).not.toHaveBeenCalled(); +}); + +it("keeps the conversation if archived threads cannot be checked", async () => { + mocks.archived.mockResolvedValue(AsyncResult.failure(Cause.fail(new Error("offline")))); + const result = await actions.deleteThread(entries[0]!.threadRef); + expect(result._tag).toBe("Failure"); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); + expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":thread:delete"))).toBe(false); +}); + +it("rechecks live references before deferred cleanup", async () => { + let cleanup: (() => Promise>) | undefined; + await actions.deleteThread(entries[0]!.threadRef, { + deferDeletion: (run) => { + cleanup = run; + }, + }); + mocks.readThreadShell.mockImplementation(({ threadId }) => ({ + ...threads.find((thread) => thread.id === threadId), + worktreePath: "/repo/one", + })); + await cleanup!(); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); +}); + +it("keeps a failed worktree's thread, finishes other deletions, and allows retry", async () => { + const failure = AsyncResult.failure(Cause.fail(new Error("worktree is locked"))); + mocks.run.mockImplementation(async (label, { input }) => + label.endsWith(":thread:delete") && + input.deleteWorktreePath && + input.deleteWorktreePath === "/repo/two" + ? failure + : AsyncResult.success({ sequence: 1 }), + ); + const result = await deleteSelectedThreadEntries({ + entries, + delete: ({ threadRef }, deletedThreadKeys, deferDeletion) => + actions.deleteThread(threadRef, { deletedThreadKeys, deferDeletion }), + }); + expect(result.firstFailure).toBe(failure); + expect(result.deletedThreadKeys).toEqual(new Set([entries[0]!.threadKey, entries[2]!.threadKey])); + expect( + mocks.run.mock.calls + .filter(([label]) => label.endsWith(":thread:delete")) + .map(([, { input }]) => input.threadId), + ).toEqual(["one", "two", "three"]); + expect( + mocks.run.mock.calls + .filter(([label]) => label.includes("terminal") && label.endsWith(":close")) + .every(([, { input }]) => input.deleteHistory === false), + ).toBe(true); + mocks.run.mockResolvedValue(AsyncResult.success({ sequence: 1 })); + expect((await actions.deleteThread(entries[1]!.threadRef))._tag).toBe("Success"); + expect( + mocks.run.mock.calls + .filter(([label]) => label.endsWith(":thread:delete")) + .map(([, { input }]) => input.threadId), + ).toEqual(["one", "two", "three", "two"]); +}); + +it("shares an in-flight deletion when the same thread is deleted again", async () => { + let finishCleanup!: () => void; + const pending = new Promise((resolve) => { + finishCleanup = resolve; + }); + let startCleanup!: () => void; + const started = new Promise((resolve) => { + startCleanup = resolve; + }); + mocks.run.mockImplementation(async (label, { input }) => { + if (label.endsWith(":thread:delete") && input.deleteWorktreePath) { + startCleanup(); + await pending; + } + return AsyncResult.success({ sequence: 1 }); + }); + const first = actions.deleteThread(entries[0]!.threadRef); + await started; + const second = actions.deleteThread(entries[0]!.threadRef); + finishCleanup(); + await Promise.all([first, second]); + expect( + mocks.run.mock.calls.filter( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toHaveLength(1); + expect(mocks.run.mock.calls.filter(([label]) => label.endsWith(":thread:delete"))).toHaveLength( + 1, + ); +}); + +it.each([":terminal:close", ":thread:stop-session"])( + "keeps the thread and files if %s fails", + async (failedOperation) => { + mocks.readThreadShell.mockImplementation(({ threadId }) => { + const thread = threads.find((thread) => thread.id === threadId); + return thread ? { ...thread, session: { status: "ready" } } : null; + }); + mocks.run.mockImplementation(async (label) => + label.endsWith(failedOperation) + ? AsyncResult.failure(Cause.fail(new Error("stop failed"))) + : AsyncResult.success({ sequence: 1 }), + ); + expect((await actions.deleteThread(entries[0]!.threadRef))._tag).toBe("Failure"); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); + expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":thread:delete"))).toBe(false); + }, +); + +it("rechecks shared references when a queued repository deletion actually starts", async () => { + let finishFirst!: () => void; + const firstPending = new Promise((resolve) => { + finishFirst = resolve; + }); + let startFirst!: () => void; + const firstStarted = new Promise((resolve) => { + startFirst = resolve; + }); + let scheduleSecond!: () => void; + const secondScheduled = new Promise((resolve) => { + scheduleSecond = resolve; + }); + mocks.run.mockImplementation(async (label, { input }) => { + if ( + label.endsWith(":thread:delete") && + input.deleteWorktreePath && + input.deleteWorktreePath === "/repo/one" + ) { + startFirst(); + await firstPending; + } + return AsyncResult.success({ sequence: 1 }); + }); + const deletion = deleteSelectedThreadEntries({ + entries: entries.slice(0, 2), + delete: async ({ threadRef }, deletedThreadKeys, deferDeletion) => { + const result = await actions.deleteThread(threadRef, { deletedThreadKeys, deferDeletion }); + if (threadRef.threadId === "two") scheduleSecond(); + return result; + }, + }); + await Promise.all([firstStarted, secondScheduled]); + mocks.readThreadShell.mockImplementation(({ threadId }) => { + const thread = threads.find((entry) => entry.id === threadId); + return threadId === "three" ? { ...thread, worktreePath: "/repo/two" } : thread; + }); + finishFirst(); + expect((await deletion).deletedThreadKeys.size).toBe(2); + expect( + mocks.run.mock.calls + .filter(([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath) + .map(([, { input }]) => input.deleteWorktreePath), + ).toEqual(["/repo/one"]); +}); + +it("retains a thread that changes worktrees while its reference check is pending", async () => { + let finishCheck!: (result: ReturnType>) => void; + const pending = new Promise>>( + (resolve) => { + finishCheck = resolve; + }, + ); + let startCheck!: () => void; + const started = new Promise((resolve) => { + startCheck = resolve; + }); + mocks.archived.mockImplementation(() => { + startCheck(); + return pending; + }); + const deletion = actions.deleteThread(entries[0]!.threadRef); + await started; + mocks.readThreadShell.mockImplementation(({ threadId }) => { + const thread = threads.find((entry) => entry.id === threadId); + return threadId === "one" + ? { ...thread, worktreePath: "/repo/new" } + : { ...thread, worktreePath: "/repo/one" }; + }); + finishCheck(AsyncResult.success({ threads: [] })); + expect((await deletion)._tag).toBe("Failure"); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); + expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":thread:delete"))).toBe(false); +}); + +it("keeps the session and terminal running when an older server cannot delete the worktree", async () => { + mocks.recoverableDeletion = false; + mocks.readThreadShell.mockImplementation(({ threadId }) => { + const thread = threads.find((entry) => entry.id === threadId); + return thread ? { ...thread, session: { status: "ready" } } : null; + }); + const result = await actions.deleteThread(entries[0]!.threadRef); + expect(result._tag).toBe("Failure"); + expect(mocks.run).not.toHaveBeenCalled(); +}); + +it("still deletes on an older server when refreshed references require keeping the worktree", async () => { + mocks.recoverableDeletion = false; + mocks.archived.mockResolvedValue( + AsyncResult.success({ + threads: [{ id: ThreadId.make("archived"), worktreePath: "/repo/one" }], + }), + ); + expect((await actions.deleteThread(entries[0]!.threadRef))._tag).toBe("Success"); + expect(mocks.run.mock.calls.find(([label]) => label.endsWith(":thread:delete"))?.[1]).toEqual({ + environmentId, + input: { threadId: threads[0]!.id }, + }); +}); + +it("reports committed deletion separately from pending cleanup and retries only cleanup", async () => { + const pending = { cwd: "/repo", path: "/repo/.t3-delete-test" }; + mocks.run.mockImplementation(async (label) => + label.endsWith(":thread:delete") + ? AsyncResult.success({ + sequence: 1, + worktreeCleanupPending: { ...pending, retryable: true }, + }) + : AsyncResult.success({ sequence: 1 }), + ); + const outcome = await actions.deleteThread(entries[0]!.threadRef); + expect(outcome._tag).toBe("Success"); + const toast = mocks.toastAdd.mock.calls[0]![0]; + expect(toast.title).toBe("Thread deleted; worktree cleanup incomplete"); + expect(toast.description).toContain(pending.path); + await toast.actionProps!.onClick(); + expect(mocks.run.mock.calls.filter(([label]) => label.endsWith(":thread:delete"))).toHaveLength( + 1, + ); + expect(mocks.run).toHaveBeenCalledWith(expect.stringContaining(":remove-worktree"), { + environmentId, + input: { ...pending, force: true }, + }); + expect(mocks.toastClose).toHaveBeenCalledWith("cleanup-toast"); + expect(mocks.run).toHaveBeenCalledWith(expect.stringContaining(":refresh-status"), { + environmentId, + input: { cwd: "/repo" }, + }); +}); + +it("does not offer destructive retry when cleanup could not be verified", async () => { + mocks.run.mockResolvedValue( + AsyncResult.success({ + sequence: 1, + worktreeCleanupPending: { cwd: "/repo", path: "/repo/.t3-delete-test", retryable: false }, + }), + ); + expect((await actions.deleteThread(entries[0]!.threadRef))._tag).toBe("Success"); + const toast = mocks.toastAdd.mock.calls[0]![0]; + expect(toast.description).toContain("cleanup could not be verified"); + expect(toast.actionProps).toBeUndefined(); +}); diff --git a/apps/web/src/hooks/useThreadActions.ts b/apps/web/src/hooks/useThreadActions.ts index 6a82b920ab31..7d4db50a143b 100644 --- a/apps/web/src/hooks/useThreadActions.ts +++ b/apps/web/src/hooks/useThreadActions.ts @@ -4,7 +4,13 @@ import { scopeThreadRef, scopedThreadKey, } from "@t3tools/client-runtime/environment"; -import { settlePromise, squashAtomCommandFailure } from "@t3tools/client-runtime/state/runtime"; +import { + type AtomCommandResult, + createAtomCommandScheduler, + executeAtomQuery, + settlePromise, + squashAtomCommandFailure, +} from "@t3tools/client-runtime/state/runtime"; import { canSnooze, threadWokeAt } from "@t3tools/client-runtime/state/thread-settled"; import { EnvironmentId, type ScopedThreadRef, ThreadId } from "@t3tools/contracts"; import { resolveWorktreeCleanup } from "@t3tools/shared/projectSettings"; @@ -16,8 +22,9 @@ import { useCallback, useMemo, useRef } from "react"; import { getFallbackThreadIdAfterDelete, pinOrderKeyBetween } from "../components/Sidebar.logic"; import { useComposerDraftStore } from "../composerDraftStore"; -import { terminalEnvironment } from "../state/terminal"; import { appAtomRegistry } from "../rpc/atomRegistry"; +import { orchestrationEnvironment } from "../state/orchestration"; +import { terminalEnvironment } from "../state/terminal"; import { environmentServerConfigsAtom } from "../state/server"; import { threadEnvironment } from "../state/threads"; import { vcsEnvironment } from "../state/vcs"; @@ -31,6 +38,7 @@ import { readEnvironmentSupportsPinReorder, readEnvironmentSupportsActiveReorder, readEnvironmentSupportsSettlement, + readEnvironmentSupportsRecoverableDeletion, readEnvironmentSupportsSnooze, readEnvironmentThreadRefs, readProject, @@ -190,6 +198,8 @@ export async function navigateAfterThreadDeletion(navigate: () => Promise) } } +const deletionScheduler = createAtomCommandScheduler(); + export function useThreadActions() { const closeTerminal = useAtomCommand(terminalEnvironment.close); const archiveThreadMutation = useAtomCommand(threadEnvironment.archive, { @@ -232,9 +242,7 @@ export function useThreadActions() { const removeWorktree = useAtomCommand(vcsEnvironment.removeWorktree, { reportFailure: false, }); - const refreshVcsStatus = useAtomCommand(vcsEnvironment.refreshStatus, { - reportFailure: false, - }); + const refreshVcsStatus = useAtomCommand(vcsEnvironment.refreshStatus, { reportFailure: false }); const sidebarThreadSortOrder = useClientSettings((settings) => settings.sidebarThreadSortOrder); const confirmThreadDelete = useClientSettings((settings) => settings.confirmThreadDelete); const confirmThreadUnpin = useClientSettings((settings) => settings.confirmThreadUnpin); @@ -357,10 +365,35 @@ export function useThreadActions() { ); const deleteThread = useCallback( - async (target: ScopedThreadRef, opts: { deletedThreadKeys?: ReadonlySet } = {}) => { - const resolved = resolveThreadTarget(target); + async ( + target: ScopedThreadRef, + opts: { + deletedThreadKeys?: ReadonlySet; + worktreeDeletionConfirmed?: boolean; + deferDeletion?: (deleteThread: () => Promise>) => void; + } = {}, + ) => { + let resolved = resolveThreadTarget(target); if (!resolved) { - // Thread not in main store (e.g. archived thread) — dispatch delete directly. + const archived = await executeAtomQuery( + appAtomRegistry, + orchestrationEnvironment.archivedShellSnapshot({ + environmentId: target.environmentId, + input: {}, + }), + { refresh: true, reportFailure: false }, + ); + if (archived._tag === "Failure") return archived; + const thread = archived.value.threads.find((entry) => entry.id === target.threadId); + if (thread) { + resolved = { + thread: { ...thread, environmentId: target.environmentId }, + threadRef: target, + }; + } + } + if (!resolved) { + // No live or archived shell remains; dispatch the ordinary idempotent delete. const result = await deleteThreadMutation({ environmentId: target.environmentId, input: { threadId: target.threadId }, @@ -393,7 +426,7 @@ export function useThreadActions() { ? threads.filter((entry) => entry.id === threadRef.threadId || !deletedIds.has(entry.id)) : threads; const orphanedWorktreePath = getOrphanedWorktreePathForThread( - survivingThreads, + [...survivingThreads, thread], threadRef.threadId, ); const displayWorktreePath = orphanedWorktreePath @@ -401,21 +434,21 @@ export function useThreadActions() { : null; const canDeleteWorktree = orphanedWorktreePath !== null && threadProject !== null; const localApi = readLocalApi(); - let shouldDeleteWorktree = false; + let shouldDeleteWorktree = !confirmThreadDelete || opts.worktreeDeletionConfirmed === true; const environmentSettings = appAtomRegistry .get(environmentServerConfigsAtom) .get(threadRef.environmentId)?.settings; const automaticWorktreeCleanup = environmentSettings ? resolveWorktreeCleanup(environmentSettings, thread.projectId).worktreeOnDelete : false; - if (canDeleteWorktree && localApi && !automaticWorktreeCleanup) { + if (canDeleteWorktree && !shouldDeleteWorktree && localApi && !automaticWorktreeCleanup) { const confirmationResult = await settlePromise(() => localApi.dialogs.confirm( [ "This thread is the only one linked to this worktree:", displayWorktreePath ?? orphanedWorktreePath, "", - "Delete the worktree too?", + "Delete the worktree too? Cancel keeps the worktree but still deletes the thread.", ].join("\n"), { variant: "destructive" }, ), @@ -426,122 +459,202 @@ export function useThreadActions() { shouldDeleteWorktree = confirmationResult.value; } - if (thread.session && thread.session.status !== "stopped") { - await stopThreadSession({ + const completeDeletion = async (): Promise> => { + let deleteWorktreePath: string | undefined; + if (shouldDeleteWorktree && orphanedWorktreePath && threadProject) { + // Archived threads are absent from the sidebar; refresh them before removing files. + const archived = await executeAtomQuery( + appAtomRegistry, + orchestrationEnvironment.archivedShellSnapshot({ + environmentId: threadRef.environmentId, + input: {}, + }), + { refresh: true, reportFailure: false }, + ); + if (archived._tag === "Failure") return archived; + const currentThread = + readThreadShell(threadRef) ?? + archived.value.threads.find((entry) => entry.id === thread.id); + if (currentThread && currentThread.worktreePath?.trim() !== orphanedWorktreePath) { + return AsyncResult.failure( + Cause.fail( + new Error("The thread's worktree changed during deletion. Try deleting it again."), + ), + ); + } + const remaining = [ + ...readEnvironmentThreadRefs(threadRef.environmentId).flatMap((ref) => { + const shell = readThreadShell(ref); + return shell === null ? [] : [shell]; + }), + ...archived.value.threads, + ].filter( + (entry) => + !opts.deletedThreadKeys?.has( + scopedThreadKey(scopeThreadRef(threadRef.environmentId, entry.id)), + ), + ); + if ( + getOrphanedWorktreePathForThread([...remaining, thread], thread.id) === + orphanedWorktreePath + ) { + if (!readEnvironmentSupportsRecoverableDeletion(threadRef.environmentId)) { + return AsyncResult.failure( + Cause.fail( + new Error( + "Update this environment's server to delete threads and worktrees safely.", + ), + ), + ); + } + deleteWorktreePath = orphanedWorktreePath; + } + } + + if (thread.session && thread.session.status !== "stopped") { + const stopResult = await stopThreadSession({ + environmentId: threadRef.environmentId, + input: { threadId: threadRef.threadId }, + }); + if (stopResult._tag === "Failure") return stopResult; + } + + const closeResult = await closeTerminal({ environmentId: threadRef.environmentId, - input: { threadId: threadRef.threadId }, + input: { threadId: threadRef.threadId, deleteHistory: false }, }); - } - - await closeTerminal({ - environmentId: threadRef.environmentId, - input: { threadId: threadRef.threadId, deleteHistory: true }, - }); + if (closeResult._tag === "Failure") return closeResult; - const deletedThreadIds = deletedIds ?? new Set(); - const currentRouteThreadRef = getCurrentRouteThreadRef(); - const shouldNavigateToFallback = - currentRouteThreadRef?.threadId === threadRef.threadId && - currentRouteThreadRef.environmentId === threadRef.environmentId; - const fallbackThreadId = getFallbackThreadIdAfterDelete({ - threads, - deletedThreadId: threadRef.threadId, - deletedThreadIds, - sortOrder: sidebarThreadSortOrder, - }); - const deleteResult = await deleteThreadMutation({ - environmentId: threadRef.environmentId, - input: { threadId: threadRef.threadId }, - }); - if (deleteResult._tag === "Failure") { - return deleteResult; - } - refreshArchivedThreadsForEnvironment(threadRef.environmentId); - releaseComposerDraftUploads(threadRef); - clearComposerDraftForThread(threadRef); - clearProjectDraftThreadById( - scopeProjectRef(threadRef.environmentId, thread.projectId), - threadRef, - ); - clearTerminalUiState(threadRef); - - if (shouldNavigateToFallback) { - const fallbackThread = fallbackThreadId - ? readThreadShell(scopeThreadRef(threadRef.environmentId, fallbackThreadId)) - : null; - await navigateAfterThreadDeletion(() => - fallbackThread - ? router.navigate({ - to: "/$environmentId/$threadId", - params: buildThreadRouteParams( - scopeThreadRef(fallbackThread.environmentId, fallbackThread.id), - ), - replace: true, - }) - : router.navigate({ to: "/", replace: true }), + const deletedThreadIds = deletedIds ?? new Set(); + const currentRouteThreadRef = getCurrentRouteThreadRef(); + const shouldNavigateToFallback = + currentRouteThreadRef?.threadId === threadRef.threadId && + currentRouteThreadRef.environmentId === threadRef.environmentId; + const fallbackThreadId = getFallbackThreadIdAfterDelete({ + threads, + deletedThreadId: threadRef.threadId, + deletedThreadIds, + sortOrder: sidebarThreadSortOrder, + }); + const deleteResult = await deleteThreadMutation({ + environmentId: threadRef.environmentId, + input: { + threadId: threadRef.threadId, + ...(deleteWorktreePath ? { deleteWorktreePath } : {}), + }, + }).finally(() => { + if (deleteWorktreePath && threadProject) { + // This command also invalidates live and persisted worktree-ref caches. + // A refresh failure must not turn a committed deletion into a failure. + void settlePromise(() => + refreshVcsStatus({ + environmentId: threadRef.environmentId, + input: { cwd: threadProject.workspaceRoot }, + }), + ); + } + }); + if (deleteResult._tag === "Failure") { + return deleteResult; + } + const pendingCleanup = deleteResult.value.worktreeCleanupPending; + if (pendingCleanup) { + const cleanupToast = toastManager.add( + stackedThreadToast({ + type: "error", + title: "Thread deleted; worktree cleanup incomplete", + description: pendingCleanup.retryable + ? `Remaining files are at ${pendingCleanup.path}. Retry to remove them.` + : `Files are preserved at ${pendingCleanup.path}, but cleanup could not be verified. Refresh and check worktree references before removing them.`, + timeout: 0, + ...(pendingCleanup.retryable + ? { + actionProps: { + children: "Retry cleanup", + onClick: async () => { + const retry = await removeWorktree({ + environmentId: threadRef.environmentId, + input: { + cwd: pendingCleanup.cwd, + path: pendingCleanup.path, + force: true, + }, + }); + if (retry._tag === "Success") toastManager.close(cleanupToast); + else + toastManager.update(cleanupToast, { + description: `Files remain at ${pendingCleanup.path}. ${String(squashAtomCommandFailure(retry))}`, + }); + }, + }, + } + : {}), + }), + ); + } + refreshArchivedThreadsForEnvironment(threadRef.environmentId); + releaseComposerDraftUploads(threadRef); + clearComposerDraftForThread(threadRef); + clearProjectDraftThreadById( + scopeProjectRef(threadRef.environmentId, thread.projectId), + threadRef, ); - } + clearTerminalUiState(threadRef); - if (!shouldDeleteWorktree || !orphanedWorktreePath || !threadProject) { - return deleteResult; - } - - const removeResult = await removeWorktree({ - environmentId: threadRef.environmentId, - input: { - cwd: threadProject.workspaceRoot, - path: orphanedWorktreePath, - force: true, - }, - }); - const refreshResult = - removeResult._tag === "Success" - ? await refreshVcsStatus({ - environmentId: threadRef.environmentId, - input: { cwd: threadProject.workspaceRoot }, - }) - : null; - const cleanupFailure = - removeResult._tag === "Failure" - ? removeResult - : refreshResult?._tag === "Failure" - ? refreshResult + if (shouldNavigateToFallback) { + const fallbackThread = fallbackThreadId + ? readThreadShell(scopeThreadRef(threadRef.environmentId, fallbackThreadId)) : null; - if (cleanupFailure) { - const removalFailed = removeResult._tag === "Failure"; - const error = squashAtomCommandFailure(cleanupFailure); - const message = error instanceof Error ? error.message : "An error occurred."; - console.error("Worktree cleanup failed after thread deletion", { - threadId: threadRef.threadId, - projectCwd: threadProject.workspaceRoot, - worktreePath: orphanedWorktreePath, - error, - }); - toastManager.add( - stackedThreadToast({ - type: "error", - title: removalFailed - ? "Failed to delete worktree" - : "Worktree deleted, but Git status refresh failed", - description: removalFailed - ? `Could not remove ${displayWorktreePath ?? orphanedWorktreePath}. ${message}` - : message, - }), + await navigateAfterThreadDeletion(() => + fallbackThread + ? router.navigate({ + to: "/$environmentId/$threadId", + params: buildThreadRouteParams( + scopeThreadRef(fallbackThread.environmentId, fallbackThread.id), + ), + replace: true, + }) + : router.navigate({ to: "/", replace: true }), + ); + } + + return deleteResult; + }; + const runDeletion = () => + deletionScheduler.schedule( + appAtomRegistry, + { mode: "singleFlight", key: scopedThreadKey }, + threadRef, + () => + deletionScheduler.schedule( + appAtomRegistry, + shouldDeleteWorktree && threadProject && orphanedWorktreePath + ? { + mode: "serial", + key: () => + JSON.stringify([threadRef.environmentId, threadProject.workspaceRoot]), + } + : { mode: "parallel" }, + threadRef, + completeDeletion, + ), ); - // The thread was deleted. Cleanup has its own toast; returning its - // failure would make callers incorrectly report a thread deletion error. + if (opts.deferDeletion && shouldDeleteWorktree && canDeleteWorktree) { + opts.deferDeletion(runDeletion); + return AsyncResult.success(undefined); } - return deleteResult; + return runDeletion(); }, [ clearComposerDraftForThread, clearProjectDraftThreadById, clearTerminalUiState, closeTerminal, + confirmThreadDelete, deleteThreadMutation, getCurrentRouteThreadRef, - refreshVcsStatus, removeWorktree, + refreshVcsStatus, router, resolveThreadTarget, sidebarThreadSortOrder, diff --git a/apps/web/src/state/entities.ts b/apps/web/src/state/entities.ts index af977d567f2b..a56acafca96b 100644 --- a/apps/web/src/state/entities.ts +++ b/apps/web/src/state/entities.ts @@ -201,6 +201,13 @@ export function readEnvironmentSupportsSettlement(environmentId: EnvironmentId): ); } +export function readEnvironmentSupportsRecoverableDeletion(environmentId: EnvironmentId): boolean { + return ( + appAtomRegistry.get(environmentServerConfigsAtom).get(environmentId)?.environment.capabilities + .recoverableThreadDeletion === true + ); +} + /** Whether the environment's server understands thread.snooze/unsnooze. Same version-skew contract as settlement. */ export function readEnvironmentSupportsSnooze(environmentId: EnvironmentId): boolean { diff --git a/packages/contracts/src/environment.ts b/packages/contracts/src/environment.ts index 089133d878a4..096edf08e0bc 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -112,6 +112,8 @@ export const ExecutionEnvironmentCapabilities = Schema.Struct({ pre-settlement servers, so clients treat missing as unsupported and never send the commands under version skew. */ threadSettlement: Schema.optionalKey(Schema.Boolean), + /** Server stages worktree removal and restores it if thread deletion fails. */ + recoverableThreadDeletion: Schema.optionalKey(Schema.Boolean), /** Server evaluates merge and inactivity settlement without a client. */ threadAutoSettlement: Schema.optionalKey(Schema.Boolean), storageCleanup: Schema.optionalKey(Schema.Boolean), diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index 1ae704d3f666..6534be067dbf 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -1137,6 +1137,7 @@ const ThreadDeleteCommand = Schema.Struct({ type: Schema.Literal("thread.delete"), commandId: CommandId, threadId: ThreadId, + deleteWorktreePath: Schema.optional(TrimmedNonEmptyString), }); const ThreadArchiveCommand = Schema.Struct({ @@ -2266,6 +2267,13 @@ export type ProjectionPendingApprovalDecision = typeof ProjectionPendingApproval export const DispatchResult = Schema.Struct({ sequence: NonNegativeInt, + worktreeCleanupPending: Schema.optional( + Schema.Struct({ + cwd: TrimmedNonEmptyString, + path: TrimmedNonEmptyString, + retryable: Schema.Boolean, + }), + ), }); export type DispatchResult = typeof DispatchResult.Type;