From 99a4a25faad6e08783c2b7484c143ba9be479bc0 Mon Sep 17 00:00:00 2001 From: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> Date: Thu, 17 Sep 2026 05:26:27 +0200 Subject: [PATCH 1/5] fix(web): make bulk thread deletion reliable Respect deletion confirmations and retain the last thread for a worktree until cleanup succeeds. Keep failed deletions selected for retry, serialize worktree cleanup per repository, and recheck shared references before removing files. Implemented with GPT-6 Astra (Codex) in T3 Code. Signed-off-by: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> --- apps/web/src/components/LegacySidebar.tsx | 21 +- apps/web/src/components/Sidebar.logic.ts | 18 +- apps/web/src/components/Sidebar.tsx | 17 +- .../components/settings/SettingsPanels.tsx | 2 +- apps/web/src/hooks/useThreadActionMenu.ts | 9 +- .../hooks/useThreadActions.deletion.test.tsx | 354 ++++++++++++++++++ apps/web/src/hooks/useThreadActions.ts | 254 +++++++------ 7 files changed, 549 insertions(+), 126 deletions(-) create mode 100644 apps/web/src/hooks/useThreadActions.deletion.test.tsx diff --git a/apps/web/src/components/LegacySidebar.tsx b/apps/web/src/components/LegacySidebar.tsx index 7a13f01a4742..7d5d7add5451 100644 --- a/apps/web/src/components/LegacySidebar.tsx +++ b/apps/web/src/components/LegacySidebar.tsx @@ -1957,25 +1957,36 @@ 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 deletionToast = toastManager.add({ + type: "loading", + title: `Deleting ${count} thread${count === 1 ? "" : "s"}…`, + timeout: 0, + }); const { deletedThreadKeys, firstFailure } = await deleteSelectedThreadEntries({ entries: selectedThreadEntries, - delete: ({ threadRef }, deletedThreadKeys) => - deleteThread(threadRef, { deletedThreadKeys }), + delete: ({ threadRef }, deletedThreadKeys, deferDeletion) => + deleteThread(threadRef, { + deletedThreadKeys, + deferDeletion, + worktreeDeletionConfirmed: true, + }), }); + toastManager.close(deletionToast); 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.", + 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, }), ); } diff --git a/apps/web/src/components/Sidebar.logic.ts b/apps/web/src/components/Sidebar.logic.ts index 27e47d131e0d..6d0054be1749 100644 --- a/apps/web/src/components/Sidebar.logic.ts +++ b/apps/web/src/components/Sidebar.logic.ts @@ -408,14 +408,25 @@ 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 result = await input.delete(entry, deletedThreadKeys, (deleteThread) => { + deferred = true; + pendingDeletions.push( + deleteThread().then((result) => { + if (result._tag === "Success") deletedThreadKeys.add(entry.threadKey); + else if (!isAtomCommandInterrupted(result)) firstFailure ??= result; + }), + ); + }); + if (result === null || deferred) continue; if (result._tag === "Failure") { if (isAtomCommandInterrupted(result)) break; firstFailure ??= result; @@ -424,6 +435,9 @@ export async function deleteSelectedThreadEntries< deletedThreadKeys.add(entry.threadKey); } + // Worktree removal must succeed before its last thread disappears. Independent + // deletions can overlap; the VCS scheduler serializes work within each repository. + await Promise.all(pendingDeletions); return { deletedThreadKeys, firstFailure }; } diff --git a/apps/web/src/components/Sidebar.tsx b/apps/web/src/components/Sidebar.tsx index b11032987fd4..6d3638b848cd 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -3955,30 +3955,39 @@ 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 deletionToast = toastManager.add({ + type: "loading", + title: `Deleting ${count} thread${count === 1 ? "" : "s"}…`, + timeout: 0, + }); const { deletedThreadKeys, firstFailure } = await deleteSelectedThreadEntries({ entries: threadKeys.map((threadKey) => ({ threadKey })), - delete: async ({ threadKey }, deletedThreadKeys) => { + 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, }); }, }); + toastManager.close(deletionToast); 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.", + 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, }), ); } diff --git a/apps/web/src/components/settings/SettingsPanels.tsx b/apps/web/src/components/settings/SettingsPanels.tsx index 4fb847df684b..32a0793a6834 100644 --- a/apps/web/src/components/settings/SettingsPanels.tsx +++ b/apps/web/src/components/settings/SettingsPanels.tsx @@ -3005,7 +3005,7 @@ export function GeneralSettingsPanel() { ({ + confirm: vi.fn(async () => true), + run: vi.fn(), + readThreadShell: vi.fn(), + readProject: vi.fn(), + readEnvironmentThreadRefs: vi.fn(), + confirmThreadDelete: false, + archived: vi.fn(), +})); +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, +})); +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() })); + +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.readProject.mockReturnValue({ workspaceRoot: "/repo" }); + mocks.confirm.mockResolvedValue(true); + mocks.run.mockResolvedValue(AsyncResult.success(undefined)); + 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")) deleted.push(input.threadId); + if (label.endsWith(":remove-worktree")) { + removals.push(input.path); + if (removals.length === 3) startCleanup(); + await releaseCleanup; + } + return AsyncResult.success(undefined); + }); + 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(undefined), + ); + await deleteSelectedThreadEntries({ + entries, + delete: ({ threadRef }, deletedThreadKeys, deferDeletion) => + actions.deleteThread(threadRef, { deletedThreadKeys, deferDeletion }), + }); + expect( + mocks.run.mock.calls.filter(([label]) => label.endsWith(":remove-worktree")), + ).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]) => label.endsWith(":remove-worktree"))).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]) => label.endsWith(":remove-worktree"))).toBe(false); +}); + +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]) => label.endsWith(":remove-worktree"))).toBe(false); +}); + +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]) => label.endsWith(":remove-worktree"))).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]) => label.endsWith(":remove-worktree"))).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(":remove-worktree") && input.path === "/repo/two" + ? failure + : AsyncResult.success(undefined), + ); + 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", "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(undefined)); + 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", "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) => { + if (label.endsWith(":remove-worktree")) { + startCleanup(); + await pending; + } + return AsyncResult.success(undefined); + }); + 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]) => label.endsWith(":remove-worktree"))).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(undefined), + ); + expect((await actions.deleteThread(entries[0]!.threadRef))._tag).toBe("Failure"); + expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":remove-worktree"))).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(":remove-worktree") && input.path === "/repo/one") { + startFirst(); + await firstPending; + } + return AsyncResult.success(undefined); + }); + 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]) => label.endsWith(":remove-worktree")) + .map(([, { input }]) => input.path), + ).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]) => label.endsWith(":remove-worktree"))).toBe(false); + expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":thread:delete"))).toBe(false); +}); diff --git a/apps/web/src/hooks/useThreadActions.ts b/apps/web/src/hooks/useThreadActions.ts index 1d3aa4c3abba..c545630e2e79 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 * as Cause from "effect/Cause"; @@ -15,6 +21,8 @@ import { useCallback, useMemo, useRef } from "react"; import { getFallbackThreadIdAfterDelete, pinOrderKeyBetween } from "../components/Sidebar.logic"; import { useComposerDraftStore } from "../composerDraftStore"; +import { appAtomRegistry } from "../rpc/atomRegistry"; +import { orchestrationEnvironment } from "../state/orchestration"; import { terminalEnvironment } from "../state/terminal"; import { threadEnvironment } from "../state/threads"; import { vcsEnvironment } from "../state/vcs"; @@ -172,6 +180,8 @@ export async function navigateAfterThreadDeletion(navigate: () => Promise) } } +const deletionScheduler = createAtomCommandScheduler(); + export function useThreadActions() { const closeTerminal = useAtomCommand(terminalEnvironment.close); const archiveThreadMutation = useAtomCommand(threadEnvironment.archive, { @@ -211,9 +221,6 @@ export function useThreadActions() { const removeWorktree = useAtomCommand(vcsEnvironment.removeWorktree, { 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); @@ -311,7 +318,14 @@ export function useThreadActions() { ); const deleteThread = useCallback( - async (target: ScopedThreadRef, opts: { deletedThreadKeys?: ReadonlySet } = {}) => { + async ( + target: ScopedThreadRef, + opts: { + deletedThreadKeys?: ReadonlySet; + worktreeDeletionConfirmed?: boolean; + deferDeletion?: (deleteThread: () => Promise>) => void; + } = {}, + ) => { const resolved = resolveThreadTarget(target); if (!resolved) { // Thread not in main store (e.g. archived thread) — dispatch delete directly. @@ -355,15 +369,15 @@ export function useThreadActions() { : null; const canDeleteWorktree = orphanedWorktreePath !== null && threadProject !== null; const localApi = readLocalApi(); - let shouldDeleteWorktree = false; - if (canDeleteWorktree && localApi) { + let shouldDeleteWorktree = !confirmThreadDelete || opts.worktreeDeletionConfirmed === true; + if (canDeleteWorktree && !shouldDeleteWorktree && localApi) { 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" }, ), @@ -374,121 +388,149 @@ export function useThreadActions() { shouldDeleteWorktree = confirmationResult.value; } - if (thread.session && thread.session.status !== "stopped") { - await stopThreadSession({ + const completeDeletion = async (): Promise> => { + 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 }, }); - } + if (closeResult._tag === "Failure") return closeResult; - await closeTerminal({ - environmentId: threadRef.environmentId, - input: { threadId: threadRef.threadId, deleteHistory: true }, - }); + 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 + ) { + const removeResult = await removeWorktree({ + environmentId: threadRef.environmentId, + input: { + cwd: threadProject.workspaceRoot, + path: orphanedWorktreePath, + force: true, + }, + }); + if (removeResult._tag === "Failure") return removeResult; + } + } - 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 }, + }); + if (deleteResult._tag === "Failure") { + return deleteResult; + } + refreshArchivedThreadsForEnvironment(threadRef.environmentId); + releaseComposerDraftUploads(threadRef); + clearComposerDraftForThread(threadRef); + clearProjectDraftThreadById( + scopeProjectRef(threadRef.environmentId, thread.projectId), + threadRef, ); - } - - if (!shouldDeleteWorktree || !orphanedWorktreePath || !threadProject) { - return deleteResult; - } + clearTerminalUiState(threadRef); - 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, router, resolveThreadTarget, From 7fd630b1c840241b6051aa25a29ce3e7f1815f50 Mon Sep 17 00:00:00 2001 From: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> Date: Thu, 17 Sep 2026 06:08:41 +0200 Subject: [PATCH 2/5] fix: preserve worktrees when thread deletion fails Stage worktrees on the server before committing thread deletion, restore them on failure, and recheck references at commit time. Report post-commit cleanup separately and avoid destructive retries when cleanup cannot be verified. Drain rejected bulk deletions and always close their loading notifications. Add focused recovery and concurrency tests. Implemented with GPT-6 Astra (Codex) in T3 Code. Signed-off-by: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> --- .../src/environment/ServerEnvironment.ts | 1 + apps/server/src/orchestration/decider.ts | 22 +- .../threadWorktreeDeletion.test.ts | 401 ++++++++++++++++++ .../orchestration/threadWorktreeDeletion.ts | 153 +++++++ apps/server/src/ws.ts | 26 +- apps/web/src/components/LegacySidebar.tsx | 51 +-- apps/web/src/components/Sidebar.logic.test.ts | 43 ++ apps/web/src/components/Sidebar.logic.ts | 26 +- apps/web/src/components/Sidebar.tsx | 57 +-- .../hooks/useThreadActions.deletion.test.tsx | 168 ++++++-- apps/web/src/hooks/useThreadActions.ts | 74 +++- apps/web/src/state/entities.ts | 7 + packages/contracts/src/environment.ts | 2 + packages/contracts/src/orchestration.ts | 8 + 14 files changed, 935 insertions(+), 104 deletions(-) create mode 100644 apps/server/src/orchestration/threadWorktreeDeletion.test.ts create mode 100644 apps/server/src/orchestration/threadWorktreeDeletion.ts diff --git a/apps/server/src/environment/ServerEnvironment.ts b/apps/server/src/environment/ServerEnvironment.ts index 64d8dfab1733..90740051a73a 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* () { pullRequests: true, inlineMessageContext: true, threadSettlement: true, + recoverableThreadDeletion: true, threadAutoSettlement: true, threadRestartContinuation: true, projectSettingsOverrides: true, diff --git a/apps/server/src/orchestration/decider.ts b/apps/server/src/orchestration/decider.ts index 0119c0e8599a..8392e6ee0a2f 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..e3dd22a2bf72 --- /dev/null +++ b/apps/server/src/orchestration/threadWorktreeDeletion.test.ts @@ -0,0 +1,401 @@ +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 { withThreadWorktreeDeletion } from "./threadWorktreeDeletion.ts"; +import { PersistenceSqlError } from "../persistence/Errors.ts"; + +const TestLayer = Git.layer.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 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( + Layer.mock(ProjectionSnapshotQuery)({ + getCommandReadModel: () => + snapshotError ? Effect.fail(snapshotError) : Effect.succeed(snapshot), + }), + ), + ); + + return { + fs, + path, + git, + run, + cwd, + root, + worktree, + command, + execute, + 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]); + yield* f.git.removeWorktree({ cwd: f.cwd, path: staged, force: true }); + expect(yield* f.fs.exists(staged)).toBe(false); + }), + ); + + 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..e46c6fa844cb --- /dev/null +++ b/apps/server/src/orchestration/threadWorktreeDeletion.ts @@ -0,0 +1,153 @@ +import * as NodeCrypto from "node:crypto"; +import { + OrchestrationDispatchCommandError, + type DispatchResult, + type OrchestrationCommand, +} 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"; + +const locks = new Map(); + +/** 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); + } + // 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/ws.ts b/apps/server/src/ws.ts index fe1b202599de..b7e1f89ef6cd 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -102,6 +102,7 @@ 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 { withThreadWorktreeDeletion } from "./orchestration/threadWorktreeDeletion.ts"; import { observeRpcEffect as instrumentRpcEffect, observeRpcStream as instrumentRpcStream, @@ -550,6 +551,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; @@ -1852,7 +1856,27 @@ 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( + ProjectionSnapshotQuery.ProjectionSnapshotQuery, + projectionSnapshotQuery, + ), + Effect.mapError((cause) => + toDispatchCommandError(cause, "Failed to delete thread and worktree"), + ), + ) + : dispatch + ).pipe( Effect.tapError(() => cleanupFailedUploadedAttachments(command, normalizedCommand)), ); yield* recordClientCommandAnalytics(normalizedCommand); diff --git a/apps/web/src/components/LegacySidebar.tsx b/apps/web/src/components/LegacySidebar.tsx index 7d5d7add5451..068e8e83a666 100644 --- a/apps/web/src/components/LegacySidebar.tsx +++ b/apps/web/src/components/LegacySidebar.tsx @@ -1969,33 +1969,36 @@ const SidebarProjectItem = memo(function SidebarProjectItem(props: SidebarProjec title: `Deleting ${count} thread${count === 1 ? "" : "s"}…`, timeout: 0, }); - const { deletedThreadKeys, firstFailure } = await deleteSelectedThreadEntries({ - entries: selectedThreadEntries, - delete: ({ threadRef }, deletedThreadKeys, deferDeletion) => - deleteThread(threadRef, { - deletedThreadKeys, - deferDeletion, - worktreeDeletionConfirmed: true, - }), - }); - toastManager.close(deletionToast); - 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, + 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 0a1ea308d128..12dac2ac766e 100644 --- a/apps/web/src/components/Sidebar.logic.test.ts +++ b/apps/web/src/components/Sidebar.logic.test.ts @@ -150,6 +150,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 6d0054be1749..e04546838532 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 type { ContextMenuItem } from "@t3tools/contracts"; @@ -417,15 +418,19 @@ export async function deleteSelectedThreadEntries< for (const entry of input.entries) { let deferred = false; - const result = await input.delete(entry, deletedThreadKeys, (deleteThread) => { - deferred = true; - pendingDeletions.push( - deleteThread().then((result) => { - if (result._tag === "Success") deletedThreadKeys.add(entry.threadKey); - else if (!isAtomCommandInterrupted(result)) firstFailure ??= result; - }), - ); - }); + 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; @@ -435,8 +440,7 @@ export async function deleteSelectedThreadEntries< deletedThreadKeys.add(entry.threadKey); } - // Worktree removal must succeed before its last thread disappears. Independent - // deletions can overlap; the VCS scheduler serializes work within each repository. + // 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 6d3638b848cd..3e9d361cc50c 100644 --- a/apps/web/src/components/Sidebar.tsx +++ b/apps/web/src/components/Sidebar.tsx @@ -3967,36 +3967,39 @@ export default function Sidebar() { title: `Deleting ${count} thread${count === 1 ? "" : "s"}…`, timeout: 0, }); - 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, - }); - }, - }); - toastManager.close(deletionToast); - 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, + 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/hooks/useThreadActions.deletion.test.tsx b/apps/web/src/hooks/useThreadActions.deletion.test.tsx index 33194877f37a..269abb4ae16d 100644 --- a/apps/web/src/hooks/useThreadActions.deletion.test.tsx +++ b/apps/web/src/hooks/useThreadActions.deletion.test.tsx @@ -17,7 +17,17 @@ const mocks = vi.hoisted(() => ({ readProject: vi.fn(), readEnvironmentThreadRefs: vi.fn(), confirmThreadDelete: false, + recoverableDeletion: true, archived: vi.fn(), + toastAdd: vi.fn( + (_toast: { + title: string; + description: string; + actionProps?: { onClick: () => Promise }; + }) => "cleanup-toast", + ), + toastClose: vi.fn(), + toastUpdate: vi.fn(), })); vi.mock("@t3tools/client-runtime/state/runtime", async (importOriginal) => ({ ...(await importOriginal()), @@ -31,6 +41,7 @@ vi.mock("../state/entities", () => ({ readThreadShell: mocks.readThreadShell, readEnvironmentThreadRefs: mocks.readEnvironmentThreadRefs, readProject: mocks.readProject, + readEnvironmentSupportsRecoverableDeletion: () => mocks.recoverableDeletion, })); vi.mock("./useSettings", () => ({ useClientSettings: (select: (settings: object) => unknown) => @@ -46,6 +57,10 @@ 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) => ({ @@ -76,9 +91,10 @@ beforeEach(() => { vi.clearAllMocks(); vi.stubGlobal("IS_REACT_ACT_ENVIRONMENT", true); mocks.confirmThreadDelete = false; + mocks.recoverableDeletion = true; mocks.readProject.mockReturnValue({ workspaceRoot: "/repo" }); mocks.confirm.mockResolvedValue(true); - mocks.run.mockResolvedValue(AsyncResult.success(undefined)); + 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( @@ -107,13 +123,13 @@ it("keeps conversations until worktree cleanup succeeds without blocking the res const deleted: string[] = []; const removals: string[] = []; mocks.run.mockImplementation(async (label, { input }) => { - if (label.endsWith(":thread:delete")) deleted.push(input.threadId); - if (label.endsWith(":remove-worktree")) { - removals.push(input.path); + if (label.endsWith(":thread:delete") && input.deleteWorktreePath) { + removals.push(input.deleteWorktreePath); if (removals.length === 3) startCleanup(); await releaseCleanup; } - return AsyncResult.success(undefined); + if (label.endsWith(":thread:delete")) deleted.push(input.threadId); + return AsyncResult.success({ sequence: 1 }); }); const deletion = deleteSelectedThreadEntries({ entries, @@ -138,7 +154,7 @@ it.each([false, true])( mocks.run.mockImplementation(async (label, { input }) => failFirst && label.endsWith(":thread:delete") && input.threadId === "one" ? AsyncResult.failure(Cause.fail(new Error("delete failed"))) - : AsyncResult.success(undefined), + : AsyncResult.success({ sequence: 1 }), ); await deleteSelectedThreadEntries({ entries, @@ -146,7 +162,9 @@ it.each([false, true])( actions.deleteThread(threadRef, { deletedThreadKeys, deferDeletion }), }); expect( - mocks.run.mock.calls.filter(([label]) => label.endsWith(":remove-worktree")), + mocks.run.mock.calls.filter( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), ).toHaveLength(failFirst ? 0 : 1); }, ); @@ -156,7 +174,11 @@ it("does not repeat a worktree confirmation already included in the bulk confirm act(() => renderer.update()); await actions.deleteThread(entries[0]!.threadRef, { worktreeDeletionConfirmed: true }); expect(mocks.confirm).not.toHaveBeenCalled(); - expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":remove-worktree"))).toBe(true); + 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 () => { @@ -166,7 +188,11 @@ it("still offers to keep a worktree for single deletion when confirmations are o 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]) => label.endsWith(":remove-worktree"))).toBe(false); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); }); it("keeps a worktree still used by an archived thread", async () => { @@ -176,14 +202,22 @@ it("keeps a worktree still used by an archived thread", async () => { }), ); await actions.deleteThread(entries[0]!.threadRef); - expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":remove-worktree"))).toBe(false); + expect( + mocks.run.mock.calls.some( + ([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath, + ), + ).toBe(false); }); 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]) => label.endsWith(":remove-worktree"))).toBe(false); + 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); }); @@ -199,15 +233,21 @@ it("rechecks live references before deferred cleanup", async () => { worktreePath: "/repo/one", })); await cleanup!(); - expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":remove-worktree"))).toBe(false); + 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(":remove-worktree") && input.path === "/repo/two" + label.endsWith(":thread:delete") && + input.deleteWorktreePath && + input.deleteWorktreePath === "/repo/two" ? failure - : AsyncResult.success(undefined), + : AsyncResult.success({ sequence: 1 }), ); const result = await deleteSelectedThreadEntries({ entries, @@ -220,19 +260,19 @@ it("keeps a failed worktree's thread, finishes other deletions, and allows retry mocks.run.mock.calls .filter(([label]) => label.endsWith(":thread:delete")) .map(([, { input }]) => input.threadId), - ).toEqual(["one", "three"]); + ).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(undefined)); + 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", "three", "two"]); + ).toEqual(["one", "two", "three", "two"]); }); it("shares an in-flight deletion when the same thread is deleted again", async () => { @@ -244,21 +284,23 @@ it("shares an in-flight deletion when the same thread is deleted again", async ( const started = new Promise((resolve) => { startCleanup = resolve; }); - mocks.run.mockImplementation(async (label) => { - if (label.endsWith(":remove-worktree")) { + mocks.run.mockImplementation(async (label, { input }) => { + if (label.endsWith(":thread:delete") && input.deleteWorktreePath) { startCleanup(); await pending; } - return AsyncResult.success(undefined); + 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]) => label.endsWith(":remove-worktree"))).toHaveLength( - 1, - ); + 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, ); @@ -274,10 +316,14 @@ it.each([":terminal:close", ":thread:stop-session"])( mocks.run.mockImplementation(async (label) => label.endsWith(failedOperation) ? AsyncResult.failure(Cause.fail(new Error("stop failed"))) - : AsyncResult.success(undefined), + : AsyncResult.success({ sequence: 1 }), ); expect((await actions.deleteThread(entries[0]!.threadRef))._tag).toBe("Failure"); - expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":remove-worktree"))).toBe(false); + 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); }, ); @@ -296,11 +342,15 @@ it("rechecks shared references when a queued repository deletion actually starts scheduleSecond = resolve; }); mocks.run.mockImplementation(async (label, { input }) => { - if (label.endsWith(":remove-worktree") && input.path === "/repo/one") { + if ( + label.endsWith(":thread:delete") && + input.deleteWorktreePath && + input.deleteWorktreePath === "/repo/one" + ) { startFirst(); await firstPending; } - return AsyncResult.success(undefined); + return AsyncResult.success({ sequence: 1 }); }); const deletion = deleteSelectedThreadEntries({ entries: entries.slice(0, 2), @@ -319,8 +369,8 @@ it("rechecks shared references when a queued repository deletion actually starts expect((await deletion).deletedThreadKeys.size).toBe(2); expect( mocks.run.mock.calls - .filter(([label]) => label.endsWith(":remove-worktree")) - .map(([, { input }]) => input.path), + .filter(([label, { input }]) => label.endsWith(":thread:delete") && input.deleteWorktreePath) + .map(([, { input }]) => input.deleteWorktreePath), ).toEqual(["/repo/one"]); }); @@ -349,6 +399,64 @@ it("retains a thread that changes worktrees while its reference check is pending }); finishCheck(AsyncResult.success({ threads: [] })); expect((await deletion)._tag).toBe("Failure"); - expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":remove-worktree"))).toBe(false); + 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("does not use destructive client cleanup against an older server", async () => { + mocks.recoverableDeletion = false; + const result = await actions.deleteThread(entries[0]!.threadRef); + expect(result._tag).toBe("Failure"); + expect( + mocks.run.mock.calls.some( + ([label]) => label.endsWith(":thread:delete") || label.endsWith(":remove-worktree"), + ), + ).toBe(false); +}); + +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 c545630e2e79..1b2a2a19a702 100644 --- a/apps/web/src/hooks/useThreadActions.ts +++ b/apps/web/src/hooks/useThreadActions.ts @@ -35,6 +35,7 @@ import { readEnvironmentSupportsPinReorder, readEnvironmentSupportsActiveReorder, readEnvironmentSupportsSettlement, + readEnvironmentSupportsRecoverableDeletion, readEnvironmentSupportsSnooze, readEnvironmentThreadRefs, readProject, @@ -221,6 +222,7 @@ export function useThreadActions() { const removeWorktree = useAtomCommand(vcsEnvironment.removeWorktree, { 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); @@ -389,6 +391,7 @@ export function useThreadActions() { } const completeDeletion = async (): Promise> => { + let deleteWorktreePath: string | undefined; if (thread.session && thread.session.status !== "stopped") { const stopResult = await stopThreadSession({ environmentId: threadRef.environmentId, @@ -440,15 +443,16 @@ export function useThreadActions() { getOrphanedWorktreePathForThread([...remaining, thread], thread.id) === orphanedWorktreePath ) { - const removeResult = await removeWorktree({ - environmentId: threadRef.environmentId, - input: { - cwd: threadProject.workspaceRoot, - path: orphanedWorktreePath, - force: true, - }, - }); - if (removeResult._tag === "Failure") return removeResult; + if (!readEnvironmentSupportsRecoverableDeletion(threadRef.environmentId)) { + return AsyncResult.failure( + Cause.fail( + new Error( + "Update this environment's server to delete threads and worktrees safely.", + ), + ), + ); + } + deleteWorktreePath = orphanedWorktreePath; } } @@ -465,11 +469,60 @@ export function useThreadActions() { }); const deleteResult = await deleteThreadMutation({ environmentId: threadRef.environmentId, - input: { threadId: threadRef.threadId }, + 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); @@ -532,6 +585,7 @@ export function useThreadActions() { deleteThreadMutation, getCurrentRouteThreadRef, removeWorktree, + refreshVcsStatus, router, resolveThreadTarget, sidebarThreadSortOrder, diff --git a/apps/web/src/state/entities.ts b/apps/web/src/state/entities.ts index d9610e20717f..8c0bffbc0369 100644 --- a/apps/web/src/state/entities.ts +++ b/apps/web/src/state/entities.ts @@ -193,6 +193,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 6d5690187104..f47b6f40d1a8 100644 --- a/packages/contracts/src/environment.ts +++ b/packages/contracts/src/environment.ts @@ -103,6 +103,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), /** Server persists the opt-in for continuing interrupted threads after restarts. */ diff --git a/packages/contracts/src/orchestration.ts b/packages/contracts/src/orchestration.ts index cdcc5b7e437a..11309b3ade3d 100644 --- a/packages/contracts/src/orchestration.ts +++ b/packages/contracts/src/orchestration.ts @@ -1111,6 +1111,7 @@ const ThreadDeleteCommand = Schema.Struct({ type: Schema.Literal("thread.delete"), commandId: CommandId, threadId: ThreadId, + deleteWorktreePath: Schema.optional(TrimmedNonEmptyString), }); const ThreadArchiveCommand = Schema.Struct({ @@ -2235,6 +2236,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; From 0634558ca5cbf1259bd040a79bc9ac2cd7c043f5 Mon Sep 17 00:00:00 2001 From: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> Date: Thu, 17 Sep 2026 06:41:04 +0200 Subject: [PATCH 3/5] fix: protect archived deletion and worktree cleanup retries Signed-off-by: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> --- .../threadWorktreeDeletion.test.ts | 53 +++++++++++++++---- .../orchestration/threadWorktreeDeletion.ts | 33 ++++++++++++ apps/server/src/ws.ts | 14 ++++- .../hooks/useThreadActions.deletion.test.tsx | 42 +++++++++++++++ apps/web/src/hooks/useThreadActions.ts | 24 +++++++-- 5 files changed, 151 insertions(+), 15 deletions(-) diff --git a/apps/server/src/orchestration/threadWorktreeDeletion.test.ts b/apps/server/src/orchestration/threadWorktreeDeletion.test.ts index e3dd22a2bf72..f17e55b7b80d 100644 --- a/apps/server/src/orchestration/threadWorktreeDeletion.test.ts +++ b/apps/server/src/orchestration/threadWorktreeDeletion.test.ts @@ -19,7 +19,7 @@ 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 { withThreadWorktreeDeletion } from "./threadWorktreeDeletion.ts"; +import { removeUnusedWorktree, withThreadWorktreeDeletion } from "./threadWorktreeDeletion.ts"; import { PersistenceSqlError } from "../persistence/Errors.ts"; const TestLayer = Git.layer.pipe( @@ -91,6 +91,10 @@ const fixture = Effect.gen(function* () { } 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, @@ -121,14 +125,7 @@ const fixture = Effect.gen(function* () { }), ), ), - ).pipe( - Effect.provide( - Layer.mock(ProjectionSnapshotQuery)({ - getCommandReadModel: () => - snapshotError ? Effect.fail(snapshotError) : Effect.succeed(snapshot), - }), - ), - ); + ).pipe(Effect.provide(snapshots)); return { fs, @@ -140,6 +137,10 @@ const fixture = Effect.gen(function* () { 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: () => { @@ -320,11 +321,43 @@ it.layer(TestLayer)("recoverable worktree deletion", (it) => { const staged = result.worktreeCleanupPending!.path; expect(yield* f.fs.readFileString(f.path.join(staged, "ignored"))).toBe("ignored contents"); yield* f.run(["worktree", "unlock", staged]); - yield* f.git.removeWorktree({ cwd: f.cwd, path: staged, force: true }); + 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; diff --git a/apps/server/src/orchestration/threadWorktreeDeletion.ts b/apps/server/src/orchestration/threadWorktreeDeletion.ts index e46c6fa844cb..1040a7a00ea1 100644 --- a/apps/server/src/orchestration/threadWorktreeDeletion.ts +++ b/apps/server/src/orchestration/threadWorktreeDeletion.ts @@ -1,8 +1,10 @@ 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"; @@ -15,6 +17,37 @@ import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.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 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. */ diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index b7e1f89ef6cd..6960dd030b11 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -102,7 +102,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 { withThreadWorktreeDeletion } from "./orchestration/threadWorktreeDeletion.ts"; +import { + removeUnusedWorktree, + withThreadWorktreeDeletion, +} from "./orchestration/threadWorktreeDeletion.ts"; import { observeRpcEffect as instrumentRpcEffect, observeRpcStream as instrumentRpcStream, @@ -3300,7 +3303,14 @@ 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( + ProjectionSnapshotQuery.ProjectionSnapshotQuery, + projectionSnapshotQuery, + ), + Effect.tap(() => refreshGitStatus(input.cwd)), + ), { "rpc.aggregate": "vcs" }, ), [WS_METHODS.vcsCreateRef]: (input) => diff --git a/apps/web/src/hooks/useThreadActions.deletion.test.tsx b/apps/web/src/hooks/useThreadActions.deletion.test.tsx index 269abb4ae16d..8abe2571db67 100644 --- a/apps/web/src/hooks/useThreadActions.deletion.test.tsx +++ b/apps/web/src/hooks/useThreadActions.deletion.test.tsx @@ -209,6 +209,48 @@ it("keeps a worktree still used by an archived thread", async () => { ).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); diff --git a/apps/web/src/hooks/useThreadActions.ts b/apps/web/src/hooks/useThreadActions.ts index 1b2a2a19a702..6fa95a0a6373 100644 --- a/apps/web/src/hooks/useThreadActions.ts +++ b/apps/web/src/hooks/useThreadActions.ts @@ -328,9 +328,27 @@ export function useThreadActions() { deferDeletion?: (deleteThread: () => Promise>) => void; } = {}, ) => { - const resolved = resolveThreadTarget(target); + let resolved = resolveThreadTarget(target); + if (!resolved) { + 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) { - // Thread not in main store (e.g. archived thread) — dispatch delete directly. + // No live or archived shell remains; dispatch the ordinary idempotent delete. const result = await deleteThreadMutation({ environmentId: target.environmentId, input: { threadId: target.threadId }, @@ -363,7 +381,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 From 3b138609937a9ff081607d0411152596ecde254f Mon Sep 17 00:00:00 2001 From: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> Date: Thu, 17 Sep 2026 07:13:54 +0200 Subject: [PATCH 4/5] fix: prevent worktree adoption during cleanup Signed-off-by: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> --- .../Layers/OrchestrationEngine.test.ts | 81 +++++++++++ .../Layers/OrchestrationEngine.ts | 63 ++++++++- .../Layers/ProviderCommandReactor.test.ts | 1 + .../Services/OrchestrationEngine.ts | 6 + .../threadWorktreeDeletion.test.ts | 6 +- .../orchestration/threadWorktreeDeletion.ts | 129 ++++++++++-------- .../src/project/AgentSessionImporter.test.ts | 10 ++ .../src/relay/AgentAwarenessRelay.test.ts | 6 + apps/server/src/server.test.ts | 1 + .../serverRuntimeStartup.reconcile.test.ts | 4 + apps/server/src/serverRuntimeStartup.test.ts | 8 ++ ...serverRuntimeStartup.worktreeSetup.test.ts | 2 + apps/server/src/ws.ts | 8 ++ 13 files changed, 264 insertions(+), 61 deletions(-) diff --git a/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts b/apps/server/src/orchestration/Layers/OrchestrationEngine.test.ts index 4a2e99f3ebc0..8fc7e0655167 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 fb2fadde5e63..9d7eafd1d2f0 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 9bc701af0837..56c743502569 100644 --- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts +++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts @@ -423,6 +423,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/threadWorktreeDeletion.test.ts b/apps/server/src/orchestration/threadWorktreeDeletion.test.ts index f17e55b7b80d..c1b5b49d0d5c 100644 --- a/apps/server/src/orchestration/threadWorktreeDeletion.test.ts +++ b/apps/server/src/orchestration/threadWorktreeDeletion.test.ts @@ -21,8 +21,12 @@ 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 = Git.layer.pipe( +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), ); diff --git a/apps/server/src/orchestration/threadWorktreeDeletion.ts b/apps/server/src/orchestration/threadWorktreeDeletion.ts index 1040a7a00ea1..509ac4f63baa 100644 --- a/apps/server/src/orchestration/threadWorktreeDeletion.ts +++ b/apps/server/src/orchestration/threadWorktreeDeletion.ts @@ -14,6 +14,7 @@ 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(); @@ -31,21 +32,29 @@ export const removeUnusedWorktree = Effect.fn("removeUnusedWorktree")(function* cwd: input.cwd, detail, }); - 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; + 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. @@ -132,48 +141,56 @@ export const withThreadWorktreeDeletion = Effect.fn("withThreadWorktreeDeletion" }); return yield* Effect.failCause(result.cause); } - // 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, + 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; + }), ); - 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), diff --git a/apps/server/src/project/AgentSessionImporter.test.ts b/apps/server/src/project/AgentSessionImporter.test.ts index 4eb03a5cc036..1971b7f08ad0 100644 --- a/apps/server/src/project/AgentSessionImporter.test.ts +++ b/apps/server/src/project/AgentSessionImporter.test.ts @@ -226,6 +226,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({ @@ -331,6 +333,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({ @@ -396,6 +400,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({ @@ -467,6 +473,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), }); @@ -505,6 +513,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 1e7cdc3f00e6..7a2ba8fa989a 100644 --- a/apps/server/src/relay/AgentAwarenessRelay.test.ts +++ b/apps/server/src/relay/AgentAwarenessRelay.test.ts @@ -514,6 +514,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; @@ -740,6 +742,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 d11129a1842a..25d16f162b08 100644 --- a/apps/server/src/server.test.ts +++ b/apps/server/src/server.test.ts @@ -983,6 +983,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 37fd210ee6da..5b84a9ab785e 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( @@ -721,6 +723,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 4a8aca46bdda..1de4828c4a73 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), @@ -340,6 +342,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), @@ -412,6 +416,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), @@ -476,6 +482,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 6960dd030b11..bdfe32253815 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -1870,6 +1870,10 @@ const makeWsRpcLayer = ( ), ).pipe( Effect.provide(worktreeDeletionContext), + Effect.provideService( + OrchestrationEngine.OrchestrationEngineService, + orchestrationEngine, + ), Effect.provideService( ProjectionSnapshotQuery.ProjectionSnapshotQuery, projectionSnapshotQuery, @@ -3305,6 +3309,10 @@ const makeWsRpcLayer = ( WS_METHODS.vcsRemoveWorktree, removeUnusedWorktree(input, gitWorkflow.removeWorktree(input)).pipe( Effect.provide(worktreeDeletionContext), + Effect.provideService( + OrchestrationEngine.OrchestrationEngineService, + orchestrationEngine, + ), Effect.provideService( ProjectionSnapshotQuery.ProjectionSnapshotQuery, projectionSnapshotQuery, From 23dccb34870d9acc8d30b53fce5ce26e8f3eead1 Mon Sep 17 00:00:00 2001 From: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> Date: Thu, 17 Sep 2026 07:55:36 +0200 Subject: [PATCH 5/5] fix(web): check deletion support before stopping resources Signed-off-by: 0x00-sys <71561164+0x00-sys@users.noreply.github.com> --- .../hooks/useThreadActions.deletion.test.tsx | 26 +++++++++++++---- apps/web/src/hooks/useThreadActions.ts | 28 +++++++++---------- 2 files changed, 34 insertions(+), 20 deletions(-) diff --git a/apps/web/src/hooks/useThreadActions.deletion.test.tsx b/apps/web/src/hooks/useThreadActions.deletion.test.tsx index 8abe2571db67..ec13df1d7167 100644 --- a/apps/web/src/hooks/useThreadActions.deletion.test.tsx +++ b/apps/web/src/hooks/useThreadActions.deletion.test.tsx @@ -449,15 +449,29 @@ it("retains a thread that changes worktrees while its reference check is pending expect(mocks.run.mock.calls.some(([label]) => label.endsWith(":thread:delete"))).toBe(false); }); -it("does not use destructive client cleanup against an older server", async () => { +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.mock.calls.some( - ([label]) => label.endsWith(":thread:delete") || label.endsWith(":remove-worktree"), - ), - ).toBe(false); + 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 () => { diff --git a/apps/web/src/hooks/useThreadActions.ts b/apps/web/src/hooks/useThreadActions.ts index 6fa95a0a6373..9d6245e3c901 100644 --- a/apps/web/src/hooks/useThreadActions.ts +++ b/apps/web/src/hooks/useThreadActions.ts @@ -410,20 +410,6 @@ export function useThreadActions() { const completeDeletion = async (): Promise> => { let deleteWorktreePath: string | undefined; - 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, deleteHistory: false }, - }); - if (closeResult._tag === "Failure") return closeResult; - if (shouldDeleteWorktree && orphanedWorktreePath && threadProject) { // Archived threads are absent from the sidebar; refresh them before removing files. const archived = await executeAtomQuery( @@ -474,6 +460,20 @@ export function useThreadActions() { } } + 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, deleteHistory: false }, + }); + if (closeResult._tag === "Failure") return closeResult; + const deletedThreadIds = deletedIds ?? new Set(); const currentRouteThreadRef = getCurrentRouteThreadRef(); const shouldNavigateToFallback =