FN-151: make task reset always complete
Make reset a reliable description-only fresh start by fencing active planning work and cleaning discarded execution state. - Add reset lifecycle disposal for planning sessions, locks, artifacts, and reset worktrees. - Preserve operator-provided descriptions while clearing stale execution and planning projections. - Expose reset routes and add core, dashboard, and engine regression coverage. Files changed: .changeset/fn-151-task-reset-always-completes.md | 7 + docs/architecture.md | 3 + docs/dashboard-guide.md | 6 +- .../postgres/task-reset-publication.pg.test.ts | 71 +++++++++- .../core/src/__tests__/task-move-disposer.test.ts | 38 ++++++ packages/core/src/index.gate.ts | 2 + packages/core/src/index.ts | 2 + packages/core/src/task-store/reset-lifecycle.ts | 38 +++++- packages/core/src/tasks/task-move-disposer.ts | 36 +++++- .../src/__tests__/task-reset-lifecycle.test.ts | 110 +++++++++++++++- .../src/routes/register-task-workflow-routes.ts | 29 +++-- .../src/__tests__/agent-document-tools.test.ts | 49 +++++++ .../src/__tests__/reset-worktree-removal.test.ts | 77 +++++++++++ .../__tests__/triage-reset-planning-fence.test.ts | 118 +++++++++++++++++ packages/engine/src/agent-tools.ts | 60 +++++++-- packages/engine/src/agents/planning-liveness.ts | 19 +++ packages/engine/src/index.ts | 4 + packages/engine/src/plan-artifact-writeback.ts | 27 +++- packages/engine/src/planning-reset-fence.ts | 26 +++ packages/engine/src/triage.ts | 143 +++++++++++++++++++-- packages/engine/src/worktree/remove-reset-worktree.ts | 65 ++++++++++ 21 files changed, 882 insertions(+), 48 deletions(-) Fusion-Task-Id: FN-151 Fusion-Task-Lineage: 1c25d58e-2d23-4340-a489-51f3e2d8a7f4 Co-authored-by: Fusion <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-151-task-reset-always-completes.md
Normal file
7
.changeset/fn-151-task-reset-always-completes.md
Normal file
@@ -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.
|
||||
@@ -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.
|
||||
|
||||
<!-- FNXC:TaskReset 2026-08-22-18:15: Reset holds this non-reentrant lock while it cancels planner ownership and clears discarded output, so planner publications re-check their attempt generation inside the authoritative lock-held mutation. The reset disposer only aborts/disposes and unregisters its own planning paths; it never waits for a finalizer that needs this lock. -->
|
||||
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.
|
||||
|
||||
@@ -74,9 +74,11 @@ Both actions are irreversible; there is no undo after confirming. The dialog clo
|
||||
|
||||
<!-- FNXC:TaskResetDocs 2026-08-19-06:45: Task Reset is a destructive fresh-planning boundary. The guide must explain the cancellation fence, filesystem cleanup, retained history, atomic intake publication, and retry behavior for incomplete cleanup. -->
|
||||
|
||||
**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/<task-id>/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/<task-id>/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
|
||||
|
||||
|
||||
@@ -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` })]);
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<void>((resolve) => { releaseMove = resolve; }));
|
||||
const reset = vi.fn(() => new Promise<void>((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();
|
||||
});
|
||||
|
||||
@@ -643,7 +643,9 @@ export {
|
||||
disposeTaskBeforeMove,
|
||||
disposeTaskBeforeReset,
|
||||
getTaskMoveDisposer,
|
||||
getTaskResetDisposer,
|
||||
registerTaskMoveDisposer,
|
||||
registerTaskResetDisposer,
|
||||
type TaskMoveDisposer,
|
||||
type TaskResetDisposer,
|
||||
type TaskMoveDisposalInput,
|
||||
|
||||
@@ -773,7 +773,9 @@ export {
|
||||
disposeTaskBeforeMove,
|
||||
disposeTaskBeforeReset,
|
||||
getTaskMoveDisposer,
|
||||
getTaskResetDisposer,
|
||||
registerTaskMoveDisposer,
|
||||
registerTaskResetDisposer,
|
||||
type TaskMoveDisposer,
|
||||
type TaskResetDisposer,
|
||||
type TaskMoveDisposalInput,
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -20,6 +20,11 @@ export interface TaskMoveDisposalInput {
|
||||
* owned by another store. A set preserves every live owner during overlap.
|
||||
*/
|
||||
const disposers = new WeakMap<TaskStore, Set<TaskMoveDisposer>>();
|
||||
/*
|
||||
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<TaskStore, Set<TaskResetDisposer>>();
|
||||
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<TaskResetDisposer>();
|
||||
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<void> {
|
||||
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<typeof setTimeout> | undefined;
|
||||
try {
|
||||
await Promise.race([
|
||||
disposer(task),
|
||||
disposal,
|
||||
new Promise<void>((_resolve, reject) => {
|
||||
timeout = setTimeout(() => {
|
||||
reject(new Error(`Timed out stopping active work for ${task.id} before resetting the task`));
|
||||
|
||||
@@ -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<typeof actual.removeTaskResetWorktree>[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");
|
||||
|
||||
@@ -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}`);
|
||||
}
|
||||
|
||||
@@ -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<void>((resolve) => { releaseReset = resolve; });
|
||||
const queuedBehindReset = new Promise<void>((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<void>) => {
|
||||
enteredQueue();
|
||||
await resetPublication;
|
||||
await write();
|
||||
}),
|
||||
withTaskLock: vi.fn(async (_id: string, write: () => Promise<void>) => 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: "" });
|
||||
|
||||
77
packages/engine/src/__tests__/reset-worktree-removal.test.ts
Normal file
77
packages/engine/src/__tests__/reset-worktree-removal.test.ts
Normal file
@@ -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<Parameters<typeof removeTaskResetWorktree>[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(); }
|
||||
});
|
||||
});
|
||||
@@ -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<void>((resolve) => { releaseReset = resolve; });
|
||||
const entered = new Promise<void>((resolve) => { enteredQueue = resolve; });
|
||||
const updateTaskUnlocked = vi.fn();
|
||||
const store = {
|
||||
on: vi.fn(),
|
||||
off: vi.fn(),
|
||||
withPlanningLifecycleLock: vi.fn(async (_id: string, work: () => Promise<void>) => {
|
||||
enteredQueue();
|
||||
await resetPublication;
|
||||
await work();
|
||||
}),
|
||||
withTaskLock: vi.fn(async (_id: string, work: () => Promise<unknown>) => 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<boolean>;
|
||||
}).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<never>(() => 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<string, unknown> }).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<string, unknown> }).activeSessions.has(task.id)).toBe(false);
|
||||
} finally {
|
||||
processor.stop();
|
||||
activeSessionRegistry.unregisterPath(worktree);
|
||||
activeSessionRegistry.unregisterPath(`${worktree}-executor`);
|
||||
activeSessionRegistry.unregisterPath(`${worktree}-foreign`);
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -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: {},
|
||||
|
||||
19
packages/engine/src/agents/planning-liveness.ts
Normal file
19
packages/engine/src/agents/planning-liveness.ts
Normal file
@@ -0,0 +1,19 @@
|
||||
export type PlanningLivenessProbe = (taskId: string) => boolean;
|
||||
const probes = new Set<PlanningLivenessProbe>();
|
||||
/*
|
||||
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);
|
||||
@@ -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,
|
||||
|
||||
@@ -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<boolean>;
|
||||
/** Optional caller-owned plan-document writer subject to the same publication fence. */
|
||||
mirrorAuthoritativePlan?: (content: string, author?: string) => Promise<boolean>;
|
||||
}
|
||||
|
||||
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<ReconcileWorktreePlanArtifactResult & { mirrored: boolean }> {
|
||||
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 };
|
||||
}
|
||||
|
||||
26
packages/engine/src/planning-reset-fence.ts
Normal file
26
packages/engine/src/planning-reset-fence.ts
Normal file
@@ -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<string, number>();
|
||||
private readonly holds = new Map<string, number>();
|
||||
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); }
|
||||
}
|
||||
@@ -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<string, { dispose: () => void }>();
|
||||
private readonly resetFence = new PlanningResetFence();
|
||||
/** Captured attempt generations let a delayed finalizer reject a reset-fenced planner. */
|
||||
private readonly activePlanningGenerations = new Map<string, number>();
|
||||
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<boolean> {
|
||||
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<void>; 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<void> {
|
||||
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
|
||||
|
||||
65
packages/engine/src/worktree/remove-reset-worktree.ts
Normal file
65
packages/engine/src/worktree/remove-reset-worktree.ts
Normal file
@@ -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<Settings>;
|
||||
taskId: string;
|
||||
audit?: Parameters<typeof removeWorktree>[0]["audit"];
|
||||
liveOwnerProbe?: LiveBindingProbe;
|
||||
remove?: (input: Parameters<typeof removeWorktree>[0]) => Promise<WorktreeRemoveOutcome>;
|
||||
now?: () => number;
|
||||
wait?: (ms: number) => Promise<void>;
|
||||
}
|
||||
|
||||
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<WorktreeRemoveOutcome> {
|
||||
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<void>((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);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user