diff --git a/.changeset/fn-151-task-reset-always-completes.md b/.changeset/fn-151-task-reset-always-completes.md new file mode 100644 index 0000000000..8dd4a60c50 --- /dev/null +++ b/.changeset/fn-151-task-reset-always-completes.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Make task reset safely fence active planning sessions. +category: fix +dev: Reset adds planner reset disposers, releases held symbol locks, and clears discarded-run projections while retaining operator input. diff --git a/docs/architecture.md b/docs/architecture.md index 28147beb0a..b3e17ae11e 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -32,6 +32,9 @@ The dashboard completes an ordinary approval through the public `@fusion/engine` The conservative stranded-shape recovery claims only a complete non-seed prompt after planner staleness grace, with no live/finalizing session, approval park, handoff, graph continuation/work item evidence, or any persisted workflow-run step instance; those step instances are graph-owned evidence even before a result or continuation exists. Automatic unplanned release refusals atomically claim a project/task episode and append one durable task-log diagnostic; changed prompt, dependency, status, fingerprint, or Plan Review node state starts a new episode. Direct-session validation compares PostgreSQL server identity as well as database name, preventing a same-named database on another cluster from becoming the lifecycle lock namespace. + +A task Reset is therefore a planner-inclusive but lock-free cancellation fence. A surviving self-owned worktree registration is reconciled only after normal live-owner and idle checks, while live or foreign holders remain conflicts rather than force-removal candidates. + ## 1) Overview Fusion is an AI-orchestrated task board. It takes tasks through a structured lifecycle (`planning → todo → in-progress → in-review → done → archived`) and automates planning, execution, review, merge, and operational recovery. diff --git a/docs/dashboard-guide.md b/docs/dashboard-guide.md index e27b3ec75b..a5acc75c4d 100644 --- a/docs/dashboard-guide.md +++ b/docs/dashboard-guide.md @@ -74,9 +74,11 @@ Both actions are irreversible; there is no undo after confirming. The dialog clo -**Reset** is destructive and has no undo. After confirmation, Fusion fences active task work, removes only the task-owned standard worktree and the current `.fusion/tasks//PROMPT.md` plan, then atomically returns the same task to its workflow's **Planning/intake** column with pending steps and `needs-replan`. The task ID, title, description, dependencies, workflow selection, attachments, documents, logs/audit history, and branch history remain intact. +**Reset** is destructive and has no undo. After confirmation, Fusion fences executor and planner work without waiting for a planner that needs the reset-held planning lock. It removes only the task-owned standard worktree and current `.fusion/tasks//PROMPT.md` plan, then atomically returns the same task to its workflow's **Planning/intake** column with pending steps and `needs-replan`. -Reset does not support workspace tasks or external/operator-owned, foreign, unsafe, or project-root worktrees. If cancellation, worktree removal, plan removal, runtime finalization, or durable publication fails, Fusion reports incomplete cleanup and does not claim success or expose the task to Planning; retry **Reset** after the reported problem is resolved. A plan-removal failure may leave the worktree already removed, but the stored task lifecycle state remains unchanged until a retry completes. Once the atomic publication commits, a task-file mirror problem is repaired separately and does not turn the successful reset into a false failure. +A stale self-owned session registration is reconciled only under the ordinary liveness and idle gates, then removal is retried once. A live planner, executor claim, or foreign holder is reported as an actionable conflict; Reset never forces deletion over live work. Workspace tasks and external/operator-owned, foreign, unsafe, or project-root worktrees remain unsupported. + +Reset retains the task ID, title, description, dependencies, workflow selection, comments, attachments, attachment-backed artifacts, operator-authored documents and their revisions, spec-lock history, commit associations, logs, and audit history. It clears agent-only documents and run-produced planning, verification, merge, and artifact projections; held symbol locks are released as history. If cancellation, cleanup, or publication fails, Fusion reports incomplete cleanup and does not expose the task to Planning. Once publication commits, a task-file mirror problem is repaired separately and does not turn the successful reset into a false failure. ## Keyboard shortcuts diff --git a/packages/core/src/__tests__/postgres/task-reset-publication.pg.test.ts b/packages/core/src/__tests__/postgres/task-reset-publication.pg.test.ts index 882053a260..ccf3f100e3 100644 --- a/packages/core/src/__tests__/postgres/task-reset-publication.pg.test.ts +++ b/packages/core/src/__tests__/postgres/task-reset-publication.pg.test.ts @@ -1,4 +1,6 @@ import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it } from "vitest"; +import { eq } from "drizzle-orm"; +import * as schema from "../../postgres/schema/index.js"; import { __setResetPublicationFailureForTesting } from "../../task-store/reset-lifecycle.js"; import { pgDescribe, @@ -27,10 +29,10 @@ pgDescribe("TaskStore reset publication", () => { status: "failed", worktree: "/tmp/owned-worktree", branch: "fusion/fn-reset", + branchWriteOrigin: "engine", checkedOutBy: "agent-reset", workflowIrPin: "pin-before-reset", workflowStepResults: [{ workflowStepId: "plan-review", status: "failed" }], - reviewState: { status: "changes-requested" }, awaitingApprovalReason: "plan-review-replan-cap", } as never); const continuation = await store.upsertWorkflowWorkItem({ @@ -68,7 +70,7 @@ pgDescribe("TaskStore reset publication", () => { expect(reset.branch).toBeUndefined(); expect(reset.checkedOutBy).toBeUndefined(); expect(reset.workflowIrPin).toBeUndefined(); - expect(reset.workflowStepResults).toEqual([]); + expect(reset.workflowStepResults ?? []).toEqual([]); expect(reset.reviewState).toBeUndefined(); expect(reset.awaitingApprovalReason).toBeUndefined(); expect(await store.listWorkflowWorkItemsForTask(task.id, { kinds: ["task"] })).toEqual([ @@ -77,9 +79,60 @@ pgDescribe("TaskStore reset publication", () => { expect(await store.hasWorkflowRunStepInstancesForTask(task.id)).toBe(false); }); + it("clears run projections while retaining operator documents, attachments, and released symbol-lock history", async () => { + const { store, task } = await seedPopulatedResetState(); + const originalTitle = task.title; + const originalDescription = task.description; + const now = new Date().toISOString(); + await store.upsertTaskDocument(task.id, { key: "agent-only", content: "discard", author: "agent" }); + await store.upsertTaskDocument(task.id, { key: "operator", content: "keep", author: "user" }); + await store.upsertTaskDocument(task.id, { key: "operator", content: "agent update", author: "agent" }); + await store.updateTask(task.id, { declaredSymbols: ["src/reset.ts#freshStart"] }); + await store.acquireSymbolLocks(["src/reset.ts#freshStart"], { ownerTaskId: task.id }, 60_000); + const db = h.layer().db; + await db.insert(schema.project.artifacts).values([ + { id: `${task.id}-run`, type: "document", title: "run", authorId: "agent", taskId: task.id, createdAt: now, updatedAt: now }, + { id: `${task.id}-attachment`, type: "image", title: "attachment", authorId: "user", taskId: task.id, metadata: { source: "attachment" }, createdAt: now, updatedAt: now }, + ]); + await db.insert(schema.project.currentPlanEvidence).values({ taskId: task.id, version: 1, sourceRevision: 1, sourceHash: `${task.id}-plan`, capturedAt: now, snapshot: {} }); + await db.insert(schema.project.specDriftReports).values({ taskId: task.id, reportHash: `${task.id}-drift`, executionHash: "execution", report: {}, createdAt: now }); + await db.insert(schema.project.taskVerificationRequests).values({ taskId: task.id, requestId: `${task.id}-verify`, status: "pending", profile: "default", command: "true", scope: "package", requestedBy: "agent", requestedAt: now }); + await db.insert(schema.project.unplannedExecutionBlocks).values({ taskId: task.id, episode: "episode", createdAt: now }); + await db.insert(schema.project.completionHandoffMarkers).values({ taskId: task.id, acceptedAt: now, source: "test" }); + await db.insert(schema.project.mergeQueue).values({ taskId: task.id, enqueuedAt: now }); + await db.insert(schema.project.mergeRequests).values({ taskId: task.id, state: "pending", createdAt: now, updatedAt: now }); + + const reset = await store.resetTaskPublication(task.id, "todo"); + + expect(reset.title).toBe(originalTitle); + expect(reset.description).toBe(originalDescription); + expect(reset.declaredSymbols ?? []).toEqual([]); + expect((await db.select().from(schema.project.taskDocuments).where(eq(schema.project.taskDocuments.taskId, task.id))).map((row) => row.key)).toEqual(["operator"]); + expect((await db.select().from(schema.project.taskDocumentRevisions).where(eq(schema.project.taskDocumentRevisions.taskId, task.id))).map((row) => row.key)).toEqual(["operator"]); + expect(await db.select().from(schema.project.artifacts).where(eq(schema.project.artifacts.taskId, task.id))).toEqual([expect.objectContaining({ id: `${task.id}-attachment` })]); + await expect(db.select().from(schema.project.currentPlanEvidence).where(eq(schema.project.currentPlanEvidence.taskId, task.id))).resolves.toEqual([]); + await expect(db.select().from(schema.project.specDriftReports).where(eq(schema.project.specDriftReports.taskId, task.id))).resolves.toEqual([]); + await expect(db.select().from(schema.project.taskVerificationRequests).where(eq(schema.project.taskVerificationRequests.taskId, task.id))).resolves.toEqual([]); + await expect(db.select().from(schema.project.unplannedExecutionBlocks).where(eq(schema.project.unplannedExecutionBlocks.taskId, task.id))).resolves.toEqual([]); + await expect(db.select().from(schema.project.completionHandoffMarkers).where(eq(schema.project.completionHandoffMarkers.taskId, task.id))).resolves.toEqual([]); + await expect(db.select().from(schema.project.mergeQueue).where(eq(schema.project.mergeQueue.taskId, task.id))).resolves.toEqual([]); + await expect(db.select().from(schema.project.mergeRequests).where(eq(schema.project.mergeRequests.taskId, task.id))).resolves.toEqual([]); + expect(await db.select().from(schema.project.symbolLocks).where(eq(schema.project.symbolLocks.ownerTaskId, task.id))).toEqual([expect.objectContaining({ status: "released" })]); + }); + it("rolls back every publication participant after workflow mutation failure", async () => { const { store, task, continuation } = await seedPopulatedResetState(); + const now = new Date().toISOString(); + await store.upsertTaskDocument(task.id, { key: "agent-only", content: "must survive rollback", author: "agent" }); + await h.layer().db.insert(schema.project.artifacts).values({ id: `${task.id}-rollback-artifact`, type: "document", title: "rollback", authorId: "agent", taskId: task.id, createdAt: now, updatedAt: now }); + const [beforeFailure] = await h.layer().db.select({ + column: schema.project.tasks.column, + status: schema.project.tasks.status, + worktree: schema.project.tasks.worktree, + }).from(schema.project.tasks).where(eq(schema.project.tasks.id, task.id)); + let injected = false; const release = __setResetPublicationFailureForTesting(() => { + injected = true; throw new Error("injected reset publication failure"); }); try { @@ -88,12 +141,16 @@ pgDescribe("TaskStore reset publication", () => { release(); } - const durable = await store.getTask(task.id); - expect(durable?.column).toBe("in-progress"); - expect(durable?.status).toBe("failed"); - expect(durable?.worktree).toBe("/tmp/owned-worktree"); - expect(durable?.workflowStepResults).toHaveLength(1); + expect(injected).toBe(true); + const [durable] = await h.layer().db.select({ + column: schema.project.tasks.column, + status: schema.project.tasks.status, + worktree: schema.project.tasks.worktree, + }).from(schema.project.tasks).where(eq(schema.project.tasks.id, task.id)); + expect(durable).toEqual(beforeFailure); expect((await store.listWorkflowWorkItemsForTask(task.id, { kinds: ["task"] })).find((item) => item.id === continuation.id)?.state).toBe("running"); expect(await store.hasWorkflowRunStepInstancesForTask(task.id)).toBe(true); + expect(await h.layer().db.select().from(schema.project.taskDocuments).where(eq(schema.project.taskDocuments.taskId, task.id))).toEqual([expect.objectContaining({ key: "agent-only" })]); + expect(await h.layer().db.select().from(schema.project.artifacts).where(eq(schema.project.artifacts.taskId, task.id))).toEqual([expect.objectContaining({ id: `${task.id}-rollback-artifact` })]); }); }); diff --git a/packages/core/src/__tests__/task-move-disposer.test.ts b/packages/core/src/__tests__/task-move-disposer.test.ts index 5776abd625..228f9078f3 100644 --- a/packages/core/src/__tests__/task-move-disposer.test.ts +++ b/packages/core/src/__tests__/task-move-disposer.test.ts @@ -4,7 +4,9 @@ import { __setTaskMoveDisposalTimeoutForTesting, disposeTaskBeforeMove, disposeTaskBeforeReset, + getTaskResetDisposer, registerTaskMoveDisposer, + registerTaskResetDisposer, } from "../tasks/task-move-disposer.js"; describe("task move disposer", () => { @@ -162,6 +164,42 @@ describe("task move disposer", () => { expect(resetReady).toBe(true); }); + it("starts move and reset owners concurrently before awaiting either", async () => { + const store = {} as never; + let releaseMove!: () => void; + let releaseReset!: () => void; + const move = vi.fn(() => new Promise((resolve) => { releaseMove = resolve; })); + const reset = vi.fn(() => new Promise((resolve) => { releaseReset = resolve; })); + registerTaskMoveDisposer(store, move); + registerTaskResetDisposer(store, reset); + const pending = disposeTaskBeforeReset(store, { id: "FN-BOTH" } as never); + await Promise.resolve(); + expect(move).toHaveBeenCalledOnce(); + expect(reset).toHaveBeenCalledOnce(); + releaseMove(); + releaseReset(); + await expect(pending).resolves.toBeUndefined(); + }); + + it("supports store-scoped reset disposers and independent unregistration", async () => { + const first = {} as never; + const second = {} as never; + const disposer = vi.fn().mockResolvedValue(undefined); + const unregister = registerTaskResetDisposer(first, disposer); + await disposeTaskBeforeReset(second, { id: "FN-OTHER" } as never); + expect(disposer).not.toHaveBeenCalled(); + await disposeTaskBeforeReset(first, { id: "FN-FIRST" } as never); + expect(disposer).toHaveBeenCalledOnce(); + unregister(); + expect(getTaskResetDisposer(first)).toBeUndefined(); + }); + + it("propagates reset-only disposer rejection", async () => { + const store = {} as never; + registerTaskResetDisposer(store, vi.fn().mockRejectedValue(new Error("planner still active"))); + await expect(disposeTaskBeforeReset(store, { id: "FN-RESET-ONLY" } as never)).rejects.toThrow("planner still active"); + }); + it("is a no-op when reset has no registered runtime owners", async () => { await expect(disposeTaskBeforeReset({} as never, { id: "FN-RESET-NO-OWNER" } as never)).resolves.toBeUndefined(); }); diff --git a/packages/core/src/index.gate.ts b/packages/core/src/index.gate.ts index fe48cfed9f..b62d911611 100644 --- a/packages/core/src/index.gate.ts +++ b/packages/core/src/index.gate.ts @@ -643,7 +643,9 @@ export { disposeTaskBeforeMove, disposeTaskBeforeReset, getTaskMoveDisposer, + getTaskResetDisposer, registerTaskMoveDisposer, + registerTaskResetDisposer, type TaskMoveDisposer, type TaskResetDisposer, type TaskMoveDisposalInput, diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 9331e06227..521f0dfc5c 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -773,7 +773,9 @@ export { disposeTaskBeforeMove, disposeTaskBeforeReset, getTaskMoveDisposer, + getTaskResetDisposer, registerTaskMoveDisposer, + registerTaskResetDisposer, type TaskMoveDisposer, type TaskResetDisposer, type TaskMoveDisposalInput, diff --git a/packages/core/src/task-store/reset-lifecycle.ts b/packages/core/src/task-store/reset-lifecycle.ts index be3a948fb0..9854396745 100644 --- a/packages/core/src/task-store/reset-lifecycle.ts +++ b/packages/core/src/task-store/reset-lifecycle.ts @@ -1,4 +1,4 @@ -import { and, eq, inArray } from "drizzle-orm"; +import { and, eq, inArray, sql } from "drizzle-orm"; import type { ColumnId, Task, TaskStep } from "../types.js"; import * as schema from "../postgres/schema/index.js"; import { projectScopeFor } from "../postgres/data-layer.js"; @@ -7,6 +7,7 @@ import { withTaskWorkflowSerialization } from "./async/async-workflow-workitems. import { readTaskRowInTransaction, upsertTaskRowInTransaction } from "./async/async-persistence.js"; import type { TaskStore } from "../store.js"; import { createLogger } from "../process/logger.js"; +import { resolveTaskSymbolsForTask } from "../tasks/task-symbol-resolution.js"; const resetLog = createLogger("task-store-reset-lifecycle"); const ACTIVE_TASK_CONTINUATION_STATES = ["runnable", "running", "held", "retrying"] as const; @@ -138,6 +139,14 @@ export async function resetTaskPublicationImpl( throw new Error("Atomic task reset publication requires the PostgreSQL backend"); } const projectId = layer.projectId; + const beforeReset = await store.getTask(taskId); + if (!beforeReset) throw new Error(`Task ${taskId} not found`); + const symbols = resolveTaskSymbolsForTask(beforeReset); + /* + FNXC:TaskReset 2026-08-22-04:45: + Symbol release is intentionally before publication: it owns a separate transaction, while publication clears declaredSymbols. Releasing preserves audit history instead of leaking held rows until expiry. + */ + if (store.backendMode && symbols.resolvable) await store.releaseSymbolLocks(symbols.symbols, taskId); let published!: Task; await layer.transactionImmediate(async (tx) => { @@ -161,6 +170,33 @@ export async function resetTaskPublicationImpl( .where(and(scope, inArray(schema.project.workflowWorkItems.id, active.map((row) => row.id)))); } await resetPublicationFailureForTesting?.(); + const documentScope = projectScopeFor(schema.project.taskDocuments.projectId, projectId); + const revisionScope = projectScopeFor(schema.project.taskDocumentRevisions.projectId, projectId); + const documents = await tx.select({ key: schema.project.taskDocuments.key, author: schema.project.taskDocuments.author }) + .from(schema.project.taskDocuments).where(and(documentScope, eq(schema.project.taskDocuments.taskId, taskId))); + const revisions = await tx.select({ key: schema.project.taskDocumentRevisions.key, author: schema.project.taskDocumentRevisions.author }) + .from(schema.project.taskDocumentRevisions).where(and(revisionScope, eq(schema.project.taskDocumentRevisions.taskId, taskId))); + /* + FNXC:TaskReset 2026-08-22-04:45: + Reset retains user-authored documents and their complete revision history. Agent-only documents and run projections are discarded, while attachments, spec-locks, commit associations, and audit history remain operator history. + */ + const userTouchedKeys = new Set([...documents, ...revisions].filter((row) => row.author === "user").map((row) => row.key)); + const removableKeys = documents.filter((row) => !userTouchedKeys.has(row.key)).map((row) => row.key); + if (removableKeys.length) { + await tx.delete(schema.project.taskDocumentRevisions).where(and(revisionScope, eq(schema.project.taskDocumentRevisions.taskId, taskId), inArray(schema.project.taskDocumentRevisions.key, removableKeys))); + await tx.delete(schema.project.taskDocuments).where(and(documentScope, eq(schema.project.taskDocuments.taskId, taskId), inArray(schema.project.taskDocuments.key, removableKeys))); + } + await tx.delete(schema.project.currentPlanEvidence).where(and(projectScopeFor(schema.project.currentPlanEvidence.projectId, projectId), eq(schema.project.currentPlanEvidence.taskId, taskId))); + await tx.delete(schema.project.specDriftReports).where(and(projectScopeFor(schema.project.specDriftReports.projectId, projectId), eq(schema.project.specDriftReports.taskId, taskId))); + await tx.delete(schema.project.taskVerificationRequests).where(and(projectScopeFor(schema.project.taskVerificationRequests.projectId, projectId), eq(schema.project.taskVerificationRequests.taskId, taskId))); + await tx.delete(schema.project.unplannedExecutionBlocks).where(and(projectScopeFor(schema.project.unplannedExecutionBlocks.projectId, projectId), eq(schema.project.unplannedExecutionBlocks.taskId, taskId))); + await tx.delete(schema.project.completionHandoffMarkers).where(and(projectScopeFor(schema.project.completionHandoffMarkers.projectId, projectId), eq(schema.project.completionHandoffMarkers.taskId, taskId))); + await tx.delete(schema.project.mergeQueue).where(and(projectScopeFor(schema.project.mergeQueue.projectId, projectId), eq(schema.project.mergeQueue.taskId, taskId))); + await tx.delete(schema.project.mergeRequests).where(and(projectScopeFor(schema.project.mergeRequests.projectId, projectId), eq(schema.project.mergeRequests.taskId, taskId))); + await tx.delete(schema.project.artifacts).where(and( + projectScopeFor(schema.project.artifacts.projectId, projectId), eq(schema.project.artifacts.taskId, taskId), + sql`coalesce(${schema.project.artifacts.metadata}->>'source', '') <> 'attachment'`, + )); await tx.delete(schema.project.workflowRunStepInstances).where(and( projectScopeFor(schema.project.workflowRunStepInstances.projectId, projectId), eq(schema.project.workflowRunStepInstances.taskId, taskId), diff --git a/packages/core/src/tasks/task-move-disposer.ts b/packages/core/src/tasks/task-move-disposer.ts index 2cb900a929..77ab49b0a1 100644 --- a/packages/core/src/tasks/task-move-disposer.ts +++ b/packages/core/src/tasks/task-move-disposer.ts @@ -20,6 +20,11 @@ export interface TaskMoveDisposalInput { * owned by another store. A set preserves every live owner during overlap. */ const disposers = new WeakMap>(); +/* +FNXC:TaskReset 2026-08-22-04:32: +Reset fences planner owners in addition to move owners because a reset discards task state while an ordinary move does not. Reset disposers run under the caller's held, non-reentrant planning-lifecycle lock, so they must never wait on work that needs that lock. +*/ +const resetDisposers = new WeakMap>(); const TASK_MOVE_DISPOSAL_TIMEOUT_MS = 30_000; let taskMoveDisposalTimeoutMs = TASK_MOVE_DISPOSAL_TIMEOUT_MS; @@ -48,6 +53,25 @@ export function getTaskMoveDisposer(store: TaskStore): TaskMoveDisposer | undefi }; } +export function registerTaskResetDisposer(store: TaskStore, disposer: TaskResetDisposer): () => void { + const registered = resetDisposers.get(store) ?? new Set(); + registered.add(disposer); + resetDisposers.set(store, registered); + return () => { + const current = resetDisposers.get(store); + current?.delete(disposer); + if (current?.size === 0) resetDisposers.delete(store); + }; +} + +export function getTaskResetDisposer(store: TaskStore): TaskResetDisposer | undefined { + const registered = resetDisposers.get(store); + if (!registered?.size) return undefined; + return async (task) => { + await Promise.all([...registered].map((disposer) => disposer(task))); + }; +} + /** * FNXC:WorkflowLifecycle 2026-07-18-14:32: * A user move from active execution back to Todo is a hard cancel. Await every @@ -133,13 +157,19 @@ FNXC:TaskReset 2026-08-19-06:30: Reset is a destructive fresh-planning boundary, so it fences every registered runtime owner regardless of the task's current column before filesystem cleanup begins. The timeout is fail-closed: worktree and plan deletion never starts while an executor, agent, CLI, or planner still holds the task. */ export async function disposeTaskBeforeReset(store: TaskStore, task: Task): Promise { - const disposer = getTaskMoveDisposer(store); - if (!disposer) return; + const moveDisposer = getTaskMoveDisposer(store); + const resetDisposer = getTaskResetDisposer(store); + if (!moveDisposer && !resetDisposer) return; + // Start both domains before awaiting either so cancellation cannot serialize owners. + const disposal = Promise.all([ + moveDisposer?.(task), + resetDisposer?.(task), + ]); let timeout: ReturnType | undefined; try { await Promise.race([ - disposer(task), + disposal, new Promise((_resolve, reject) => { timeout = setTimeout(() => { reject(new Error(`Timed out stopping active work for ${task.id} before resetting the task`)); diff --git a/packages/dashboard/src/__tests__/task-reset-lifecycle.test.ts b/packages/dashboard/src/__tests__/task-reset-lifecycle.test.ts index c1a58210ec..5819e59755 100644 --- a/packages/dashboard/src/__tests__/task-reset-lifecycle.test.ts +++ b/packages/dashboard/src/__tests__/task-reset-lifecycle.test.ts @@ -5,8 +5,8 @@ import { mkdir, mkdtemp, readFile, stat, writeFile } from "node:fs/promises"; import { join } from "node:path"; import { tmpdir } from "node:os"; import type { Task, TaskStore } from "@fusion/core"; -import { registerTaskMoveDisposer } from "@fusion/core"; -import { getRegisteredWorktreeBranches } from "@fusion/engine"; +import { registerTaskMoveDisposer, registerTaskResetDisposer } from "@fusion/core"; +import { activeSessionRegistry, getRegisteredWorktreeBranches, registerPlanningLivenessProbe } from "@fusion/engine"; import { createApiRoutes } from "../routes.js"; import { request as performRequest } from "../test-request.js"; @@ -19,6 +19,14 @@ vi.mock("@fusion/engine", async () => { await rm(input.worktreePath, { recursive: true, force: true }); return { removed: true, classification: "removed" }; }), + removeTaskResetWorktree: vi.fn(async (input: Parameters[0]) => await actual.removeTaskResetWorktree({ + ...input, + remove: async ({ worktreePath }) => { + const { rm } = await import("node:fs/promises"); + await rm(worktreePath, { recursive: true, force: true }); + return { removed: true, classification: "removed" }; + }, + })), pruneWorktreeAdminEntries: vi.fn().mockResolvedValue(undefined), getRegisteredWorktreeBranches: vi.fn().mockResolvedValue([]), }; @@ -133,6 +141,104 @@ describe("POST /tasks/:id/reset", () => { } }); + it("resets a planning-owned worktree after the planner reset fence releases it", async () => { + const root = await mkdtemp(join(tmpdir(), "fusion-reset-route-planning-")); + const worktree = join(root, ".worktrees", "fn-400"); + const taskDir = join(root, ".fusion", "tasks", "FN-400"); + await mkdir(worktree, { recursive: true }); + await mkdir(taskDir, { recursive: true }); + await writeFile(join(taskDir, "PROMPT.md"), "# Discarded plan\n"); + const task = taskFixture(worktree); + vi.mocked(getRegisteredWorktreeBranches).mockResolvedValue([{ branch: task.branch!, worktreePath: worktree }]); + activeSessionRegistry.registerPath(worktree, { taskId: task.id, kind: "planning", ownerKey: `planning:${task.id}` }); + const publication = vi.fn().mockResolvedValue({ ...task, column: "triage", status: "needs-replan", worktree: undefined, branch: undefined }); + const store = createStore(root, task, [], publication); + const unregister = registerTaskResetDisposer(store, async () => activeSessionRegistry.unregisterPath(worktree)); + try { + const res = await performRequest(createApp(store), "POST", "/api/tasks/FN-400/reset", JSON.stringify({ confirm: true }), { "content-type": "application/json" }); + expect(res.status).toBe(200); + expect(publication).toHaveBeenCalledOnce(); + expect(activeSessionRegistry.lookupByPath(worktree)).toBeNull(); + await expect(readFile(join(taskDir, "PROMPT.md"), "utf8")).rejects.toMatchObject({ code: "ENOENT" }); + await expect(stat(worktree)).rejects.toMatchObject({ code: "ENOENT" }); + expect(JSON.stringify(res.body)).not.toContain("cannot remove active-session worktree"); + } finally { + unregister(); + activeSessionRegistry.unregisterPath(worktree); + } + }); + + it("reconciles an aged orphaned planning registration when no disposer owns it", async () => { + const root = await mkdtemp(join(tmpdir(), "fusion-reset-route-stale-planning-")); + const worktree = join(root, ".worktrees", "fn-400"); + await mkdir(worktree, { recursive: true }); + await mkdir(join(root, ".fusion", "tasks", "FN-400"), { recursive: true }); + await writeFile(join(root, ".fusion", "tasks", "FN-400", "PROMPT.md"), "# Discarded plan\n"); + const task = taskFixture(worktree); + vi.mocked(getRegisteredWorktreeBranches).mockResolvedValue([{ branch: task.branch!, worktreePath: worktree }]); + activeSessionRegistry.registerPath(worktree, { taskId: task.id, kind: "planning", ownerKey: `planning:${task.id}` }); + (activeSessionRegistry.lookupByPath(worktree) as { registeredAt: number }).registeredAt = 0; + const publication = vi.fn().mockResolvedValue({ ...task, column: "triage", status: "needs-replan", worktree: undefined, branch: undefined }); + const store = createStore(root, task, [], publication); + try { + const res = await performRequest(createApp(store), "POST", "/api/tasks/FN-400/reset", JSON.stringify({ confirm: true }), { "content-type": "application/json" }); + expect(res.status).toBe(200); + expect(activeSessionRegistry.lookupByPath(worktree)).toBeNull(); + expect(publication).toHaveBeenCalledOnce(); + } finally { + activeSessionRegistry.unregisterPath(worktree); + } + }); + + it("reports a live planner as an actionable conflict without deleting its plan", async () => { + const root = await mkdtemp(join(tmpdir(), "fusion-reset-route-live-planning-")); + const worktree = join(root, ".worktrees", "fn-400"); + const promptPath = join(root, ".fusion", "tasks", "FN-400", "PROMPT.md"); + await mkdir(worktree, { recursive: true }); + await mkdir(join(root, ".fusion", "tasks", "FN-400"), { recursive: true }); + await writeFile(promptPath, "# Keep plan\n"); + const task = taskFixture(worktree); + vi.mocked(getRegisteredWorktreeBranches).mockResolvedValue([{ branch: task.branch!, worktreePath: worktree }]); + activeSessionRegistry.registerPath(worktree, { taskId: task.id, kind: "planning", ownerKey: `planning:${task.id}` }); + (activeSessionRegistry.lookupByPath(worktree) as { registeredAt: number }).registeredAt = 0; + const unregisterProbe = registerPlanningLivenessProbe((id) => id === task.id); + const publication = vi.fn(); + const store = createStore(root, task, [], publication); + try { + const res = await performRequest(createApp(store), "POST", "/api/tasks/FN-400/reset", JSON.stringify({ confirm: true }), { "content-type": "application/json" }); + expect(res.status).toBe(409); + expect(res.body.error).toMatch(/active task FN-400 \(planning\).*stop or finish/i); + expect(publication).not.toHaveBeenCalled(); + await expect(readFile(promptPath, "utf8")).resolves.toBe("# Keep plan\n"); + } finally { + unregisterProbe(); + activeSessionRegistry.unregisterPath(worktree); + } + }); + + it("reports a foreign session holder as an actionable conflict", async () => { + const root = await mkdtemp(join(tmpdir(), "fusion-reset-route-foreign-session-")); + const worktree = join(root, ".worktrees", "fn-400"); + const promptPath = join(root, ".fusion", "tasks", "FN-400", "PROMPT.md"); + await mkdir(worktree, { recursive: true }); + await mkdir(join(root, ".fusion", "tasks", "FN-400"), { recursive: true }); + await writeFile(promptPath, "# Keep plan\n"); + const task = taskFixture(worktree); + vi.mocked(getRegisteredWorktreeBranches).mockResolvedValue([{ branch: task.branch!, worktreePath: worktree }]); + activeSessionRegistry.registerPath(worktree, { taskId: "FN-OTHER", kind: "planning", ownerKey: "planning:FN-OTHER" }); + const publication = vi.fn(); + const store = createStore(root, task, [], publication); + try { + const res = await performRequest(createApp(store), "POST", "/api/tasks/FN-400/reset", JSON.stringify({ confirm: true }), { "content-type": "application/json" }); + expect(res.status).toBe(409); + expect(res.body.error).toMatch(/active task FN-OTHER \(planning\).*stop or finish/i); + expect(publication).not.toHaveBeenCalled(); + await expect(readFile(promptPath, "utf8")).resolves.toBe("# Keep plan\n"); + } finally { + activeSessionRegistry.unregisterPath(worktree); + } + }); + it("keeps durable state non-replannable when prompt removal fails after worktree cleanup", async () => { const root = await mkdtemp(join(tmpdir(), "fusion-reset-route-failure-")); const worktree = join(root, ".worktrees", "fn-400"); diff --git a/packages/dashboard/src/routes/register-task-workflow-routes.ts b/packages/dashboard/src/routes/register-task-workflow-routes.ts index 40d5a13c84..a01004100c 100644 --- a/packages/dashboard/src/routes/register-task-workflow-routes.ts +++ b/packages/dashboard/src/routes/register-task-workflow-routes.ts @@ -114,8 +114,9 @@ import { // FN-8004 follow-up: shared with SelfHealingManager.recoverStaleMergingStatus so the manual // Retry gate and the automatic sweep agree on when a merge-active stamp is orphaned. isStaleMergeActiveStatus, - removeWorktree, - RemovalReason, + removeTaskResetWorktree, + ResetWorktreeForeignSessionError, + ActiveSessionWorktreeRemovalError, getRegisteredWorktreeBranches, pruneWorktreeAdminEntries, isInsideConfiguredWorktreesDir, @@ -3815,14 +3816,22 @@ export function registerTaskWorkflowRoutes(ctx: ApiRoutesContext, deps: TaskWork if (worktreePath) { if (existsSync(worktreePath)) { - const removal = await removeWorktree({ - worktreePath, - rootDir, - settings, - reason: RemovalReason.TaskReset, - taskId: req.params.id, - expectedOwnerTaskId: req.params.id, - }); + /* + FNXC:TaskReset 2026-08-22-04:32: + Reset has fenced planner and executor owners while holding the planning lock. The helper only reconciles proven-stale self-owned registrations under the normal staleness gates; it never forces a live session. + */ + let removal; + try { + removal = await removeTaskResetWorktree({ worktreePath, rootDir, settings, taskId: req.params.id }); + } catch (error) { + if (error instanceof ResetWorktreeForeignSessionError || error instanceof ActiveSessionWorktreeRemovalError) { + const message = error instanceof ResetWorktreeForeignSessionError + ? `Reset is blocked by active task ${error.details.holderTaskId} (${error.details.holderKind}); stop or finish it before retrying Reset` + : `Reset is blocked by active task ${error.details.taskId} (${error.details.kind}); stop or finish it before retrying Reset`; + throw conflict(message); + } + throw error; + } if (!removal.removed && existsSync(worktreePath)) { throw conflict(`Reset incomplete; worktree removal failed for ${req.params.id}`); } diff --git a/packages/engine/src/__tests__/agent-document-tools.test.ts b/packages/engine/src/__tests__/agent-document-tools.test.ts index 0c2f79fea0..91beb1300d 100644 --- a/packages/engine/src/__tests__/agent-document-tools.test.ts +++ b/packages/engine/src/__tests__/agent-document-tools.test.ts @@ -207,6 +207,55 @@ describe("task_prompt_write tool", () => { expect(getText(result)).toBe(`Updated PROMPT.md for ${TASK_ID}.`); }); + it("refuses a reset-fenced planning write before durable prompt persistence", async () => { + const updateTask = vi.fn(); + const store = { updateTask, getTask: vi.fn().mockResolvedValue({ id: TASK_ID }) } as unknown as TaskStore; + + const result = await runTool( + createTaskPromptWriteTool(store, TASK_ID, undefined, () => false), + "call-reset-fenced", + { content: "# Stale plan" }, + ); + + expect(getText(result)).toContain("was reset before PROMPT.md could be persisted"); + expect(updateTask).not.toHaveBeenCalled(); + }); + + it("refuses a pre-reset planner write after it queues behind the reset lifecycle lock", async () => { + let releaseReset!: () => void; + let enteredQueue!: () => void; + const resetPublication = new Promise((resolve) => { releaseReset = resolve; }); + const queuedBehindReset = new Promise((resolve) => { enteredQueue = resolve; }); + let generationCurrent = true; + const updateTask = vi.fn(); + const store = { + getTask: vi.fn().mockResolvedValue({ id: TASK_ID }), + updateTask, + withPlanningLifecycleLock: vi.fn(async (_id: string, write: () => Promise) => { + enteredQueue(); + await resetPublication; + await write(); + }), + withTaskLock: vi.fn(async (_id: string, write: () => Promise) => await write()), + updateTaskUnlocked: vi.fn(), + isBackendMode: vi.fn(() => false), + } as unknown as TaskStore; + + const resultPromise = runTool( + createTaskPromptWriteTool(store, TASK_ID, undefined, () => generationCurrent), + "call-queued-before-reset", + { content: "# Discarded plan" }, + ); + await queuedBehindReset; + generationCurrent = false; + releaseReset(); + + const result = await resultPromise; + expect(getText(result)).toContain("was reset before PROMPT.md could be persisted"); + expect(store.updateTaskUnlocked).not.toHaveBeenCalled(); + expect(updateTask).not.toHaveBeenCalled(); + }); + it("fails closed when the authoritative prompt read-back is missing or different", async () => { const updateTask = vi.fn().mockResolvedValue({}); const getTask = vi.fn().mockResolvedValue({ id: TASK_ID, prompt: "" }); diff --git a/packages/engine/src/__tests__/reset-worktree-removal.test.ts b/packages/engine/src/__tests__/reset-worktree-removal.test.ts new file mode 100644 index 0000000000..058e7dd558 --- /dev/null +++ b/packages/engine/src/__tests__/reset-worktree-removal.test.ts @@ -0,0 +1,77 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; +import { DEFAULT_SELF_OWNED_MIN_IDLE_MS, activeSessionRegistry, executingTaskLock } from "../agents/active-session-registry.js"; +import { registerPlanningLivenessProbe } from "../agents/planning-liveness.js"; +import { ActiveSessionWorktreeRemovalError } from "../worktree/worktree-backend.js"; +import { removeTaskResetWorktree, ResetWorktreeForeignSessionError } from "../worktree/remove-reset-worktree.js"; + +const PATH = "/tmp/fn-151-reset-worktree"; +const TASK = "FN-151"; + +afterEach(() => { + activeSessionRegistry.clear(); + executingTaskLock.release(TASK); +}); + +function input(overrides: Partial[0]> = {}) { + return { + worktreePath: PATH, + rootDir: "/tmp", + settings: {}, + taskId: TASK, + now: () => Date.now() + DEFAULT_SELF_OWNED_MIN_IDLE_MS + 1, + remove: vi.fn().mockResolvedValue({ removed: true, classification: "removed" }), + ...overrides, + }; +} + +describe("removeTaskResetWorktree", () => { + it("reconciles an aged dead self-owned planning entry before removal", async () => { + activeSessionRegistry.registerPath(PATH, { taskId: TASK, kind: "planning", ownerKey: `planning:${TASK}` }); + const options = input(); + await expect(removeTaskResetWorktree(options)).resolves.toMatchObject({ removed: true }); + expect(activeSessionRegistry.lookupByPath(PATH)).toBeNull(); + expect(options.remove).toHaveBeenCalledOnce(); + }); + + it("waits once for a recent stale entry and then removes it", async () => { + activeSessionRegistry.registerPath(PATH, { taskId: TASK, kind: "planning", ownerKey: `planning:${TASK}` }); + let now = Date.now(); + const wait = vi.fn().mockImplementation(async (ms: number) => { now += ms + 1; }); + const options = input({ now: () => now, wait }); + await removeTaskResetWorktree(options); + expect(wait).toHaveBeenCalledWith(expect.any(Number)); + expect(wait.mock.calls[0][0]).toBeLessThanOrEqual(DEFAULT_SELF_OWNED_MIN_IDLE_MS); + expect(options.remove).toHaveBeenCalledOnce(); + }); + + it("refuses an aged live planning holder without removing it", async () => { + activeSessionRegistry.registerPath(PATH, { taskId: TASK, kind: "planning", ownerKey: `planning:${TASK}` }); + const unregister = registerPlanningLivenessProbe((id) => id === TASK); + const options = input(); + try { + await expect(removeTaskResetWorktree(options)).rejects.toBeInstanceOf(ActiveSessionWorktreeRemovalError); + expect(activeSessionRegistry.lookupByPath(PATH)?.kind).toBe("planning"); + expect(options.remove).not.toHaveBeenCalled(); + } finally { unregister(); } + }); + + it("refuses a foreign holder without calling removal", async () => { + activeSessionRegistry.registerPath(PATH, { taskId: "FN-OTHER", kind: "planning", ownerKey: "planning:FN-OTHER" }); + const options = input(); + await expect(removeTaskResetWorktree(options)).rejects.toBeInstanceOf(ResetWorktreeForeignSessionError); + expect(options.remove).not.toHaveBeenCalled(); + }); + + it("re-applies liveness gates before its single post-removal race retry", async () => { + const remove = vi.fn() + .mockImplementationOnce(async () => { + activeSessionRegistry.registerPath(PATH, { taskId: TASK, kind: "planning", ownerKey: `planning:${TASK}` }); + throw new ActiveSessionWorktreeRemovalError({ worktreePath: PATH, taskId: TASK, kind: "planning", ownerKey: `planning:${TASK}`, reason: "task-reset" as never }); + }); + const unregister = registerPlanningLivenessProbe((id) => id === TASK); + try { + await expect(removeTaskResetWorktree(input({ remove }))).rejects.toBeInstanceOf(ActiveSessionWorktreeRemovalError); + expect(remove).toHaveBeenCalledOnce(); + } finally { unregister(); } + }); +}); diff --git a/packages/engine/src/__tests__/triage-reset-planning-fence.test.ts b/packages/engine/src/__tests__/triage-reset-planning-fence.test.ts new file mode 100644 index 0000000000..dd9938b59e --- /dev/null +++ b/packages/engine/src/__tests__/triage-reset-planning-fence.test.ts @@ -0,0 +1,118 @@ +import { describe, expect, it, vi } from "vitest"; +import { disposeTaskBeforeReset, type Task, type TaskStore } from "@fusion/core"; +import { activeSessionRegistry } from "../agents/active-session-registry.js"; +import { PLANNING_RESET_HOLD_MS, PlanningResetFence } from "../planning-reset-fence.js"; +import { TriageProcessor } from "../triage.js"; + +/* +FNXC:TaskReset 2026-08-22-18:10: +The planner fence is deliberately synchronous: Reset must cancel a planner that may be waiting for +its non-reentrant lifecycle lock without awaiting that planner's settlement. +*/ +describe("Triage planning reset fence", () => { + it("invalidates a captured attempt and holds new admission until publication clears it", () => { + let now = 1_000; + const fence = new PlanningResetFence(() => now); + const generation = fence.currentGeneration("FN-151"); + + fence.cancelPlanning("FN-151"); + + expect(fence.isStale("FN-151", generation)).toBe(true); + expect(fence.isResetHoldActive("FN-151")).toBe(true); + now += PLANNING_RESET_HOLD_MS + 1; + expect(fence.isResetHoldActive("FN-151")).toBe(false); + fence.clearHold("FN-151"); + expect(fence.isResetHoldActive("FN-151")).toBe(false); + }); + + it("does not republish a recovered worktree artifact queued behind Reset", async () => { + const task = { + id: "FN-151-QUEUED-ARTIFACT", + description: "reset planner", + column: "triage", + dependencies: [], + steps: [], + currentStep: 0, + log: [], + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + } as Task; + let releaseReset!: () => void; + let enteredQueue!: () => void; + const resetPublication = new Promise((resolve) => { releaseReset = resolve; }); + const entered = new Promise((resolve) => { enteredQueue = resolve; }); + const updateTaskUnlocked = vi.fn(); + const store = { + on: vi.fn(), + off: vi.fn(), + withPlanningLifecycleLock: vi.fn(async (_id: string, work: () => Promise) => { + enteredQueue(); + await resetPublication; + await work(); + }), + withTaskLock: vi.fn(async (_id: string, work: () => Promise) => await work()), + updateTaskUnlocked, + isBackendMode: vi.fn(() => false), + } as unknown as TaskStore; + const processor = new TriageProcessor(store, "/tmp/fn-151-root"); + const generation = (processor as unknown as { resetFence: PlanningResetFence }).resetFence.currentGeneration(task.id); + + try { + const publish = (processor as unknown as { + persistResetFencedPlanningArtifact: (task: Task, generation: number, content: string, mirror: boolean) => Promise; + }).persistResetFencedPlanningArtifact(task, generation, "# Discarded plan", false); + await entered; + await disposeTaskBeforeReset(store, task); + releaseReset(); + + await expect(publish).resolves.toBe(false); + expect(updateTaskUnlocked).not.toHaveBeenCalled(); + } finally { + processor.stop(); + } + }); + + it("fences a live TriageProcessor session without awaiting its lifecycle-lock finalize", async () => { + const task = { + id: "FN-151-RESET-FENCE", + description: "reset planner", + column: "triage", + dependencies: [], + steps: [], + currentStep: 0, + log: [], + createdAt: new Date().toISOString(), + updatedAt: new Date().toISOString(), + } as Task; + const worktree = "/tmp/fn-151-reset-fence"; + const neverSettles = new Promise(() => undefined); + const store = { + on: vi.fn(), + off: vi.fn(), + withPlanningLifecycleLock: vi.fn(() => neverSettles), + } as unknown as TaskStore; + const processor = new TriageProcessor(store, "/tmp/fn-151-root"); + const abort = vi.fn().mockResolvedValue(undefined); + const dispose = vi.fn(() => { throw new Error("synthetic dispose failure"); }); + (processor as unknown as { activeSessions: Map }).activeSessions.set(task.id, { abort, dispose }); + activeSessionRegistry.registerPath(worktree, { taskId: task.id, kind: "planning", ownerKey: `planning:${task.id}` }); + activeSessionRegistry.registerPath(`${worktree}-executor`, { taskId: task.id, kind: "executor", ownerKey: task.id }); + activeSessionRegistry.registerPath(`${worktree}-foreign`, { taskId: "FN-OTHER", kind: "planning", ownerKey: "planning:FN-OTHER" }); + + try { + await expect(disposeTaskBeforeReset(store, task)).resolves.toBeUndefined(); + expect(abort).toHaveBeenCalledOnce(); + expect(dispose).toHaveBeenCalledOnce(); + expect(store.withPlanningLifecycleLock).not.toHaveBeenCalled(); + expect(activeSessionRegistry.lookupByPath(worktree)).toBeNull(); + expect(activeSessionRegistry.lookupByPath(`${worktree}-executor`)?.kind).toBe("executor"); + expect(activeSessionRegistry.lookupByPath(`${worktree}-foreign`)?.taskId).toBe("FN-OTHER"); + expect((processor as unknown as { activeSessions: Map }).activeSessions.has(task.id)).toBe(false); + } finally { + processor.stop(); + activeSessionRegistry.unregisterPath(worktree); + activeSessionRegistry.unregisterPath(`${worktree}-executor`); + activeSessionRegistry.unregisterPath(`${worktree}-foreign`); + } + }); +}); diff --git a/packages/engine/src/agent-tools.ts b/packages/engine/src/agent-tools.ts index a05dc7cd85..42078e8090 100644 --- a/packages/engine/src/agent-tools.ts +++ b/packages/engine/src/agent-tools.ts @@ -2150,7 +2150,12 @@ function parsePlanRepositoryScope(content: string, configured: readonly string[] return [...new Set(repositories)].sort(); } -export function createTaskPromptWriteTool(store: TaskStore, taskId: string, runContext?: RunMutationContext): ToolDefinition { +export function createTaskPromptWriteTool( + store: TaskStore, + taskId: string, + runContext?: RunMutationContext, + canPersist?: () => boolean, +): ToolDefinition { return { name: "fn_task_prompt_write", label: "Write PROMPT.md", @@ -2202,10 +2207,45 @@ export function createTaskPromptWriteTool(store: TaskStore, taskId: string, runC extensions: current?.repositoryScope?.extensions, } : undefined; - await store.updateTask(taskId, { + /* + FNXC:TaskReset 2026-08-22-04:49: + A triage attempt captured before Reset can outlive the route's non-reentrant planning lock. + Check its generation immediately before the authoritative write so a stale planner cannot + recreate PROMPT.md after Reset publishes the description-only state. + */ + if (canPersist && !canPersist()) { + throw new Error(`Planning for ${taskId} was reset before PROMPT.md could be persisted`); + } + const promptUpdate = { prompt: params.content, ...(repositoryScope ? { repositoryScope } : {}), - }, runContext); + }; + if (canPersist) { + /* + FNXC:TaskReset 2026-08-22-18:07: + Reset serializes its publication with the non-reentrant planning lifecycle lock. A + generation check before `updateTask` is insufficient because that writer can queue behind + Reset and recreate PROMPT.md after reset commits. Check and write under the same lock, + using the lock-held TaskStore variants so this tool never re-enters the advisory lock. + */ + await store.withPlanningLifecycleLock(taskId, async () => { + if (!canPersist()) { + throw new Error(`Planning for ${taskId} was reset before PROMPT.md could be persisted`); + } + const updated = await store.withTaskLock(taskId, () => store.updateTaskUnlocked(taskId, promptUpdate, runContext)); + if (store.isBackendMode()) { + await store.reconcileSpecDriftWhilePlanningLocked(updated).catch((error: unknown) => { + log.warn(`[spec-lock] deferred drift reconciliation for ${updated.id}: ${error instanceof Error ? error.message : String(error)}`); + }); + } + /* FNXC:TaskReset 2026-08-22-18:15: Reset clears agent-authored plan documents under this lifecycle lock. */ + await mirrorPlanToProjectDb(store, taskId, params.content, { + author: runContext?.agentId ?? "agent", + }); + }); + } else { + await store.updateTask(taskId, promptUpdate, runContext); + } /* FNXC:PlanArtifactPersistence 2026-08-22-03:37: FN-094 repointed this fail-closed check at updateTask's task-row return value after adding @@ -2225,13 +2265,15 @@ export function createTaskPromptWriteTool(store: TaskStore, taskId: string, runC /* FNXC:PlanArtifactPersistence 2026-07-26-03:55: `updateTask({ prompt })` writes the project-root PROMPT.md and task.json, but `project.tasks` has - no `prompt` column — the spec would live only as a file in the project checkout. Mirror it into the - `plan` task document so the plan is durable in the project database too. Best-effort: a mirror - failure must not fail a write whose authoritative persistence was just verified above. + no `prompt` column — the spec would live only as a file in the project checkout. Reset-fenced + planners mirror it inside the serialized prompt publication above, preventing Reset from clearing + the document and then observing a stale mirror recreated after its transaction commits. */ - await mirrorPlanToProjectDb(store, taskId, params.content, { - author: runContext?.agentId ?? "agent", - }); + if (!canPersist) { + await mirrorPlanToProjectDb(store, taskId, params.content, { + author: runContext?.agentId ?? "agent", + }); + } return { content: [{ type: "text" as const, text: `Updated PROMPT.md for ${taskId}.` }], details: {}, diff --git a/packages/engine/src/agents/planning-liveness.ts b/packages/engine/src/agents/planning-liveness.ts new file mode 100644 index 0000000000..721b8d4782 --- /dev/null +++ b/packages/engine/src/agents/planning-liveness.ts @@ -0,0 +1,19 @@ +export type PlanningLivenessProbe = (taskId: string) => boolean; +const probes = new Set(); +/* +FNXC:TaskReset 2026-08-22-04:32: +Planning probes are process-wide and multi-project safe. Any true or throwing probe refuses removal: ambiguity must preserve a possibly live planner; no probe means no in-process planner exists. +*/ +export const planningLivenessRegistry = { + registerPlanningLivenessProbe(probe: PlanningLivenessProbe): () => void { + probes.add(probe); + return () => probes.delete(probe); + }, + isPlanningLive(taskId: string): boolean { + for (const probe of probes) { try { if (probe(taskId)) return true; } catch { return true; } } + return false; + }, + hasRegisteredProbe(): boolean { return probes.size > 0; }, +}; +export const registerPlanningLivenessProbe = planningLivenessRegistry.registerPlanningLivenessProbe.bind(planningLivenessRegistry); +export const isPlanningLive = planningLivenessRegistry.isPlanningLive.bind(planningLivenessRegistry); diff --git a/packages/engine/src/index.ts b/packages/engine/src/index.ts index 6284ebbcc5..8e9132172f 100644 --- a/packages/engine/src/index.ts +++ b/packages/engine/src/index.ts @@ -1,4 +1,8 @@ export { AgentLogger, type AgentLoggerOptions, summarizeToolArgs } from "./agents/agent-logger.js"; +export { PlanningResetFence, PLANNING_RESET_HOLD_MS } from "./planning-reset-fence.js"; +export { removeTaskResetWorktree, ResetWorktreeForeignSessionError } from "./worktree/remove-reset-worktree.js"; +export { ActiveSessionWorktreeRemovalError } from "./worktree/worktree-backend.js"; +export { planningLivenessRegistry, registerPlanningLivenessProbe, isPlanningLive } from "./agents/planning-liveness.js"; export { classifyReportHealth, type ReportHealthBucket, diff --git a/packages/engine/src/plan-artifact-writeback.ts b/packages/engine/src/plan-artifact-writeback.ts index 3b09c4d79e..b1aab86ce8 100644 --- a/packages/engine/src/plan-artifact-writeback.ts +++ b/packages/engine/src/plan-artifact-writeback.ts @@ -43,6 +43,13 @@ export interface ReconcileWorktreePlanArtifactOptions { /** cwd the planning session ran in. Equal to `rootDir` when planning did not get a worktree. */ planningCwd: string; logger?: PlanWritebackLogger; + /** + * Optional caller-owned authoritative writer. Triage uses this to serialize a recovered + * worktree artifact with its reset-generation fence before it can recreate PROMPT.md. + */ + writeAuthoritativePrompt?: (content: string) => Promise; + /** Optional caller-owned plan-document writer subject to the same publication fence. */ + mirrorAuthoritativePlan?: (content: string, author?: string) => Promise; } export type PlanWritebackOutcome = @@ -57,7 +64,9 @@ export type PlanWritebackOutcome = /** Worktree copy was copied back into the project `.fusion/` folder. */ | "recovered" /** A worktree copy existed but persisting it failed; the root copy is untouched. */ - | "recovery-failed"; + | "recovery-failed" + /** A reset fenced this planning attempt before its worktree copy could be published. */ + | "recovery-fenced"; export interface ReconcileWorktreePlanArtifactResult { outcome: PlanWritebackOutcome; @@ -109,7 +118,10 @@ export async function reconcileWorktreePlanArtifact( try { // The single validated persistence path: File Scope validation + root PROMPT.md write + task.json sync. - await store.updateTask(taskId, { prompt: worktreeContent }); + const persisted = options.writeAuthoritativePrompt + ? await options.writeAuthoritativePrompt(worktreeContent) + : (await store.updateTask(taskId, { prompt: worktreeContent }), true); + if (!persisted) return { outcome: "recovery-fenced", content: rootContent ?? undefined }; logger?.log?.( `${taskId}: recovered a worktree-local PROMPT.md into the project .fusion folder (${planningCwd})`, ); @@ -171,11 +183,14 @@ export async function persistPlanArtifact( options: ReconcileWorktreePlanArtifactOptions & { author?: string }, ): Promise { const result = await reconcileWorktreePlanArtifact(options); + if (result.outcome === "recovery-fenced") return { ...result, mirrored: false }; const mirrored = result.content - ? await mirrorPlanToProjectDb(options.store, options.taskId, result.content, { - author: options.author, - logger: options.logger, - }) + ? options.mirrorAuthoritativePlan + ? await options.mirrorAuthoritativePlan(result.content, options.author) + : await mirrorPlanToProjectDb(options.store, options.taskId, result.content, { + author: options.author, + logger: options.logger, + }) : false; return { ...result, mirrored }; } diff --git a/packages/engine/src/planning-reset-fence.ts b/packages/engine/src/planning-reset-fence.ts new file mode 100644 index 0000000000..9f0fc072c5 --- /dev/null +++ b/packages/engine/src/planning-reset-fence.ts @@ -0,0 +1,26 @@ +export const PLANNING_RESET_HOLD_MS = 30_000; + +/** + * FNXC:TaskReset 2026-08-22-04:32: + * A reset invalidates planner attempts by generation so an aborted session cannot recreate a worktree or prompt after reset publication. + */ +export class PlanningResetFence { + private readonly generations = new Map(); + private readonly holds = new Map(); + constructor(private readonly now: () => number = Date.now) {} + currentGeneration(taskId: string): number { return this.generations.get(taskId) ?? 0; } + cancelPlanning(taskId: string): number { + const generation = this.currentGeneration(taskId) + 1; + this.generations.set(taskId, generation); + this.holds.set(taskId, this.now() + PLANNING_RESET_HOLD_MS); + return generation; + } + isStale(taskId: string, capturedGeneration: number): boolean { return this.currentGeneration(taskId) !== capturedGeneration; } + isResetHoldActive(taskId: string, now = this.now()): boolean { + const expiresAt = this.holds.get(taskId); + if (!expiresAt) return false; + if (expiresAt <= now) { this.holds.delete(taskId); return false; } + return true; + } + clearHold(taskId: string): void { this.holds.delete(taskId); } +} diff --git a/packages/engine/src/triage.ts b/packages/engine/src/triage.ts index 69cc398601..37f714d876 100644 --- a/packages/engine/src/triage.ts +++ b/packages/engine/src/triage.ts @@ -187,6 +187,8 @@ import { AgentLogger } from "./agents/agent-logger.js"; import { attachAgentUsageTelemetry, emitAgentSessionStart } from "./agents/agent-usage-telemetry.js"; import { emitApprovalMail } from "./agents/approval-mail.js"; import { acquireActiveSessionPath, activeSessionRegistry } from "./agents/active-session-registry.js"; +import { PlanningResetFence } from "./planning-reset-fence.js"; +import { registerPlanningLivenessProbe } from "./agents/planning-liveness.js"; import { resolveAgentInstructions, resolveAgentInstructionsWithRatings, @@ -457,6 +459,11 @@ export class TriageProcessor { private lastPlanThrottleSignature: string | null = null; /** Active agent sessions per task, used to terminate on pause. */ private activeSessions = new Map void }>(); + private readonly resetFence = new PlanningResetFence(); + /** Captured attempt generations let a delayed finalizer reject a reset-fenced planner. */ + private readonly activePlanningGenerations = new Map(); + private unregisterPlanningLiveness?: () => void; + private unregisterResetDisposer?: () => void; /** * Reviewer subagent sessions per task. The spec reviewer (`reviewer.ts`) * creates its own AgentSession that isn't part of `activeSessions`, so @@ -615,12 +622,58 @@ export class TriageProcessor { }; } + /** + * FNXC:TaskReset 2026-08-22-18:15: + * Worktree-artifact recovery and duplicate-marker recovery are planner publications, not raw + * filesystem cleanup. Reset holds this non-reentrant lifecycle lock while clearing its output, + * so each write rechecks the captured generation inside the authoritative serialized mutation; + * a pre-lock check alone can queue behind Reset and republish discarded planning output. + */ + private async persistResetFencedPlanningArtifact( + task: Task, + planningGeneration: number, + content: string, + mirrorPlan: boolean, + ): Promise { + let persisted = false; + await this.store.withPlanningLifecycleLock(task.id, async () => { + if (this.resetFence.isStale(task.id, planningGeneration)) return; + const updated = await this.store.withTaskLock(task.id, () => this.store.updateTaskUnlocked(task.id, { prompt: content })); + if (this.store.isBackendMode()) { + await this.store.reconcileSpecDriftWhilePlanningLocked(updated).catch((error: unknown) => { + planLog.warn(`[spec-lock] deferred drift reconciliation for ${updated.id}: ${error instanceof Error ? error.message : String(error)}`); + }); + } + if (mirrorPlan) { + await mirrorPlanToProjectDb(this.store, task.id, content, { + author: "triage", + logger: { log: (message: string) => planLog.log(message), warn: (message: string) => planLog.warn(message) }, + }); + } + persisted = true; + }); + return persisted; + } + constructor( private store: TaskStore, private rootDir: string, private options: TriageProcessorOptions = {}, ) { this.workflowAgentCapacity = new WorkflowAgentCapacity(this.options.agentStore); + this.unregisterPlanningLiveness = registerPlanningLivenessProbe((taskId) => this.getPlanningTaskIds().has(taskId)); + this.unregisterResetDisposer = fusionCore.registerTaskResetDisposer(this.store, async (task) => { + /* + FNXC:TaskReset 2026-08-22-04:32: + The route already holds the non-reentrant planning lock, so this disposer is synchronous/lock-free: waiting for finalize would deadlock reset against the planner being cancelled. + */ + this.resetFence.cancelPlanning(task.id); + this.abortAndDisposePlanningSessionForTask(task.id, "task reset"); + for (const path of activeSessionRegistry.pathsForTask(task.id)) { + const record = activeSessionRegistry.lookupByPath(path); + if (record?.ownerKey === `planning:${task.id}`) activeSessionRegistry.unregisterPath(path); + } + }); this.unregisterAdmissionProvider = projectAdmissionCoordinator.registerProvider(`specify:${this.rootDir}`, { projectId: this.rootDir, refresh: async () => { @@ -761,6 +814,8 @@ export class TriageProcessor { */ this.taskColumnWakeHandler = (task: Task, meta?: { lanes?: TaskMoveLanes }) => { if (!task?.id) return; + // Reset publication emits this durable state after its held lock commits, so a fresh planner is not delayed by the conservative reset TTL. + if (task.status === "needs-replan") this.resetFence.clearHold(task.id); const isPlannerWakeColumn = meta?.lanes ? task.column === meta.lanes.hold || task.column === meta.lanes.intake : LEGACY_PLANNER_WAKE_COLUMNS.has(task.column); @@ -1044,6 +1099,10 @@ export class TriageProcessor { // Tear down any in-flight specify sessions and reviewer subagents so they // don't keep streaming LLM tokens / tool calls past engine shutdown. this.abortAndDisposeActiveSessions("engine stop"); + this.unregisterPlanningLiveness?.(); + this.unregisterPlanningLiveness = undefined; + this.unregisterResetDisposer?.(); + this.unregisterResetDisposer = undefined; planLog.log("Processor stopped"); } @@ -1056,6 +1115,35 @@ export class TriageProcessor { * interrupts any in-flight LLM stream / tool call; dispose() then * releases session resources. */ + private abortAndDisposePlanningSessionForTask(taskId: string, reason: string): void { + try { + this.disposeSubagentsForTask(taskId, reason); + } catch (error) { + planLog.warn(`${taskId}: failed to dispose planning subagents: ${error instanceof Error ? error.message : String(error)}`); + } + const session = this.activeSessions.get(taskId); + if (!session) return; + this.pauseAborted.add(taskId); + this.options.stuckTaskDetector?.untrackTask(taskId); + const abortable = session as { abort?: () => Promise; dispose: () => void }; + try { + if (typeof abortable.abort === "function") void Promise.resolve(abortable.abort()).catch(() => undefined); + this.recordTriageSessionTokenUsageSoon(taskId, session as AgentSession); + abortable.dispose(); + } catch (error) { + /* + FNXC:TaskReset 2026-08-22-18:10: + Reset must finish fencing its own registry records even when an agent session's synchronous + disposer is faulty. The route's cleanup remains guarded by ownership and staleness checks; + this catch only prevents an already-cancelled planner from turning Reset into a raw failure. + */ + planLog.warn(`${taskId}: failed to dispose planning session: ${error instanceof Error ? error.message : String(error)}`); + } finally { + // A failed disposer cannot retain an active-session bookkeeping claim past Reset. + this.activeSessions.delete(taskId); + } + } + private abortAndDisposeActiveSessions(reason: string): void { for (const taskId of [...this.activeSubagentSessions.keys()]) { this.disposeSubagentsForTask(taskId, reason); @@ -2473,6 +2561,12 @@ export class TriageProcessor { } async specifyTask(task: Task): Promise { + if (this.resetFence.isResetHoldActive(task.id)) { + if (dropPreHeldExecutorSlot(task.id)) this.options.semaphore?.release(); + return; + } + const planningGeneration = this.resetFence.currentGeneration(task.id); + this.activePlanningGenerations.set(task.id, planningGeneration); /* FNXC:TriageStuckKill 2026-07-18-21:05: Refuse a second planner when finalize/Plan Review is still live even if @@ -2763,7 +2857,12 @@ export class TriageProcessor { ...this.createTriageTools({ parentTaskId: task.id }), createTaskDocumentWriteTool(this.store, task.id), createTaskDocumentReadTool(this.store, task.id), - createTaskPromptWriteTool(this.store, task.id, triageRunContext), + createTaskPromptWriteTool( + this.store, + task.id, + triageRunContext, + () => !this.resetFence.isStale(task.id, planningGeneration), + ), createWorkflowListTool(this.store), createWorkflowSelectTool(this.store, task.id), ...(isResearchToolSurfaceEnabled(settings) @@ -3017,6 +3116,7 @@ export class TriageProcessor { bypasses it. `contended` means a genuinely live foreign holder, so planning falls back to the shared checkout rather than running in a worktree someone else owns. */ + if (this.resetFence.isStale(task.id, planningGeneration)) return; const acquired = acquireActiveSessionPath(activeSessionRegistry, planningCwd, { taskId: task.id, kind: "planning", @@ -3376,6 +3476,13 @@ export class TriageProcessor { authoritative into the project database (PROMPT.md has no `tasks` column and is otherwise filesystem-only). Both halves are best-effort — validation below still owns the verdict. */ + /* + FNXC:TaskReset 2026-08-22-04:49: + Reset can commit while a planner session drains. Fence worktree artifact recovery before + it writes the project-root PROMPT.md, otherwise a pre-reset generic-file write recreates + planning output after the description-only publication. + */ + if (this.resetFence.isStale(task.id, planningGeneration)) return; const planPersistence = await persistPlanArtifact({ store: this.store, taskId: task.id, @@ -3383,7 +3490,20 @@ export class TriageProcessor { planningCwd, author: "triage", logger: { log: (m: string) => planLog.log(m), warn: (m: string) => planLog.warn(m) }, + writeAuthoritativePrompt: async (content) => await this.persistResetFencedPlanningArtifact( + task, + planningGeneration, + content, + false, + ), + mirrorAuthoritativePlan: async (content) => await this.persistResetFencedPlanningArtifact( + task, + planningGeneration, + content, + true, + ), }); + if (planPersistence.outcome === "recovery-fenced") return; if (planPersistence.outcome === "recovered") { await this.store.logEntry( task.id, @@ -3422,13 +3542,16 @@ export class TriageProcessor { const recoveredMarker = fusionCore.parseDuplicateMarkerFromSessionText(sessionTextTail); if (recoveredMarker) { const markerBody = `DUPLICATE: ${recoveredMarker.canonicalId}\n`; - const recovered = await writeFile(join(this.rootDir, promptPath), markerBody, "utf-8") - .then(() => true) - .catch((err: unknown) => { - const msg = err instanceof Error ? err.message : String(err); - planLog.warn(`${task.id}: failed to persist recovered duplicate marker: ${msg}`); - return false; - }); + const recovered = await this.persistResetFencedPlanningArtifact( + task, + planningGeneration, + markerBody, + false, + ).catch((err: unknown) => { + const msg = err instanceof Error ? err.message : String(err); + planLog.warn(`${task.id}: failed to persist recovered duplicate marker: ${msg}`); + return false; + }); if (recovered) { written = markerBody; planLog.log(`${task.id}: recovered duplicate verdict ${recoveredMarker.canonicalId} from the planner's reply (no PROMPT.md was written)`); @@ -3876,6 +3999,7 @@ export class TriageProcessor { this.options.onSpecifyError?.(task, err instanceof Error ? err : new Error(errorMessage)); } } finally { + this.activePlanningGenerations.delete(task.id); // FNXC:ConcurrencyAdmission 2026-08-03-10:00: a coordinator reservation // can exist before planner setup reaches takePreHeldExecutorSlot(). Every // early setup failure must return that untransferred host slot; after a @@ -4366,6 +4490,9 @@ export class TriageProcessor { after acquiring it so a dependency invalidation committed first fences this stale finalizer before it can restore approval or continuation handoff data. */ + // Reset owns this same non-reentrant lock; only after it releases may a finalizer enter, and its captured generation must then be fenced before any durable handoff. + const planningGeneration = this.activePlanningGenerations.get(task.id); + if (planningGeneration !== undefined && this.resetFence.isStale(task.id, planningGeneration)) return report; const reRead = await Promise.resolve(this.store.getTask(task.id)).catch(() => null); // Older pure unit-test adapters expose a no-op getTask; production returns // a Task or rejects. Preserve that fixture seam without treating a failed diff --git a/packages/engine/src/worktree/remove-reset-worktree.ts b/packages/engine/src/worktree/remove-reset-worktree.ts new file mode 100644 index 0000000000..907d05c84a --- /dev/null +++ b/packages/engine/src/worktree/remove-reset-worktree.ts @@ -0,0 +1,65 @@ +import type { Settings } from "@fusion/core"; +import { activeSessionRegistry, executingTaskLock, reconcileSelfOwnedActiveSessionForRemoval, DEFAULT_SELF_OWNED_MIN_IDLE_MS, type LiveBindingProbe } from "../agents/active-session-registry.js"; +import { isPlanningLive } from "../agents/planning-liveness.js"; +import { ActiveSessionWorktreeRemovalError, removeWorktree, RemovalReason, type WorktreeRemoveOutcome } from "./worktree-backend.js"; + +export interface RemoveTaskResetWorktreeInput { + worktreePath: string; + rootDir: string; + settings: Partial; + taskId: string; + audit?: Parameters[0]["audit"]; + liveOwnerProbe?: LiveBindingProbe; + remove?: (input: Parameters[0]) => Promise; + now?: () => number; + wait?: (ms: number) => Promise; +} + +export class ResetWorktreeForeignSessionError extends Error { + constructor(public readonly details: { worktreePath: string; holderTaskId: string; holderKind: string }) { + super(`worktree ${details.worktreePath} is held by ${details.holderTaskId} (${details.holderKind})`); + } +} + +/* +FNXC:TaskReset 2026-08-22-04:45: +The reset fence proves only registrations it released. A surviving self-owned entry still needs the normal idle and liveness gates: treating it as dead would let Reset remove a live planner worktree. +*/ +export async function removeTaskResetWorktree(input: RemoveTaskResetWorktreeInput): Promise { + const probe = input.liveOwnerProbe ?? ((_path: string, id: string) => executingTaskLock.has(id) || isPlanningLive(id)); + const processActiveProbe = (id: string) => executingTaskLock.has(id); + const reconcile = () => reconcileSelfOwnedActiveSessionForRemoval( + activeSessionRegistry, input.worktreePath, input.taskId, probe, { processActiveProbe, now: input.now }, + ); + const reconcileForRemoval = async () => { + let outcome = reconcile(); + if (outcome.action === "too-recent-refuses") { + const waitMs = Math.max(0, Math.min(DEFAULT_SELF_OWNED_MIN_IDLE_MS, outcome.minIdleMs - outcome.ageMs)); + await (input.wait ?? ((ms) => new Promise((resolve) => setTimeout(resolve, ms))))(waitMs); + outcome = reconcile(); + } + if (outcome.action === "foreign-task") { + const record = activeSessionRegistry.lookupByPath(input.worktreePath); + throw new ResetWorktreeForeignSessionError({ worktreePath: input.worktreePath, holderTaskId: outcome.ownerTaskId, holderKind: record?.kind ?? "unknown" }); + } + if (outcome.action === "live-binding-refuses" || outcome.action === "process-active-refuses" || outcome.action === "too-recent-refuses") { + const record = activeSessionRegistry.lookupByPath(input.worktreePath); + throw new ActiveSessionWorktreeRemovalError({ worktreePath: input.worktreePath, taskId: outcome.ownerTaskId, kind: record?.kind ?? "unknown", ownerKey: record?.ownerKey ?? "unknown", reason: RemovalReason.TaskReset }); + } + }; + await reconcileForRemoval(); + const remove = input.remove ?? removeWorktree; + const options = { + worktreePath: input.worktreePath, rootDir: input.rootDir, settings: input.settings, taskId: input.taskId, audit: input.audit, + reason: RemovalReason.TaskReset, expectedOwnerTaskId: input.taskId, liveOwnerProbe: probe, processActiveProbe, + }; + try { + return await remove(options); + } catch (error) { + if (!(error instanceof ActiveSessionWorktreeRemovalError) || error.details.taskId !== input.taskId) throw error; + // A new holder can appear after the first reconcile. Re-apply every normal gate before + // the sole retry; ignoring this result would turn a live or fresh registration into deletion. + await reconcileForRemoval(); + return await remove(options); + } +}