From 0440ae4bb94db07e570a6c77cf6418082e7da639 Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Fri, 26 Jun 2026 14:55:01 -0700 Subject: [PATCH] FN-7074: make task ID reservation commits atomic Task creation now commits or rolls back distributed task ID reservations with the task-row transaction. - Add transaction-participating reservation commit and rollback helpers. - Wire create, duplicate, and refinement task paths to commit reservations inside task insertion. - Record rollback audit events and preserve burned reservation rows after failed creates. - Cover atomicity, rollback, and allocator behavior with reservation-focused tests and docs. Files changed: .changeset/fn-7074-reservation-atomicity.md | 7 + AGENTS.md | 1 + docs/architecture.md | 7 +- docs/storage.md | 4 +- .../core/src/__tests__/distributed-task-id.test.ts | 25 ++- .../__tests__/store-reservation-atomicity.test.ts | 224 +++++++++++++++++++++ packages/core/src/distributed-task-id.ts | 208 +++++++++++++------ packages/core/src/store.ts | 93 +++++++-- 8 files changed, 479 insertions(+), 90 deletions(-) Fusion-Task-Id: FN-7074 Fusion-Task-Lineage: 464a98cf-e903-4257-93ac-424c7129412b --- .changeset/fn-7074-reservation-atomicity.md | 7 + AGENTS.md | 1 + docs/architecture.md | 7 +- docs/storage.md | 4 +- .../src/__tests__/distributed-task-id.test.ts | 25 +- .../store-reservation-atomicity.test.ts | 224 ++++++++++++++++++ packages/core/src/distributed-task-id.ts | 208 +++++++++++----- packages/core/src/store.ts | 93 ++++++-- 8 files changed, 479 insertions(+), 90 deletions(-) create mode 100644 .changeset/fn-7074-reservation-atomicity.md create mode 100644 packages/core/src/__tests__/store-reservation-atomicity.test.ts diff --git a/.changeset/fn-7074-reservation-atomicity.md b/.changeset/fn-7074-reservation-atomicity.md new file mode 100644 index 0000000000..dd16f7ff13 --- /dev/null +++ b/.changeset/fn-7074-reservation-atomicity.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": patch +--- + +summary: Task creation no longer leaves orphaned reserved-ID records when a create fails partway. +category: fix +dev: createTaskWithDistributedReservation now commits the distributed_task_id_reservations row in the same SQLite transaction as the tasks-row insert, and a rollback guard reverts both the row and the reservation if post-insert task.json/PROMPT.md materialization or create validation fails, preventing committed-reservation-without-task phantoms. Adds transaction-participating allocator helpers for commit and failed-create rollback. diff --git a/AGENTS.md b/AGENTS.md index 6185d3bd75..0ac9a7d7ce 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -225,6 +225,7 @@ Scoped exception (FN-5819): shared-branch-group members (`branchContext.assignme - FN-6736: self-healing emits `task:reclaim-phantom-executor-binding` when it proves an in-memory executor-active binding is stale, clears the binding, and requeues the in-progress task with worktree/progress preserved. - FN-6783: task-store open and self-healing housekeeping emit `task:reconcile-orphaned-task-dir` when they non-destructively re-import a valid live `.fusion/tasks/{ID}/task.json` directory that has no task row anywhere, preserving soft-deleted/archived/tombstoned IDs. - FN-7069: task-store open and self-healing housekeeping emit `task:reconcile-phantom-committed-reservation` when they prune orphaned child rows for a committed task-ID reservation that has no live/soft-deleted/archived task row and no task directory, while preserving the committed reservation so the ID is never reused. +- FN-7074: task creation emits `task:reservation-commit-rolled-back` when a distributed reservation was committed atomically with a `tasks` row but a later create materialization step failed; metadata includes `reservationId`, `nodeId`, `reason: "failed-create"`, and `error`, and the reservation is moved to aborted so the sequence remains burned. - FN-6782/FN-6796: self-healing emits `task:auto-recover-paused-abort-park` when it clears a benign pause-abort operator park, requeueing safe `todo`/`in-progress` rows or preserving a clean auto-merge-eligible `in-review` row for review progression. - FN-6793/FN-6797: self-healing emits `task:reconcile-in-review-unmet-dependencies` when it rebounds an `in-review` task whose declared dependencies are still unmet, and `task:reconcile-in-review-unmet-dependencies-no-action` when pause/user-pause, `autoMerge:false`, live execution/checkout proof, or a failed rebound mutation blocks that backward move. - Workspace (Phase D U1): self-healing emits `task:reconcile-workspace-partial-land` when it re-enqueues a partial/zero-landed workspace task's per-repo land (or parks it `failed` when a sub-repo's `fusion/` branch is gone with no `landedSha`), and `task:reconcile-workspace-partial-land-no-action` when `autoMerge:false`, user-pause, or a live sub-repo worktree (workspace-aware liveness) blocks that backward move. diff --git a/docs/architecture.md b/docs/architecture.md index d1745fde1e..d94cf11240 100644 --- a/docs/architecture.md +++ b/docs/architecture.md @@ -736,6 +736,7 @@ Guardrails: this routine does **not** retry merges, does **not** apply to mixed/ - `RunAudit` (`run-audit.ts`) — mutation audit tracking (DB/git/filesystem) - FN-6782/FN-6796: `task:auto-recover-paused-abort-park` records self-healing recovery of pause-abort operator parks. Metadata includes the source column and whether recovery preserved a clean `in-review` row instead of requeueing to `todo`. - FN-7069: `task:reconcile-phantom-committed-reservation` records task-store startup or self-healing cleanup of committed-reservation-without-task phantoms. Metadata includes `reservationStatus: "committed"` plus pruned `activityLog` and `agents` counts; `runAuditEvents` and the committed reservation are intentionally retained for auditability and ID permanence. + - FN-7074: `task:reservation-commit-rolled-back` records preventive create-path rollback when a distributed reservation was committed with the task-row insert but a later create materialization step failed. Metadata includes `{ reservationId, nodeId, reason: "failed-create", error }`; the task row/partial directory are removed and the reservation is moved to `aborted` so FN-7069 should not need to clean up a new phantom. - FN-4956: Layer 3 merge-conflict arbitration now scope-partitions conflicted files before AI resolution. Out-of-scope conflicts are deterministically resolved to the integration branch (`git checkout --ours`) and unstaged, while only in-scope conflicts flow to AI. Integration branch defaults are resolved via `resolveIntegrationBranch(rootDir, settings)`. Audit events: `merge:layer3:foreign-file-skipped` and `merge:layer3:scope-override-bypass`. - FN-5655 goal anchoring observability adds `database`-domain mutation types `goal:injection-applied`, `goal:injection-skipped`, and `goal:retrieval-invoked` so Slice 2 cite-rate tracking has a prompt-independent signal. Metadata uses counts/IDs only (`count`, `lane`, `toolName`, optional `truncated`/`reason`/`notFound`) and never stores prompt bodies or goal titles/descriptions. These events surface through `GET /api/agents/:id/runs/:runId/audit` and support the existing `startTime`/`endTime` filters. @@ -797,13 +798,13 @@ A lease is recoverable only when there is **no active local executor session for - Existing `/api/mesh/sync` and settings-sync payloads remain the active exchange primitives while follow-on runtime tasks implement full v1 coordinator/quorum behavior. - Distributed task-ID allocation (`packages/core/src/distributed-task-id.ts`) is the first mesh-aware coordinated write primitive. - Durable state lives in SQLite tables `distributed_task_id_state` (prefix sequence + authoritative committed count) and `distributed_task_id_reservations` (reservation lifecycle rows). - - Reserve/commit/abort execute under a process-local lock and a single SQLite transaction. Lazy reservation expiry cleanup runs inside those same transactions. - - Default reservation TTL is `15 * 60 * 1000` ms (15 minutes). Expired/aborted reservations are **burned IDs** and are never reissued. + - Reserve/commit/abort execute under a process-local lock and a single SQLite transaction. Lazy reservation expiry cleanup runs inside those same transactions. `TaskStore` also uses a non-locking commit core inside its own `BEGIN IMMEDIATE` create transaction so the reservation `committed` flip and authoritative `tasks` row insert share one SQLite durability point. + - Default reservation TTL is `15 * 60 * 1000` ms (15 minutes). Expired/aborted reservations are **burned IDs** and are never reissued. If a post-insert create step fails after the reservation was committed (for example `task.json`/`PROMPT.md` disk materialization, file-scope validation, or duplicate-intake tombstone checks), the failed-create rollback deletes the just-created task row/partial directory, moves the reservation to `aborted`, recomputes committed reservation counters, and emits `task:reservation-commit-rolled-back`; the sequence stays burned for FN-5105 ID permanence. - `committedClusterTaskCount` from allocator state is the only authoritative cluster-wide committed-task count. Local task-row counts and ID suffix math are not authoritative. - Store open reconciles every known prefix in `distributed_task_id_state` to `max(current nextSequence, max(tasks suffix)+1, max(archivedTasks suffix)+1, max(reservation sequence)+1)`. This self-heals stale counters before ordinary task creation resumes. - Mesh allocator write routes (`/api/mesh/task-ids/reserve|commit|abort`) return `503` when the coordinator node is unreachable; they never fall back to local-only cluster ID issuance. - Cluster task creation now uses a strong-write reserve → create → replicate → commit/abort sequence. - - Ordinary local task creation (`TaskStore.createTask()`, duplicate, and refine flows) now allocates IDs through the same distributed reserve/commit/abort lifecycle owned by `TaskStore`. + - Ordinary local task creation (`TaskStore.createTask()`, duplicate, and refine flows) now allocates IDs through the same distributed reserve/commit/abort lifecycle owned by `TaskStore`; the invariant is `distributed_task_id_reservations.status = 'committed'` iff a live durable `tasks` row and task directory landed for that ID. `applyReplicatedTaskCreate(...)` remains a direct reserved-ID apply path and does not require a local reservation row. - `POST /api/tasks` uses the store-owned allocator path for local creates rather than maintaining a separate route-local allocator implementation. - `POST /api/tasks` reserves a distributed ID, creates the authoritative local task with that reserved ID, then POSTs authenticated replication payloads to peer nodes. - All create-class writes now use conflict-raising inserts, not SQLite `ON CONFLICT ... DO UPDATE`. Existing task rows and `.fusion/tasks/{id}` contents always win over stale counters or colliding reservations. diff --git a/docs/storage.md b/docs/storage.md index 4d062b9955..dbef183f75 100644 --- a/docs/storage.md +++ b/docs/storage.md @@ -3,10 +3,10 @@ ## Task-ID allocator authority and compatibility - `distributed_task_id_state` is the authoritative local task-ID allocator state. `nextSequence` is the active high-water mark used for local ID reservations. -- `distributed_task_id_reservations` tracks reserve/commit/abort lifecycle entries. Aborted/expired reservations are burned and never reissued. +- `distributed_task_id_reservations` tracks reserve/commit/abort lifecycle entries. Aborted/expired reservations are burned and never reissued. Create-class writes commit the reservation in the same SQLite transaction as the `tasks` row insert, then roll back the row/partial directory and move the reservation to aborted if post-insert `task.json`/`PROMPT.md` materialization or create validation fails. - `config.nextId` is retained only as a deprecated legacy compatibility field and optional one-time seed source. Fusion still reads it during reconciliation, but runtime task creation and settings writes no longer mutate it. - Startup/store-open allocator reconciliation bumps each active prefix sequence to `max(current nextSequence, max(tasks suffix)+1, max(archivedTasks suffix)+1, max(reservation sequence)+1)` so stale allocator rows self-heal before local task creation resumes. -- Create-class task persistence is intentionally non-destructive: new tasks use plain `INSERT` semantics, while `ON CONFLICT(id) DO UPDATE` remains update-only. If counters drift and a reserved ID still collides, the create fails and the existing SQLite row / task directory stays intact. +- Create-class task persistence is intentionally non-destructive: new tasks use plain `INSERT` semantics, while `ON CONFLICT(id) DO UPDATE` remains update-only. If counters drift and a reserved ID still collides, the create fails and the existing SQLite row / task directory stays intact. A `committed` distributed reservation is valid only with a matching durable task row/directory; failed creates burn the reservation as `aborted` instead of leaving a committed-reservation-without-task phantom. ## Soft-deleted tasks (FN-5105) diff --git a/packages/core/src/__tests__/distributed-task-id.test.ts b/packages/core/src/__tests__/distributed-task-id.test.ts index a5a4948702..c6e7e0aff4 100644 --- a/packages/core/src/__tests__/distributed-task-id.test.ts +++ b/packages/core/src/__tests__/distributed-task-id.test.ts @@ -1,6 +1,11 @@ import { describe, expect, it } from "vitest"; import { Database } from "../db.js"; -import { createDistributedTaskIdAllocator, DistributedTaskIdError, reconcileTaskIdState } from "../distributed-task-id.js"; +import { + createDistributedTaskIdAllocator, + DistributedTaskIdError, + reconcileTaskIdState, + rollbackDistributedTaskIdReservationForFailedCreateInExistingTransaction, +} from "../distributed-task-id.js"; describe("distributed-task-id allocator", () => { const createAllocator = () => { @@ -43,6 +48,24 @@ describe("distributed-task-id allocator", () => { expect(state.burnedReservationCount).toBe(1); }); + it("rolls back a committed failed-create reservation and preserves sequence permanence", async () => { + const { db, allocator } = createAllocator(); + const first = await allocator.reserveDistributedTaskId({ prefix: "FN", nodeId: "node-a" }); + await allocator.commitDistributedTaskIdReservation({ reservationId: first.reservationId, nodeId: "node-a" }); + + const rolledBack = db.transaction(() => rollbackDistributedTaskIdReservationForFailedCreateInExistingTransaction(db, { + reservationId: first.reservationId, + nodeId: "node-a", + reason: "failed-create", + })); + const second = await allocator.reserveDistributedTaskId({ prefix: "FN", nodeId: "node-a" }); + const state = await allocator.getDistributedTaskIdState({ prefix: "FN" }); + + expect(rolledBack).toMatchObject({ taskId: "FN-001", sequence: 1, committedClusterTaskCount: 0 }); + expect(second.taskId).toBe("FN-002"); + expect(state).toMatchObject({ committedClusterTaskCount: 0, burnedReservationCount: 1, nextSequence: 3 }); + }); + it("expired reservations cannot be committed and count as burned", async () => { const { allocator } = createAllocator(); const reservation = await allocator.reserveDistributedTaskId({ diff --git a/packages/core/src/__tests__/store-reservation-atomicity.test.ts b/packages/core/src/__tests__/store-reservation-atomicity.test.ts new file mode 100644 index 0000000000..44bf115ae9 --- /dev/null +++ b/packages/core/src/__tests__/store-reservation-atomicity.test.ts @@ -0,0 +1,224 @@ +import { afterEach, beforeAll, afterAll, describe, expect, it } from "vitest"; +import { existsSync } from "node:fs"; +import { rm } from "node:fs/promises"; +import { join } from "node:path"; + +import { InvalidFileScopeError, TaskStore, TombstonedTaskResurrectionError } from "../store.js"; +import { commitDistributedTaskIdReservationInExistingTransaction } from "../distributed-task-id.js"; +import { clearInMemoryDbSnapshot, installInMemoryDbSnapshot, makeTmpDir } from "./store-test-helpers.js"; + +function reservationRows(store: TaskStore) { + return store.getDatabase().prepare( + "SELECT taskId, status, sequence FROM distributed_task_id_reservations ORDER BY sequence", + ).all() as Array<{ taskId: string; status: string; sequence: number }>; +} + +function committedReservationPhantoms(store: TaskStore) { + return store.getDatabase().prepare( + `SELECT r.taskId + FROM distributed_task_id_reservations r + LEFT JOIN tasks t ON t.id = r.taskId + WHERE r.status = 'committed' AND t.id IS NULL + ORDER BY r.taskId`, + ).all() as Array<{ taskId: string }>; +} + +function reservationTaskMismatches(store: TaskStore) { + return store.getDatabase().prepare( + `SELECT t.id AS taskId, r.status + FROM tasks t + JOIN distributed_task_id_reservations r ON r.taskId = t.id + WHERE t.deletedAt IS NULL AND r.status != 'committed' + ORDER BY t.id`, + ).all() as Array<{ taskId: string; status: string }>; +} + +function expectNoReservationTaskDivergence(store: TaskStore) { + expect(committedReservationPhantoms(store)).toEqual([]); + expect(reservationTaskMismatches(store)).toEqual([]); +} + +async function createStore(options: { inMemoryDb: boolean }) { + const rootDir = makeTmpDir(); + const globalDir = makeTmpDir(); + const store = new TaskStore(rootDir, globalDir, { inMemoryDb: options.inMemoryDb }); + await store.init(); + return { rootDir, globalDir, store }; +} + +describe("FN-7074 task-create reservation atomicity", () => { + const cleanup: Array<() => Promise> = []; + + beforeAll(() => installInMemoryDbSnapshot()); + afterAll(() => clearInMemoryDbSnapshot()); + + afterEach(async () => { + while (cleanup.length > 0) { + await cleanup.pop()?.(); + } + }); + + async function scopedStore(options: { inMemoryDb: boolean } = { inMemoryDb: true }) { + const context = await createStore(options); + cleanup.push(async () => { + context.store.stopWatching(); + await context.store.close(); + await rm(context.rootDir, { recursive: true, force: true }); + await rm(context.globalDir, { recursive: true, force: true }); + }); + return context; + } + + it.each([ + ["in-memory", true], + ["file-backed", false], + ])("commits reservation iff task row and task directory land for %s stores", async (_label, inMemoryDb) => { + const { rootDir, store } = await scopedStore({ inMemoryDb }); + + const task = await store.createTask({ description: "happy atomic create" }); + + expect(reservationRows(store)).toEqual([{ taskId: task.id, status: "committed", sequence: 1 }]); + expect(store.getDatabase().prepare("SELECT id FROM tasks WHERE id = ?").get(task.id)).toMatchObject({ id: task.id }); + expect(existsSync(join(rootDir, ".fusion", "tasks", task.id, "task.json"))).toBe(true); + expect(existsSync(join(rootDir, ".fusion", "tasks", task.id, "PROMPT.md"))).toBe(true); + expectNoReservationTaskDivergence(store); + }); + + it("aborts the reservation and leaves no task row when the tasks-row insert fails", async () => { + const { store } = await scopedStore(); + const original = (store as unknown as { insertTaskWithFtsRecovery: (...args: unknown[]) => void }).insertTaskWithFtsRecovery; + (store as unknown as { insertTaskWithFtsRecovery: (...args: unknown[]) => void }).insertTaskWithFtsRecovery = () => { + throw new Error("synthetic insert failure"); + }; + + await expect(store.createTask({ description: "insert should fail" })).rejects.toThrow("synthetic insert failure"); + (store as unknown as { insertTaskWithFtsRecovery: (...args: unknown[]) => void }).insertTaskWithFtsRecovery = original; + + expect(reservationRows(store)).toEqual([{ taskId: "FN-001", status: "aborted", sequence: 1 }]); + expect(store.getDatabase().prepare("SELECT id FROM tasks WHERE id = ?").get("FN-001")).toBeUndefined(); + expectNoReservationTaskDivergence(store); + }); + + it("rolls back the committed reservation and task row when task.json disk write fails after insert", async () => { + const { rootDir, store } = await scopedStore({ inMemoryDb: false }); + const original = (store as unknown as { writeTaskJsonFile: (...args: unknown[]) => Promise }).writeTaskJsonFile; + (store as unknown as { writeTaskJsonFile: (...args: unknown[]) => Promise }).writeTaskJsonFile = async () => { + throw new Error("synthetic task.json write failure"); + }; + + await expect(store.createTask({ description: "disk write should fail" })).rejects.toThrow("synthetic task.json write failure"); + (store as unknown as { writeTaskJsonFile: (...args: unknown[]) => Promise }).writeTaskJsonFile = original; + + expect(reservationRows(store)).toEqual([{ taskId: "FN-001", status: "aborted", sequence: 1 }]); + expect(store.getDatabase().prepare("SELECT id FROM tasks WHERE id = ?").get("FN-001")).toBeUndefined(); + expect(existsSync(join(rootDir, ".fusion", "tasks", "FN-001"))).toBe(false); + expectNoReservationTaskDivergence(store); + }); + + it("rolls back distributed create reservations when file-scope validation throws", async () => { + const { rootDir, store } = await scopedStore(); + const originalGenerate = (store as unknown as { generateSpecifiedPrompt: (task: unknown) => string }).generateSpecifiedPrompt; + (store as unknown as { generateSpecifiedPrompt: (task: unknown) => string }).generateSpecifiedPrompt = () => + "# Bad prompt\n\n## File Scope\n\n- `origin/fusion/fn-4280`\n"; + + await expect(store.createTask({ description: "bad scope", column: "todo" })).rejects.toBeInstanceOf(InvalidFileScopeError); + (store as unknown as { generateSpecifiedPrompt: (task: unknown) => string }).generateSpecifiedPrompt = originalGenerate; + + expect(reservationRows(store)).toEqual([{ taskId: "FN-001", status: "aborted", sequence: 1 }]); + expect(store.getDatabase().prepare("SELECT id FROM tasks WHERE id = ?").get("FN-001")).toBeUndefined(); + expect(existsSync(join(rootDir, ".fusion", "tasks", "FN-001"))).toBe(false); + expectNoReservationTaskDivergence(store); + }); + + it("rolls back distributed create reservations when duplicate intake hits a recent tombstone", async () => { + const { store } = await scopedStore(); + await store.updateSettings({ tombstoneStickyWindowDays: 7 }); + const original = await store.createTask({ + title: "Memory leak", + description: "Fix memory leak in merge worker", + source: { sourceType: "unknown", sourceAgentId: "agent-1" }, + }); + await store.deleteTask(original.id); + + await expect(store.createTask({ + title: "Memory leak", + description: "Fix memory leak in merge worker", + source: { sourceType: "unknown", sourceAgentId: "agent-1" }, + })).rejects.toBeInstanceOf(TombstonedTaskResurrectionError); + + const rows = reservationRows(store); + expect(rows).toEqual([ + { taskId: "FN-001", status: "committed", sequence: 1 }, + { taskId: "FN-002", status: "aborted", sequence: 2 }, + ]); + expect(store.getDatabase().prepare("SELECT id FROM tasks WHERE id = ? AND deletedAt IS NULL").get("FN-002")).toBeUndefined(); + expectNoReservationTaskDivergence(store); + }); + + it("preserves ID permanence after a committed create is rolled back", async () => { + const { store } = await scopedStore(); + const original = (store as unknown as { writeTaskJsonFile: (...args: unknown[]) => Promise }).writeTaskJsonFile; + (store as unknown as { writeTaskJsonFile: (...args: unknown[]) => Promise }).writeTaskJsonFile = async () => { + throw new Error("synthetic task.json write failure"); + }; + await expect(store.createTask({ description: "burn FN-001" })).rejects.toThrow("synthetic task.json write failure"); + (store as unknown as { writeTaskJsonFile: (...args: unknown[]) => Promise }).writeTaskJsonFile = original; + + const next = await store.createTask({ description: "next id" }); + + expect(next.id).toBe("FN-002"); + expect(reservationRows(store)).toEqual([ + { taskId: "FN-001", status: "aborted", sequence: 1 }, + { taskId: "FN-002", status: "committed", sequence: 2 }, + ]); + expectNoReservationTaskDivergence(store); + }); + + it("allows replicated direct-reserved creates without requiring a reservation row", async () => { + const { store } = await scopedStore(); + const now = new Date().toISOString(); + + const result = await store.applyReplicatedTaskCreate({ + replicationVersion: 1, + reservationId: "remote-reservation", + taskId: "FN-123", + sourceNodeId: "node-b", + input: { + id: "FN-123", + description: "replicated create", + column: "triage", + dependencies: [], + steps: [], + currentStep: 0, + log: [], + createdAt: now, + updatedAt: now, + columnMovedAt: now, + } as never, + createdAt: now, + updatedAt: now, + prompt: "# replicated\n", + }); + + expect(result.applied).toBe(true); + expect(reservationRows(store)).toEqual([]); + expect(store.getDatabase().prepare("SELECT id FROM tasks WHERE id = ?").get("FN-123")).toMatchObject({ id: "FN-123" }); + }); + + it("commits reservations inside an existing store transaction without nested transaction errors", async () => { + const { store } = await scopedStore(); + const allocator = store.getDistributedTaskIdAllocator(); + const reservation = await allocator.reserveDistributedTaskId({ prefix: "FN", nodeId: "node-a" }); + + expect(() => { + store.getDatabase().transactionImmediate(() => { + commitDistributedTaskIdReservationInExistingTransaction(store.getDatabase(), { + reservationId: reservation.reservationId, + nodeId: "node-a", + }); + }); + }).not.toThrow(); + + expect(reservationRows(store)).toEqual([{ taskId: "FN-001", status: "committed", sequence: 1 }]); + }); +}); diff --git a/packages/core/src/distributed-task-id.ts b/packages/core/src/distributed-task-id.ts index 98a5aed311..755e1ac358 100644 --- a/packages/core/src/distributed-task-id.ts +++ b/packages/core/src/distributed-task-id.ts @@ -163,8 +163,7 @@ function getNextSequenceFloor(db: Database, prefix: string): number { return nextSequence; } -function ensureStateRow(db: Database, prefix: string): void { - const nowIso = new Date().toISOString(); +function ensureStateRow(db: Database, prefix: string, nowIso = new Date().toISOString()): void { const nextSequence = getNextSequenceFloor(db, prefix); db.prepare( `INSERT OR IGNORE INTO distributed_task_id_state ( @@ -179,6 +178,148 @@ function ensureStateRow(db: Database, prefix: string): void { ).run(nextSequence, nowIso, prefix); } +function expireDistributedTaskIdReservationsInExistingTransaction(db: Database, nowIso: string): number { + const result = db.prepare( + `UPDATE distributed_task_id_reservations + SET status = 'expired', reason = 'expired', abortedAt = ? + WHERE status = 'reserved' AND expiresAt <= ?`, + ).run(nowIso, nowIso) as { changes?: number }; + return result.changes ?? 0; +} + +/** + * FNXC:TaskIdReservation 2026-06-26-00:00: + * Distributed reservation commits must be composable into the caller's task-row insert transaction. This helper intentionally takes no allocator lock and opens no SQLite transaction so TaskStore can commit `distributed_task_id_reservations` and `tasks` atomically without better-sqlite3 nested-transaction failures. + */ +export function commitDistributedTaskIdReservationInExistingTransaction( + db: Database, + input: DistributedTaskIdCommitInput, + nowIso = new Date().toISOString(), +): DistributedTaskIdCommitResult { + expireDistributedTaskIdReservationsInExistingTransaction(db, nowIso); + const row = db + .prepare( + `SELECT reservationId, prefix, nodeId, sequence, taskId, status, reason, expiresAt, committedAt, abortedAt + FROM distributed_task_id_reservations + WHERE reservationId = ?`, + ) + .get(input.reservationId) as ReservationRow | undefined; + + if (!row) { + throw new DistributedTaskIdError("reservation not found", "reservation_not_found"); + } + if (row.nodeId !== input.nodeId) { + throw new DistributedTaskIdError("reservation belongs to a different node", "reservation_not_owned"); + } + if (row.status === "expired") { + throw new DistributedTaskIdError("reservation has expired", "reservation_expired"); + } + if (row.status !== "reserved") { + throw new DistributedTaskIdError("reservation already finalized", "reservation_finalized"); + } + + db.prepare( + `UPDATE distributed_task_id_reservations + SET status = 'committed', committedAt = ?, updatedAt = ? + WHERE reservationId = ?`, + ).run(nowIso, nowIso, row.reservationId); + + ensureStateRow(db, row.prefix, nowIso); + db.prepare( + `UPDATE distributed_task_id_state + SET committedClusterTaskCount = committedClusterTaskCount + 1, + lastCommittedTaskId = ?, + updatedAt = ? + WHERE prefix = ?`, + ).run(row.taskId, nowIso, row.prefix); + + const state = db + .prepare( + "SELECT committedClusterTaskCount FROM distributed_task_id_state WHERE prefix = ?", + ) + .get(row.prefix) as { committedClusterTaskCount: number }; + db.bumpLastModified(); + + return { + reservationId: row.reservationId, + taskId: row.taskId, + sequence: row.sequence, + committedClusterTaskCount: state.committedClusterTaskCount, + committedAt: nowIso, + }; +} + +/** + * FNXC:TaskIdReservation 2026-06-26-00:00: + * Post-insert create failures must burn the distributed ID without leaving a committed reservation. This transaction-participating rollback moves reserved or committed create reservations to `aborted`, recomputes committed counters, and preserves the reservation row so FN-5105 sequence permanence still prevents ID reuse. + */ +export function rollbackDistributedTaskIdReservationForFailedCreateInExistingTransaction( + db: Database, + input: DistributedTaskIdAbortInput, + nowIso = new Date().toISOString(), +): DistributedTaskIdAbortResult { + expireDistributedTaskIdReservationsInExistingTransaction(db, nowIso); + const row = db + .prepare( + `SELECT reservationId, prefix, nodeId, sequence, taskId, status, reason, expiresAt, committedAt, abortedAt + FROM distributed_task_id_reservations + WHERE reservationId = ?`, + ) + .get(input.reservationId) as ReservationRow | undefined; + + if (!row) { + throw new DistributedTaskIdError("reservation not found", "reservation_not_found"); + } + if (row.nodeId !== input.nodeId) { + throw new DistributedTaskIdError("reservation belongs to a different node", "reservation_not_owned"); + } + if (row.status === "expired") { + throw new DistributedTaskIdError("reservation has expired", "reservation_expired"); + } + if (row.status !== "reserved" && row.status !== "committed" && row.status !== "aborted") { + throw new DistributedTaskIdError("reservation already finalized", "reservation_finalized"); + } + + if (row.status !== "aborted") { + db.prepare( + `UPDATE distributed_task_id_reservations + SET status = 'aborted', reason = ?, committedAt = NULL, abortedAt = ?, updatedAt = ? + WHERE reservationId = ?`, + ).run(input.reason, nowIso, nowIso, row.reservationId); + } + + ensureStateRow(db, row.prefix, nowIso); + const count = db + .prepare( + "SELECT COUNT(*) AS count FROM distributed_task_id_reservations WHERE prefix = ? AND status = 'committed'", + ) + .get(row.prefix) as { count: number }; + const lastCommitted = db + .prepare( + `SELECT taskId FROM distributed_task_id_reservations + WHERE prefix = ? AND status = 'committed' + ORDER BY sequence DESC + LIMIT 1`, + ) + .get(row.prefix) as { taskId: string } | undefined; + db.prepare( + `UPDATE distributed_task_id_state + SET committedClusterTaskCount = ?, + lastCommittedTaskId = ?, + updatedAt = ? + WHERE prefix = ?`, + ).run(count.count, lastCommitted?.taskId ?? null, nowIso, row.prefix); + db.bumpLastModified(); + + return { + reservationId: row.reservationId, + taskId: row.taskId, + sequence: row.sequence, + committedClusterTaskCount: count.count, + abortedAt: nowIso, + }; +} + export function reconcileTaskIdState(db: Database): string[] { const nowIso = new Date().toISOString(); return db.transaction(() => { @@ -240,14 +381,7 @@ export function createDistributedTaskIdAllocator(db: Database): DistributedTaskI } }; - const expireReservations = (nowIso: string): number => { - const result = db.prepare( - `UPDATE distributed_task_id_reservations - SET status = 'expired', reason = 'expired', abortedAt = ? - WHERE status = 'reserved' AND expiresAt <= ?`, - ).run(nowIso, nowIso) as { changes?: number }; - return result.changes ?? 0; - }; + const expireReservations = (nowIso: string): number => expireDistributedTaskIdReservationsInExistingTransaction(db, nowIso); const taskIdExists = (prefix: string, sequence: number): boolean => { const taskId = formatDistributedTaskId(prefix, sequence); @@ -319,59 +453,7 @@ export function createDistributedTaskIdAllocator(db: Database): DistributedTaskI commitDistributedTaskIdReservation: async (input) => withLock(async () => { const nowIso = new Date().toISOString(); - return db.transaction(() => { - expireReservations(nowIso); - const row = db - .prepare( - `SELECT reservationId, prefix, nodeId, sequence, taskId, status, reason, expiresAt, committedAt, abortedAt - FROM distributed_task_id_reservations - WHERE reservationId = ?`, - ) - .get(input.reservationId) as ReservationRow | undefined; - - if (!row) { - throw new DistributedTaskIdError("reservation not found", "reservation_not_found"); - } - if (row.nodeId !== input.nodeId) { - throw new DistributedTaskIdError("reservation belongs to a different node", "reservation_not_owned"); - } - if (row.status === "expired") { - throw new DistributedTaskIdError("reservation has expired", "reservation_expired"); - } - if (row.status !== "reserved") { - throw new DistributedTaskIdError("reservation already finalized", "reservation_finalized"); - } - - db.prepare( - `UPDATE distributed_task_id_reservations - SET status = 'committed', committedAt = ?, updatedAt = ? - WHERE reservationId = ?`, - ).run(nowIso, nowIso, row.reservationId); - - ensureStateRow(db, row.prefix); - db.prepare( - `UPDATE distributed_task_id_state - SET committedClusterTaskCount = committedClusterTaskCount + 1, - lastCommittedTaskId = ?, - updatedAt = ? - WHERE prefix = ?`, - ).run(row.taskId, nowIso, row.prefix); - - const state = db - .prepare( - "SELECT committedClusterTaskCount FROM distributed_task_id_state WHERE prefix = ?", - ) - .get(row.prefix) as { committedClusterTaskCount: number }; - db.bumpLastModified(); - - return { - reservationId: row.reservationId, - taskId: row.taskId, - sequence: row.sequence, - committedClusterTaskCount: state.committedClusterTaskCount, - committedAt: nowIso, - }; - }); + return db.transaction(() => commitDistributedTaskIdReservationInExistingTransaction(db, input, nowIso)); }), abortDistributedTaskIdReservation: async (input) => withLock(async () => { diff --git a/packages/core/src/store.ts b/packages/core/src/store.ts index ee43f2d87b..0c9d5a885a 100644 --- a/packages/core/src/store.ts +++ b/packages/core/src/store.ts @@ -190,7 +190,14 @@ import { assertProjectRootDir, } from "./project-root-guard.js"; import { generateTaskLineageId, normalizeTaskCommitAssociation } from "./task-lineage.js"; -import { createDistributedTaskIdAllocator, reconcileTaskIdState, resolveLocalNodeId, type DistributedTaskIdAllocator } from "./distributed-task-id.js"; +import { + commitDistributedTaskIdReservationInExistingTransaction, + createDistributedTaskIdAllocator, + reconcileTaskIdState, + resolveLocalNodeId, + rollbackDistributedTaskIdReservationForFailedCreateInExistingTransaction, + type DistributedTaskIdAllocator, +} from "./distributed-task-id.js"; import { detectStalledReview } from "./stalled-review-detector.js"; import { computeRetrySummary } from "./retry-summary.js"; import { archiveAsSameAgentDuplicate, findSameAgentDuplicates } from "./duplicate-intake.js"; @@ -3663,13 +3670,25 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} * for backward compatibility and debugging. Create paths must call this variant * so duplicate IDs fail safely instead of overwriting existing rows. */ - private async atomicCreateTaskJson(dir: string, task: Task, operation: string): Promise { + private async atomicCreateTaskJson( + dir: string, + task: Task, + operation: string, + reservationCommit?: { reservationId: string; nodeId: string }, + ): Promise { const id = this.getTaskIdFromDir(dir); let deletedAt: string | undefined; this.db.transactionImmediate(() => { deletedAt = this.getSoftDeletedWriteConflict(id, task); if (deletedAt) return; this.insertTaskWithFtsRecovery(task, operation); + if (reservationCommit) { + /* + FNXC:TaskIdReservation 2026-06-26-00:00: + A distributed reservation is `committed` iff the corresponding `tasks` row is inserted in the same SQLite transaction. Disk artifacts are guarded separately after this transaction, but the reservation flip must never be a later durability point. + */ + commitDistributedTaskIdReservationInExistingTransaction(this.db, reservationCommit); + } }); if (deletedAt) { this.throwSoftDeletedWriteBlocked(id, deletedAt, operation); @@ -4285,7 +4304,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} options?: { onSummarize?: (description: string) => Promise; settings?: { autoSummarizeTitles?: boolean }; - createTaskWithId?: (taskId: string) => Promise; + createTaskWithId?: (taskId: string, reservationCommit: { reservationId: string; nodeId: string }) => Promise; }, ): Promise { const settings = await this.getSettingsFast(); @@ -4299,24 +4318,53 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} let createdTask: Task | null = null; try { + const reservationCommit = { reservationId: reservation.reservationId, nodeId }; createdTask = options?.createTaskWithId - ? await options.createTaskWithId(reservation.taskId) - : await this.createTaskWithReservedId(input, { taskId: reservation.taskId }); - await allocator.commitDistributedTaskIdReservation({ - reservationId: reservation.reservationId, - nodeId, - }); + ? await options.createTaskWithId(reservation.taskId, reservationCommit) + : await this.createTaskWithReservedId(input, { taskId: reservation.taskId, reservationCommit }); return createdTask; } catch (error) { - await allocator.abortDistributedTaskIdReservation({ - reservationId: reservation.reservationId, - nodeId, - reason: "failed-create", - }).catch(() => undefined); + await this.rollbackFailedDistributedReservationCreate( + reservation.taskId, + { reservationId: reservation.reservationId, nodeId }, + error, + ).catch(() => undefined); throw error; } } + private async rollbackFailedDistributedReservationCreate( + taskId: string, + reservationCommit: { reservationId: string; nodeId: string }, + cause: unknown, + ): Promise { + const dir = this.taskDir(taskId); + if (this.isWatching) this.taskCache.delete(taskId); + this.db.transactionImmediate(() => { + this.deleteTaskById(taskId); + rollbackDistributedTaskIdReservationForFailedCreateInExistingTransaction(this.db, { + reservationId: reservationCommit.reservationId, + nodeId: reservationCommit.nodeId, + reason: "failed-create", + }); + this.insertRunAuditEventRow({ + taskId, + domain: "database", + mutationType: "task:reservation-commit-rolled-back", + target: taskId, + metadata: { + reservationId: reservationCommit.reservationId, + nodeId: reservationCommit.nodeId, + reason: "failed-create", + error: cause instanceof Error ? cause.message : String(cause), + }, + }); + }); + if (existsSync(dir)) { + await rm(dir, { recursive: true, force: true }); + } + } + private taskDir(id: string): string { return join(this.tasksDir, id); } @@ -4551,14 +4599,14 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} // U7c: selection seeds are optional-group node ids (not materialized // `workflow_steps` rows), so a failed task creation strands nothing to clean. const task: Task = await this.createTaskWithDistributedReservation(input, { - createTaskWithId: async (taskId) => { + createTaskWithId: async (taskId, reservationCommit) => { await this.assertNoDependencyCycle(taskId, input.dependencies ?? [], "createTask"); return this._createTaskInternal( input, title, resolvedWorkflowSteps, taskId, - { invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization }, + { invokeTaskCreatedHook: shouldInvokeTaskCreatedHook && !hasPendingSummarization, reservationCommit }, ); }, }); @@ -4651,6 +4699,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} prompt?: string; applyDefaultWorkflowSteps?: boolean; invokeTaskCreatedHook?: boolean; + reservationCommit?: { reservationId: string; nodeId: string }; }, ): Promise { if (!input.description?.trim()) { @@ -4744,6 +4793,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} updatedAt: options.updatedAt, promptOverride: options.prompt, invokeTaskCreatedHook: options.invokeTaskCreatedHook, + reservationCommit: options.reservationCommit, }); // Record the inherited workflow selection now that the task row exists. @@ -4816,6 +4866,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} updatedAt?: string; promptOverride?: string; invokeTaskCreatedHook?: boolean; + reservationCommit?: { reservationId: string; nodeId: string }; }, ): Promise { const now = options?.createdAt ?? new Date().toISOString(); @@ -4885,7 +4936,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} this.assertTaskIdAvailable(id); const dir = this.taskDir(id); - await this.atomicCreateTaskJson(dir, task, "createTask"); + await this.atomicCreateTaskJson(dir, task, "createTask", options?.reservationCommit); // Update cache if watcher is active if (this.isWatching) this.taskCache.set(id, { ...task }); @@ -5059,7 +5110,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} const now = new Date().toISOString(); return this.createTaskWithDistributedReservation({ description: sourceTask.description }, { - createTaskWithId: async (newId) => { + createTaskWithId: async (newId, reservationCommit) => { // FN-5077: duplicated drift-stripped fragments may normalize to null and should remain unset. const normalizedTitle = normalizeTitleForTaskId(sourceTask.title, newId); if (normalizedTitle.changed) { @@ -5090,7 +5141,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} this.assertTaskIdAvailable(newId); const newDir = this.taskDir(newId); - await this.atomicCreateTaskJson(newDir, newTask, "duplicateTask"); + await this.atomicCreateTaskJson(newDir, newTask, "duplicateTask", reservationCommit); const sanitizedPrompt = sanitizeFileScopeInPromptContent(sourceTask.prompt); if (sanitizedPrompt.dropped.length > 0) { storeLog.log(`[file-scope-sanitize] duplicate ${newId} from ${id}: dropped=[${sanitizedPrompt.dropped.join(",")}]`); @@ -5137,7 +5188,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} } return this.createTaskWithDistributedReservation({ description: feedback.trim() }, { - createTaskWithId: async (newId) => { + createTaskWithId: async (newId, reservationCommit) => { // FN-5077: keep deterministic "Refinement" fallback when normalized refinement label is unusable (null). const normalizedTitle = normalizeTitleForTaskId(`Refinement: ${sourceLabel}`, newId); if (normalizedTitle.changed) { @@ -5179,7 +5230,7 @@ ${TASK_UPSERT_SQL_ASSIGNMENTS} this.assertTaskIdAvailable(newId); const newDir = this.taskDir(newId); - await this.atomicCreateTaskJson(newDir, newTask, "refineTask"); + await this.atomicCreateTaskJson(newDir, newTask, "refineTask", reservationCommit); const prompt = `# ${newTask.title}\n\n${newTask.description}\n`; const sanitizedPrompt = sanitizeFileScopeInPromptContent(prompt); await mkdir(newDir, { recursive: true });