diff --git a/.changeset/quiet-tasks-stop-before-todo.md b/.changeset/quiet-tasks-stop-before-todo.md new file mode 100644 index 0000000000..e984a302d2 --- /dev/null +++ b/.changeset/quiet-tasks-stop-before-todo.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Stop active task processing before a user move to Todo becomes visible. +category: fix +dev: User-driven in-progress-to-Todo transitions now await every executor cancellation surface before persistence. diff --git a/packages/core/src/__tests__/postgres/task-move-hard-cancel-ordering.pg.test.ts b/packages/core/src/__tests__/postgres/task-move-hard-cancel-ordering.pg.test.ts new file mode 100644 index 0000000000..33b53aba4d --- /dev/null +++ b/packages/core/src/__tests__/postgres/task-move-hard-cancel-ordering.pg.test.ts @@ -0,0 +1,82 @@ +import { afterAll, afterEach, beforeAll, beforeEach, expect, it, vi } from "vitest"; +import { + __setTaskMoveDisposalTimeoutForTesting, + registerTaskMoveDisposer, +} from "../../task-move-disposer.js"; +import { readTaskRow } from "../../task-store/async-persistence.js"; +import { + createSharedPgTaskStoreTestHarness, + pgDescribe, +} from "../../__test-utils__/pg-test-harness.js"; + +/* +FNXC:TaskMovement 2026-07-18-14:32: +Surface enumeration for the hard-cancel invariant: + - A user move from in-progress to Todo waits for executor cancellation before persistence. + - The durable task stays in-progress throughout a delayed cancellation. + - A wedged cancellation times out fail-closed and releases the per-task lock. + - Engine moves, forward moves, and other destinations do not invoke the user-cancel seam + (covered by task-move-disposer.test.ts). + - Main, step, workflow, configured-command, subagent, and CLI surfaces share the executor's + awaitAbortInFlightTaskWork path (covered by executor-user-cancel.test.ts). +*/ +pgDescribe("user move to Todo hard-cancel ordering", () => { + const harness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_task_move_cancel" }); + beforeAll(harness.beforeAll); + beforeEach(harness.beforeEach); + afterEach(harness.afterEach); + afterAll(harness.afterAll); + + it("keeps the durable task in-progress until cancellation finishes", async () => { + const store = harness.store(); + const created = await store.createTask({ description: "Stop before returning to Todo" }); + await store.moveTask(created.id, "todo", { moveSource: "engine" }); + await store.moveTask(created.id, "in-progress", { moveSource: "scheduler" }); + + let resolveCancellation: (() => void) | undefined; + const cancellation = new Promise((resolve) => { + resolveCancellation = resolve; + }); + const disposer = vi.fn(() => cancellation); + const unregister = registerTaskMoveDisposer(store, disposer); + + try { + const move = store.moveTask(created.id, "todo", { moveSource: "user" }); + await vi.waitFor(() => expect(disposer).toHaveBeenCalledOnce()); + + expect((await readTaskRow(store.asyncLayer!, created.id))?.column).toBe("in-progress"); + + resolveCancellation?.(); + await expect(move).resolves.toMatchObject({ column: "todo", userPaused: true }); + expect((await store.getTask(created.id)).column).toBe("todo"); + } finally { + resolveCancellation?.(); + unregister(); + } + }); + + it("keeps the task in-progress but releases its lock when cancellation times out", async () => { + const store = harness.store(); + const created = await store.createTask({ description: "Release a wedged hard cancel" }); + await store.moveTask(created.id, "todo", { moveSource: "engine" }); + await store.moveTask(created.id, "in-progress", { moveSource: "scheduler" }); + const unregister = registerTaskMoveDisposer(store, () => new Promise(() => {})); + + __setTaskMoveDisposalTimeoutForTesting(1); + try { + await expect( + store.moveTask(created.id, "todo", { moveSource: "user" }), + ).rejects.toThrow( + `Timed out stopping active work for ${created.id} before moving to Todo`, + ); + } finally { + __setTaskMoveDisposalTimeoutForTesting(); + unregister(); + } + + expect((await readTaskRow(store.asyncLayer!, created.id))?.column).toBe("in-progress"); + await expect( + store.moveTask(created.id, "todo", { moveSource: "engine" }), + ).resolves.toMatchObject({ column: "todo" }); + }); +}); diff --git a/packages/core/src/__tests__/task-move-disposer.test.ts b/packages/core/src/__tests__/task-move-disposer.test.ts new file mode 100644 index 0000000000..367448f548 --- /dev/null +++ b/packages/core/src/__tests__/task-move-disposer.test.ts @@ -0,0 +1,101 @@ +import { describe, expect, it, vi } from "vitest"; +import { + __setTaskMoveDisposalTimeoutForTesting, + disposeTaskBeforeMove, + registerTaskMoveDisposer, +} from "../task-move-disposer.js"; + +describe("task move disposer", () => { + it("does not complete a user in-progress to todo move until cancellation settles", async () => { + const store = {} as never; + let resolveCancellation: (() => void) | undefined; + const cancellation = new Promise((resolve) => { + resolveCancellation = resolve; + }); + const disposer = vi.fn(() => cancellation); + registerTaskMoveDisposer(store, disposer); + + let moveReady = false; + const preparation = disposeTaskBeforeMove(store, { + task: { id: "FN-CANCEL" } as never, + from: "in-progress", + to: "todo", + source: "user", + }).then(() => { + moveReady = true; + }); + + await Promise.resolve(); + expect(disposer).toHaveBeenCalledOnce(); + expect(moveReady).toBe(false); + + resolveCancellation?.(); + await preparation; + expect(moveReady).toBe(true); + }); + + it("awaits every executor registered to the same store", async () => { + const store = {} as never; + const first = vi.fn().mockResolvedValue(undefined); + const second = vi.fn().mockResolvedValue(undefined); + const unregisterFirst = registerTaskMoveDisposer(store, first); + registerTaskMoveDisposer(store, second); + + await disposeTaskBeforeMove(store, { + task: { id: "FN-MULTI-OWNER" } as never, + from: "in-progress", + to: "todo", + source: "user", + }); + + expect(first).toHaveBeenCalledOnce(); + expect(second).toHaveBeenCalledOnce(); + + unregisterFirst(); + await disposeTaskBeforeMove(store, { + task: { id: "FN-ONE-OWNER" } as never, + from: "in-progress", + to: "todo", + source: "user", + }); + expect(first).toHaveBeenCalledOnce(); + expect(second).toHaveBeenCalledTimes(2); + }); + + it("fails closed and releases the move when cancellation does not settle", async () => { + __setTaskMoveDisposalTimeoutForTesting(1); + try { + const store = {} as never; + registerTaskMoveDisposer(store, () => new Promise(() => {})); + + const preparation = disposeTaskBeforeMove(store, { + task: { id: "FN-WEDGED" } as never, + from: "in-progress", + to: "todo", + source: "user", + }); + await expect(preparation).rejects.toThrow( + "Timed out stopping active work for FN-WEDGED before moving to Todo", + ); + } finally { + __setTaskMoveDisposalTimeoutForTesting(); + } + }); + + it.each([ + { from: "in-progress", to: "todo", source: "engine" }, + { from: "todo", to: "in-progress", source: "user" }, + { from: "in-progress", to: "in-review", source: "user" }, + ] as const)("does not cancel for $source $from to $to moves", async (move) => { + const store = {} as never; + const disposer = vi.fn(); + registerTaskMoveDisposer(store, disposer); + + await disposeTaskBeforeMove(store, { + task: { id: "FN-UNCHANGED" } as never, + ...move, + }); + + expect(disposer).not.toHaveBeenCalled(); + }); +}); diff --git a/packages/core/src/index.gate.ts b/packages/core/src/index.gate.ts index abbb734a1d..65685fc48b 100644 --- a/packages/core/src/index.gate.ts +++ b/packages/core/src/index.gate.ts @@ -541,6 +541,14 @@ export { type WorkspaceDisposalPlanEntry, type ArchiveWorkspaceDisposalResult, } from "./archive-worktree-disposer.js"; +export { + disposeTaskBeforeMove, + getTaskMoveDisposer, + registerTaskMoveDisposer, + type TaskMoveDisposer, + type TaskMoveDisposalInput, + type TaskMoveSource, +} from "./task-move-disposer.js"; export { acquireWorktreePathReservation, withWorktreePathReservation, diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 120098f1c2..647e8d2761 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -556,6 +556,14 @@ export { type WorkspaceDisposalPlanEntry, type ArchiveWorkspaceDisposalResult, } from "./archive-worktree-disposer.js"; +export { + disposeTaskBeforeMove, + getTaskMoveDisposer, + registerTaskMoveDisposer, + type TaskMoveDisposer, + type TaskMoveDisposalInput, + type TaskMoveSource, +} from "./task-move-disposer.js"; export { acquireWorktreePathReservation, withWorktreePathReservation, diff --git a/packages/core/src/task-move-disposer.ts b/packages/core/src/task-move-disposer.ts new file mode 100644 index 0000000000..60777b9aac --- /dev/null +++ b/packages/core/src/task-move-disposer.ts @@ -0,0 +1,73 @@ +import type { TaskStore } from "./store.js"; +import type { ColumnId, Task } from "./types.js"; + +export type TaskMoveSource = "user" | "engine" | "scheduler"; +export type TaskMoveDisposer = (task: Task) => Promise; + +export interface TaskMoveDisposalInput { + task: Task; + from: ColumnId; + to: ColumnId; + source: TaskMoveSource; +} + +/* + * Core owns task-transition ordering but cannot import the engine. Keep the + * cancellation seam store-scoped so one project's executor cannot stop work + * owned by another store. A set preserves every live owner during overlap. + */ +const disposers = new WeakMap>(); +const TASK_MOVE_DISPOSAL_TIMEOUT_MS = 30_000; +let taskMoveDisposalTimeoutMs = TASK_MOVE_DISPOSAL_TIMEOUT_MS; + +export function __setTaskMoveDisposalTimeoutForTesting( + timeoutMs = TASK_MOVE_DISPOSAL_TIMEOUT_MS, +): void { + taskMoveDisposalTimeoutMs = timeoutMs; +} + +export function registerTaskMoveDisposer(store: TaskStore, disposer: TaskMoveDisposer): () => void { + const registered = disposers.get(store) ?? new Set(); + registered.add(disposer); + disposers.set(store, registered); + return () => { + const current = disposers.get(store); + current?.delete(disposer); + if (current?.size === 0) disposers.delete(store); + }; +} + +export function getTaskMoveDisposer(store: TaskStore): TaskMoveDisposer | undefined { + const registered = disposers.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 + * registered execution surface before publishing the new column so persisted + * board state can never claim the task is idle while its agent still runs. + */ +export async function disposeTaskBeforeMove(store: TaskStore, input: TaskMoveDisposalInput): Promise { + if (input.source !== "user" || input.from !== "in-progress" || input.to !== "todo") return; + const disposer = getTaskMoveDisposer(store); + if (!disposer) return; + + let timeout: ReturnType | undefined; + try { + await Promise.race([ + disposer(input.task), + new Promise((_resolve, reject) => { + timeout = setTimeout(() => { + reject(new Error(`Timed out stopping active work for ${input.task.id} before moving to Todo`)); + }, taskMoveDisposalTimeoutMs); + timeout.unref?.(); + }), + ]); + } finally { + if (timeout) clearTimeout(timeout); + } +} diff --git a/packages/core/src/task-store/moves.ts b/packages/core/src/task-store/moves.ts index beec1bdb8d..868ed97599 100644 --- a/packages/core/src/task-store/moves.ts +++ b/packages/core/src/task-store/moves.ts @@ -30,6 +30,7 @@ import {recordRunAuditEventWithinTransaction} from "../postgres/data-layer.js"; import {getTaskMergeBlocker} from "../task-merge.js"; import {__setTaskActivityLogLimitsForTesting} from "../task-store/comments.js"; import {readTaskRow as readTaskRowAsync, readTaskRowInTransaction, upsertTaskRowInTransaction} from "../task-store/async-persistence.js"; +import {disposeTaskBeforeMove} from "../task-move-disposer.js"; /* FNXC:PostgresCutover 2026-07-05-19:50: @@ -536,6 +537,21 @@ export async function moveTaskInternalImpl(store: TaskStore, id: string, toColum } } + /* + FNXC:TaskMovement 2026-07-18-14:32: + A user in-progress -> todo move is a hard cancel. Run the engine-owned, + store-scoped disposer after every transition guard passes but before the + task object or durable row changes column. The dashboard cannot observe a + Todo row until the agent, workflow, configured command, and CLI execution + surfaces have stopped. + */ + await disposeTaskBeforeMove(store, { + task, + from: fromColumn, + to: toColumn, + source: moveSource, + }); + const movedAt = internal.now ?? new Date().toISOString(); task.column = toColumn; task.columnMovedAt = movedAt; diff --git a/packages/engine/src/__tests__/executor-user-cancel.test.ts b/packages/engine/src/__tests__/executor-user-cancel.test.ts index 8cbbba0297..a86d5c28b0 100644 --- a/packages/engine/src/__tests__/executor-user-cancel.test.ts +++ b/packages/engine/src/__tests__/executor-user-cancel.test.ts @@ -1,9 +1,97 @@ import { describe, it, expect, vi } from "vitest"; import "./executor-test-helpers.js"; +import { getTaskMoveDisposer } from "@fusion/core"; import { TaskExecutor } from "../executor.js"; import { createMockStore, resetExecutorMocks } from "./executor-test-helpers.js"; describe("TaskExecutor user cancel handling", () => { + /* + FNXC:WorkflowLifecycle 2026-07-18-14:32: + A user move from active execution to Todo must await every executor + cancellation surface so the card cannot keep processing in the background. + */ + it("registers an awaited user-move disposer that aborts all active work before Todo", async () => { + resetExecutorMocks(); + const store = createMockStore(); + const executor = new TaskExecutor(store as any, "/tmp/test"); + const terminateChildren = vi.spyOn(executor as any, "terminateAllChildren").mockResolvedValue(undefined); + let resolveAbort: (() => void) | undefined; + const abortPending = new Promise((resolve) => { + resolveAbort = resolve; + }); + const session = { + prompt: vi.fn(), + abort: vi.fn(() => abortPending), + dispose: vi.fn(), + } as any; + (executor as any).activeSessions.set("FN-AWAITED", { + session, + seenSteeringIds: new Set(), + }); + + const disposer = getTaskMoveDisposer(store as any); + expect(disposer).toBeTypeOf("function"); + let disposed = false; + const disposal = disposer!({ id: "FN-AWAITED" } as any).then(() => { + disposed = true; + }); + + await Promise.resolve(); + expect(terminateChildren).toHaveBeenCalledWith("FN-AWAITED"); + expect(session.abort).toHaveBeenCalledOnce(); + expect(disposed).toBe(false); + + resolveAbort?.(); + await disposal; + expect(session.dispose).toHaveBeenCalledOnce(); + expect((executor as any).userCanceledTaskIds.has("FN-AWAITED")).toBe(true); + }); + + it("does not let late cancellation cleanup touch a replacement execution", async () => { + resetExecutorMocks(); + let resolveChildStateUpdate: (() => void) | undefined; + const childStateUpdate = new Promise((resolve) => { + resolveChildStateUpdate = resolve; + }); + const agentStore = { + updateAgentState: vi.fn(() => childStateUpdate), + deleteAgent: vi.fn().mockResolvedValue(undefined), + }; + const store = createMockStore(); + const executor = new TaskExecutor(store as any, "/tmp/test", { agentStore } as any); + const oldSession = { + prompt: vi.fn(), + abort: vi.fn().mockResolvedValue(undefined), + dispose: vi.fn(), + } as any; + const oldChildSession = { dispose: vi.fn() } as any; + (executor as any).activeSessions.set("FN-GENERATION", { + session: oldSession, + seenSteeringIds: new Set(), + }); + (executor as any).spawnedAgents.set("FN-GENERATION", new Set(["old-child"])); + (executor as any).childSessions.set("old-child", oldChildSession); + + const disposer = getTaskMoveDisposer(store as any)!; + const disposal = disposer({ id: "FN-GENERATION" } as any); + + const replacementSession = { dispose: vi.fn() } as any; + const replacementChildren = new Set(["new-child"]); + (executor as any).activeSessions.set("FN-GENERATION", { + session: replacementSession, + seenSteeringIds: new Set(), + }); + (executor as any).spawnedAgents.set("FN-GENERATION", replacementChildren); + + resolveChildStateUpdate?.(); + await disposal; + + expect(oldSession.abort).toHaveBeenCalledOnce(); + expect(oldChildSession.dispose).toHaveBeenCalledOnce(); + expect((executor as any).activeSessions.get("FN-GENERATION")?.session).toBe(replacementSession); + expect((executor as any).spawnedAgents.get("FN-GENERATION")).toBe(replacementChildren); + }); + it("aborts before dispose when user moves in-progress task back to todo", async () => { resetExecutorMocks(); const store = createMockStore(); diff --git a/packages/engine/src/executor.ts b/packages/engine/src/executor.ts index 1ef1f8ff4f..d2b0377425 100644 --- a/packages/engine/src/executor.ts +++ b/packages/engine/src/executor.ts @@ -134,7 +134,7 @@ import { deriveRepoScopeSubset, normalizeRepoRelPath } from "./workspace-paths.j import { RemovalReason, classifyTaskWorktree, describeRegisteredWorktrees, detectGitRepository, detectNestedWorktreeRoot, getRegisteredWorktreePaths, isInsideWorktreesDir, isRegisteredGitWorktree, removeWorktree, type GitRepoDetection, type WorktreePool } from "./worktree-pool.js"; import { attemptBranchAutocorrect } from "./branch-autocorrect.js"; import { ActiveSessionWorktreeRemovalError } from "./worktree-backend.js"; -import {canonicalizeWorktreePath, registerArchiveWorkspaceWorktreeDisposer, registerArchiveWorktreeDisposer} from "@fusion/core"; +import {canonicalizeWorktreePath, registerArchiveWorkspaceWorktreeDisposer, registerArchiveWorktreeDisposer, registerTaskMoveDisposer} from "@fusion/core"; import { activeSessionRegistry, executingTaskLock, @@ -1749,6 +1749,7 @@ export class TaskExecutor { * this so a fast re-dispatch (task:moved → in-progress) awaits the prior * session being fully reaped before creating/acquiring a new worktree. */ private pendingTaskDisposals = new Map>(); + private unregisterTaskMoveDisposer: (() => void) | undefined; private unregisterArchiveWorktreeDisposer: (() => void) | undefined; private unregisterArchiveWorkspaceWorktreeDisposer: (() => void) | undefined; /** Active agent sessions per task, used to terminate on pause and inject steering. */ @@ -3026,6 +3027,16 @@ export class TaskExecutor { private options: TaskExecutorOptions = {}, ) { executorLog.log(`TaskExecutor constructed (rootDir=${rootDir}, hasSemaphore=${!!options.semaphore}, hasStuckDetector=${!!options.stuckTaskDetector})`); + this.unregisterTaskMoveDisposer = registerTaskMoveDisposer(store, async (task) => { + // Start both paths without awaiting between them. Each synchronously + // detaches its current targets before its first await, fencing late + // cleanup from a replacement execution after the move timeout expires. + const children = this.terminateAllChildren(task.id); + const activeWork = this.awaitAbortInFlightTaskWork(task.id, "user moved task from in-progress to todo", { + userCanceled: true, + }); + await Promise.all([children, activeWork]); + }); /* FNXC:WorkflowLifecycle 2026-07-16-10:00: Executor replaces the baseline only for its own TaskStore, so archive awaits abort/sweep/removal before branch deletion without cross-store coupling. */ this.unregisterArchiveWorktreeDisposer = registerArchiveWorktreeDisposer(store, async (task) => { if (!task.worktree || await canonicalizeWorktreePath(task.worktree) === await canonicalizeWorktreePath(this.rootDir)) return; @@ -18500,8 +18511,10 @@ You have access to the file system to review changes.${inlineFixBlock}${verdictB } } - /** Remove only this executor's store-scoped archive disposer registration. */ - disposeArchiveWorktreeDisposer(): void { + /** Remove only this executor's store-scoped lifecycle disposer registrations. */ + disposeStoreLifecycleDisposers(): void { + this.unregisterTaskMoveDisposer?.(); + this.unregisterTaskMoveDisposer = undefined; this.unregisterArchiveWorktreeDisposer?.(); this.unregisterArchiveWorktreeDisposer = undefined; this.unregisterArchiveWorkspaceWorktreeDisposer?.(); @@ -19444,11 +19457,11 @@ You have access to the file system to review changes.${inlineFixBlock}${verdictB if (!childIds || childIds.size === 0) return; executorLog.log(`Terminating ${childIds.size} child agents for parent ${parentTaskId}`); - - for (const childId of childIds) { - await this.terminateChildAgent(childId); - } + // Detach the parent generation before any agent-store await. A replacement + // execution may register a new set for the same task ID while cleanup is + // still settling; the old generation must never delete that new set. this.spawnedAgents.delete(parentTaskId); + await Promise.all([...childIds].map((childId) => this.terminateChildAgent(childId))); } /** diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index 074ee8b1de..9d228fb588 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -1394,7 +1394,7 @@ export class InProcessRuntime if (this.executor) { try { await this.executor.abortAllInFlight("engine stop"); - this.executor.disposeArchiveWorktreeDisposer(); + this.executor.disposeStoreLifecycleDisposers(); runtimeLog.log("Aborted in-flight executor AI sessions"); } catch (err) { runtimeLog.warn(`Failed to abort in-flight executor AI sessions: ${err}`);