From 649b901bbda388aca8a5c3d6b0f71f44b366caac Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Thu, 20 Aug 2026 00:28:59 -0700 Subject: [PATCH] FN-9178: Classify awaited core run-audit emitters Document and ratchet which core run-audit writers may remain lifecycle-coupled. - classify awaited emitters by bounded-isolation safety and preserve transactional boundaries - characterize hostile sink behavior across store- and data-layer call sites - fail closed when new awaited audit writers lack an explicit classification Files changed: AGENTS.md | 1 + docs/run-audit.md | 6 + .../core-run-audit-emitter-isolation.test.ts | 71 ++++- .../excluded-awaited-run-audit-layer-sites.test.ts | 332 +++++++++++++++++++++ .../excluded-awaited-run-audit-store-sites.test.ts | 236 +++++++++++++++ 5 files changed, 645 insertions(+), 1 deletion(-) Fusion-Task-Id: FN-9178 Fusion-Task-Lineage: 25fed6a5-4f2f-4d4b-8724-838eff4323c8 Co-authored-by: Fusion (runfusion.ai) --- AGENTS.md | 1 + docs/run-audit.md | 6 + .../core-run-audit-emitter-isolation.test.ts | 71 +++- ...uded-awaited-run-audit-layer-sites.test.ts | 332 ++++++++++++++++++ ...uded-awaited-run-audit-store-sites.test.ts | 236 +++++++++++++ 5 files changed, 645 insertions(+), 1 deletion(-) create mode 100644 packages/core/src/__tests__/excluded-awaited-run-audit-layer-sites.test.ts create mode 100644 packages/core/src/__tests__/excluded-awaited-run-audit-store-sites.test.ts diff --git a/AGENTS.md b/AGENTS.md index da31a3a8af..31f069a69c 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -430,3 +430,4 @@ Note: the embedded main-content views Workflows (`_WorkflowEditorView`), Import ``` - FN-9177: New core best-effort emitters must use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. It deliberately mirrors the engine seam because `@fusion/core` cannot import `@fusion/engine`; transactional and deliberately awaited durability writers remain unbounded. +- FN-9178: Awaited core audit exclusions are classified with evidence in `docs/run-audit.md`: class A candidates, class B outcome-signalled candidates, and class C forensic/durability records. Transactional writers remain permanently unbounded because they share the mutation transaction. diff --git a/docs/run-audit.md b/docs/run-audit.md index 0c9cb7819a..65682638ab 100644 --- a/docs/run-audit.md +++ b/docs/run-audit.md @@ -83,3 +83,9 @@ This applies to executor, run-auditor, self-healing, merger, PR reconciliation, ### Core emit-seam policy Core best-effort emitters use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. This is a deliberate copy of the engine seam because `@fusion/core` cannot import `@fusion/engine`; it synchronously invokes a valid sink, then absorbs throws, rejection, timeout, and late settlement without making telemetry lifecycle-load-bearing. Transactional writers and explicitly awaited durability/ordering writers remain unbounded. `packages/core/src/__tests__/core-run-audit-sink-health.test.ts` and `core-run-audit-emitter-isolation.test.ts` respectively enforce hostile-sink behavior and source routing. + +### Awaited core exclusion decision + +FN-9178 classified the remaining awaited sites with hostile-sink characterization tests. `task-deleted-outbox:catch-up`, `:reconciliation-fallback`, `:lease-fenced`, `:retention-pruned`, and recall-capture audit events are class A (bound-safe) candidates. `task:workflow-switch-torn` and `task:reconcile-phantom-committed-reservation` are class B: a future bounded seam must return `recorded`, `failed`, `timed-out`, or `absent`, because their throw/result payload depends on audit outcome. `task:bypass-review`, `task:resume-step`, and both resurrection-blocked records are class C and intentionally unbounded: they claim persistence before return, destructive cleanup, or a forensic throw. + +All `recordRunAuditEventWithinTransaction(tx, ...)` calls and the `recordRunAuditEventBackend(tx, ...)` transactional call are permanently out of scope. Their audit row shares a transaction with the mutation it describes; bounding would split that atomicity. The full matrix and evidence pointers are in the FN-9178 `decision` task document; `excluded-awaited-run-audit-store-sites.test.ts`, `excluded-awaited-run-audit-layer-sites.test.ts`, and the core routing ratchet pin this boundary. diff --git a/packages/core/src/__tests__/core-run-audit-emitter-isolation.test.ts b/packages/core/src/__tests__/core-run-audit-emitter-isolation.test.ts index 3509ee5959..22ffbf352e 100644 --- a/packages/core/src/__tests__/core-run-audit-emitter-isolation.test.ts +++ b/packages/core/src/__tests__/core-run-audit-emitter-isolation.test.ts @@ -1,7 +1,36 @@ -import { readFileSync } from "node:fs"; +import { readdirSync, readFileSync } from "node:fs"; import { fileURLToPath } from "node:url"; import { describe, expect, it } from "vitest"; +/* + * FNXC:RunAudit 2026-08-20-07:16: + * FN-9178 makes every direct awaited core audit writer a named decision rather than an implicit + * exception. Class A sites are evaluated candidates, B needs a reporting bounded seam, C retains + * its ordering claim, and transactional/sink writers remain permanent atomicity boundaries. + */ +const awaitedClassifications = { + "store.ts:task:bypass-review": "C", + "store.ts:task:resume-step": "C", + "task-store/workflow-definitions.ts:task:workflow-switch-torn": "B", + "task-store/async/async-phantom-reservations.ts:task:reconcile-phantom-committed-reservation": "B", + "task-store/task-creation.ts:intake:resurrection-blocked": "C", + "task-store/task-id-integrity.ts:task:resurrection-blocked": "C", + "task-store/task-deleted-outbox-consumer.ts:task-deleted-outbox:catch-up": "A", + "task-store/task-deleted-outbox-consumer.ts:task-deleted-outbox:reconciliation-fallback": "A", + "task-store/task-deleted-outbox-consumer.ts:task-deleted-outbox:lease-fenced": "A", + "task-store/task-lifecycle-event-retention.ts:task-deleted-outbox:retention-pruned": "A", + "memory/recall-capture.ts:memory:capture-recorded|memory:capture-failed": "A", + "task-store/project-store-ops.ts:recordRunAuditEventImpl": "permanent-sink", +} as const; + +const transactionalSourceBoundaries = [ + "task-store/moves.ts", "task-store/symbol-locks.ts", "task-store/async/async-merge-coordination.ts", + "task-store/lifecycle-ops.ts", "task-store/task-creation.ts", "task-store/project-store-ops.ts", + "task-store/task-artifacts-ops.ts", "task-store/task-lifecycle-consumer-registry.ts", + "task-store/async/async-comments-attachments.ts", "task-store/async/async-workflow-workitems.ts", + "task-store/archive-lifecycle-2.ts", +] as const; + const files = [ "../planner/planner-intervention.ts", "../task-store/audit-ops.ts", "../task-store/branch-group-ops.ts", "../task-store/merge-queue-ops-2.ts", "../task-store/task-mutation-ops.ts", "../task-store/workflow-integrity.ts", @@ -9,7 +38,47 @@ const files = [ "../task-store/lifecycle-ops.ts", "../task-store/task-id-integrity.ts", ]; +const sourceRoot = fileURLToPath(new URL("..", import.meta.url)); + +function sourceFiles(directory: string, prefix = ""): string[] { + return readdirSync(directory, { withFileTypes: true }).flatMap((entry) => { + const relative = `${prefix}${entry.name}`; + if (entry.isDirectory()) return entry.name === "__tests__" ? [] : sourceFiles(`${directory}/${entry.name}`, `${relative}/`); + return entry.isFile() && entry.name.endsWith(".ts") ? [relative] : []; + }); +} + +function awaitedAuditInventory() { + const call = /await\s+(?:(?:this|store)\.)?recordRunAuditEvent(?:Async)?\s*\(/g; + return sourceFiles(sourceRoot).flatMap((relative) => { + const source = readFileSync(`${sourceRoot}/${relative}`, "utf8"); + return [...source.matchAll(call)].map((match) => { + const body = source.slice(match.index, match.index! + 1_600); + const mutation = body.match(/mutationType:\s*(?:"([^"]+)"|type)/)?.[1] + ?? (relative === "memory/recall-capture.ts" ? "memory:capture-recorded|memory:capture-failed" : undefined); + return `${relative}:${mutation ?? "recordRunAuditEventImpl"}`; + }); + }).sort(); +} + +function transactionalAuditInventory() { + const call = /await\s+(?:(?:store\.)?recordRunAuditEventWithinTransaction|store\.recordRunAuditEventBackend)\s*\(/g; + return sourceFiles(sourceRoot).flatMap((relative) => { + const source = readFileSync(`${sourceRoot}/${relative}`, "utf8"); + return [...source.matchAll(call)].map(() => relative); + }).sort(); +} + describe("core run-audit emitter isolation", () => { + it("fails closed when a direct awaited non-transactional audit emitter lacks a classification", () => { + expect(awaitedAuditInventory()).toEqual(Object.keys(awaitedClassifications).sort()); + }); + + it("keeps transactional awaited writers as named atomicity boundaries", () => { + const inventory = new Set(transactionalAuditInventory()); + expect([...inventory].sort()).toEqual([...transactionalSourceBoundaries].sort()); + }); + it("routes every best-effort core emitter through the bounded seam", () => { for (const relative of files) { const source = readFileSync(fileURLToPath(new URL(relative, import.meta.url)), "utf8"); diff --git a/packages/core/src/__tests__/excluded-awaited-run-audit-layer-sites.test.ts b/packages/core/src/__tests__/excluded-awaited-run-audit-layer-sites.test.ts new file mode 100644 index 0000000000..89a7f492ae --- /dev/null +++ b/packages/core/src/__tests__/excluded-awaited-run-audit-layer-sites.test.ts @@ -0,0 +1,332 @@ +import { afterEach, describe, expect, it, vi } from "vitest"; + +const audit = vi.hoisted(() => vi.fn()); +const asyncAudit = vi.hoisted(() => vi.fn()); +const softDelete = vi.hoisted(() => vi.fn().mockResolvedValue(undefined)); +const readTaskRow = vi.hoisted(() => vi.fn()); +const lifecycle = vi.hoisted(() => ({ + acknowledgeTaskLifecycleEvent: vi.fn(async () => true), + acquireTaskLifecycleLease: vi.fn(async () => ({ token: "lease", fencingToken: 1n })), + advanceTaskLifecycleConsumerCursor: vi.fn(async () => true), + hasTaskLifecycleConsumerReceipt: vi.fn(async () => false), + listTaskLifecycleEvents: vi.fn(async () => []), + readTaskLifecycleConsumerCursor: vi.fn(async () => ({ fencingToken: 1n, lastAckedSeq: 0n, retryAttempts: 0, updatedAt: new Date().toISOString() })), + readTaskLifecycleEventBounds: vi.fn(async () => ({ headSeq: 4n, oldestSeq: 1n })), + registerTaskLifecycleConsumer: vi.fn(async () => undefined), + releaseTaskLifecycleLease: vi.fn(async () => undefined), + renewTaskLifecycleLease: vi.fn(async () => true), + setTaskLifecycleConsumerActive: vi.fn(async () => undefined), +})); +vi.mock("../postgres/data-layer.js", async (importOriginal) => ({ + ...(await importOriginal()), + recordRunAuditEvent: audit, +})); +vi.mock("../task-store/async/async-audit.js", async (importOriginal) => ({ + ...(await importOriginal()), + recordRunAuditEvent: asyncAudit, +})); +vi.mock("../task-store/async/async-persistence.js", async (importOriginal) => ({ + ...(await importOriginal()), + softDeleteTaskRow: softDelete, + readTaskRow, +})); +vi.mock("../task-store/task-lifecycle-consumer-registry.js", async (importOriginal) => ({ + ...(await importOriginal()), + acknowledgeTaskLifecycleEvent: lifecycle.acknowledgeTaskLifecycleEvent, + acquireTaskLifecycleLease: lifecycle.acquireTaskLifecycleLease, + advanceTaskLifecycleConsumerCursor: lifecycle.advanceTaskLifecycleConsumerCursor, + hasTaskLifecycleConsumerReceipt: lifecycle.hasTaskLifecycleConsumerReceipt, + listTaskLifecycleEvents: lifecycle.listTaskLifecycleEvents, + readTaskLifecycleConsumerCursor: lifecycle.readTaskLifecycleConsumerCursor, + readTaskLifecycleEventBounds: lifecycle.readTaskLifecycleEventBounds, + registerTaskLifecycleConsumer: lifecycle.registerTaskLifecycleConsumer, + releaseTaskLifecycleLease: lifecycle.releaseTaskLifecycleLease, + renewTaskLifecycleLease: lifecycle.renewTaskLifecycleLease, + setTaskLifecycleConsumerActive: lifecycle.setTaskLifecycleConsumerActive, +})); + +import { createRecallCaptureWriter } from "../memory/recall-capture.js"; +import { resolveSameAgentDuplicateIntake } from "../task-store/task-creation.js"; +import { maybeResolveTombstonedTaskIdImpl } from "../task-store/task-id-integrity.js"; +import { TombstonedTaskResurrectionError } from "../task-store/errors.js"; +import { pruneTaskLifecycleEvents } from "../task-store/task-lifecycle-event-retention.js"; +import { TaskDeletedOutboxConsumer } from "../task-store/task-deleted-outbox-consumer.js"; + +/* + * FNXC:RunAudit 2026-08-20-06:40: + * FN-9178 invokes real helper-owned entry points with hostile audit helpers. It records existing + * awaiting/ordering behavior only; fake timers make the never-settling observation deterministic. + */ +function emptyQuery() { + const query: Record = {}; + for (const method of ["from", "where", "orderBy", "limit"]) query[method] = () => query; + query.then = (resolve: (value: unknown[]) => unknown) => resolve([]); + return query; +} + +function retentionLayer() { + return { projectId: "project", db: { select: vi.fn(() => emptyQuery()) } } as never; +} + +describe("FN-9178 awaited data-layer audit characterization", () => { + afterEach(() => { vi.clearAllMocks(); vi.useRealTimers(); }); + + function resurrectionStore() { + return { + asyncLayer: { projectId: "project", db: {} }, isWatching: true, taskCache: new Map(), + taskDir: vi.fn(() => "/definitely-absent"), getSettings: vi.fn().mockResolvedValue({ tombstoneStickyWindowDays: 7 }), + listTasksBySourceLineage: vi.fn(), + } as never; + } + + it.each([ + ["absent", () => undefined, true], + ["synchronous throw", () => { throw new Error("sync"); }, false], + ["rejection", () => Promise.reject(new Error("reject")), false], + ])("intake resurrection keeps destructive follow-up ordered after a %s audit", async (_state, sink, deletes) => { + asyncAudit.mockImplementation(sink as never); + const store = resurrectionStore(); + const deletedAt = new Date().toISOString(); + const task = { id: "FN-INTAKE", title: "new", description: "new", column: "todo", createdAt: deletedAt, sourceAgentId: "agent", sourceParentTaskId: null }; + store.listTasksBySourceLineage.mockResolvedValue([task, { ...task, id: "FN-TOMB", deletedAt, allowResurrection: false }]); + const operation = resolveSameAgentDuplicateIntake(store, task as never, task as never); + if (deletes) await expect(operation).rejects.toBeInstanceOf(TombstonedTaskResurrectionError); + else await expect(operation).resolves.toBeUndefined(); // The helper fails open when forensic audit cannot land. + expect(softDelete).toHaveBeenCalledTimes(deletes ? 1 : 0); + }); + + it("intake resurrection remains pending and cannot delete before a never-settling audit", async () => { + const store = resurrectionStore(); const deletedAt = new Date().toISOString(); + const task = { id: "FN-never-settling", title: "new", description: "new", column: "todo", createdAt: deletedAt, sourceAgentId: "agent", sourceParentTaskId: null }; + store.listTasksBySourceLineage.mockResolvedValue([task, { ...task, id: "FN-TOMB", deletedAt, allowResurrection: false }]); + asyncAudit.mockImplementation(() => new Promise(() => undefined)); + vi.useFakeTimers(); let settled = false; + void resolveSameAgentDuplicateIntake(store, task as never, task as never).finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(softDelete).not.toHaveBeenCalled(); + expect(settled).toBe(false); + }); + + it("intake resurrection throws its typed error after a late audit permits destructive cleanup", async () => { + const store = resurrectionStore(); const deletedAt = new Date().toISOString(); + const task = { id: "FN-late-settling", title: "new", description: "new", column: "todo", createdAt: deletedAt, sourceAgentId: "agent", sourceParentTaskId: null }; + store.listTasksBySourceLineage.mockResolvedValue([task, { ...task, id: "FN-TOMB", deletedAt, allowResurrection: false }]); + asyncAudit.mockImplementation(() => new Promise((resolve) => setTimeout(resolve, 2_100))); + vi.useFakeTimers(); + const operation = resolveSameAgentDuplicateIntake(store, task as never, task as never); + // FNXC:RunAudit 2026-08-20-07:16: Observe the deferred forensic rejection before fake-time advancement, then assert its real type. + void operation.catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(softDelete).toHaveBeenCalledOnce(); + await expect(operation).rejects.toBeInstanceOf(TombstonedTaskResurrectionError); + }); + + it.each([ + ["absent", () => undefined], + ["synchronous throw", () => { throw new Error("sync"); }], + ["rejection", () => Promise.reject(new Error("reject"))], + ["never-settling", () => new Promise(() => undefined)], + ["late-settling", () => new Promise((resolve) => setTimeout(resolve, 2_100))], + ])("id-integrity resurrection preserves its typed throw after %s audit behavior", async (_state, sink) => { + const deletedAt = new Date().toISOString(); + readTaskRow.mockResolvedValue({ deletedAt, allowResurrection: false }); + asyncAudit.mockImplementation(sink as never); + const store = resurrectionStore(); + if (_state.includes("settling")) vi.useFakeTimers(); + const operation = maybeResolveTombstonedTaskIdImpl(store, "FN-TOMB", {}, "createTask"); + // Attach an observer immediately: fake-time advancement may settle the late audit before the + // assertion below awaits this deliberately rejecting forensic entry point. + void operation.catch(() => undefined); + if (_state === "never-settling") { + let settled = false; void operation.finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); expect(settled).toBe(false); + } else if (_state === "late-settling") { + await vi.advanceTimersByTimeAsync(2_100); + await expect(operation).rejects.toBeInstanceOf(TombstonedTaskResurrectionError); + } else if (_state === "absent") { + await expect(operation).rejects.toBeInstanceOf(TombstonedTaskResurrectionError); + } else { + // Unlike the intake helper, this forensic pre-throw entry lets audit failure replace its + // typed resurrection error; that observable ordering is class-B evidence, not a fix. + await expect(operation).rejects.toThrow(_state === "synchronous throw" ? "sync" : "reject"); + } + }); + it.each([ + ["absent", () => undefined], + ["synchronous throw", () => { throw new Error("sync"); }], + ["rejection", () => Promise.reject(new Error("rejected"))], + ])("drives retention pruning through a %s audit helper", async (_state, sink) => { + audit.mockImplementation(sink as never); + const layer = retentionLayer(); + const operation = pruneTaskLifecycleEvents(layer, "project"); + if (_state === "absent") await expect(operation).resolves.toMatchObject({ prunedCount: 0 }); + else await expect(operation).rejects.toThrow(); + expect(audit).toHaveBeenCalledWith(layer, expect.objectContaining({ mutationType: "task-deleted-outbox:retention-pruned" })); + }); + + it("retention pruning remains pending for a never-settling audit helper", async () => { + const layer = retentionLayer(); + audit.mockImplementation(() => new Promise(() => undefined)); + vi.useFakeTimers(); + try { + let settled = false; + void pruneTaskLifecycleEvents(layer, "project").finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(settled).toBe(false); + } finally { vi.useRealTimers(); } + }); + + it("a late retention audit delays the real return until it settles", async () => { + const layer = retentionLayer(); + let resolve!: () => void; + audit.mockImplementation(() => new Promise((done) => { resolve = done; })); + audit.mockClear(); + const operation = pruneTaskLifecycleEvents(layer, "project"); + await vi.waitFor(() => expect(audit).toHaveBeenCalled()); + resolve(); + await expect(operation).resolves.toMatchObject({ prunedCount: 0 }); + }); + + it.each([ + ["absent", () => undefined, true], + ["synchronous throw", () => { throw new Error("sync"); }, false], + ["rejection", () => Promise.reject(new Error("rejected")), false], + ])("reconciliation fallback advances its cursor before a %s audit helper", async (_state, sink, completes) => { + audit.mockImplementation(sink as never); + const layer = retentionLayer(); + const store = { asyncLayer: layer, consumerId: "consumer", taskCache: new Map(), emitObservedTaskDeleted: vi.fn() } as never; + const operation = (new TaskDeletedOutboxConsumer(store) as never).reconcile(0n, { fencingToken: 1n }, "pruned-gap"); + if (completes) await expect(operation).resolves.toBe(true); else await expect(operation).rejects.toThrow(); + expect(lifecycle.advanceTaskLifecycleConsumerCursor).toHaveBeenCalledOnce(); + expect(audit).toHaveBeenCalledWith(layer, expect.objectContaining({ mutationType: "task-deleted-outbox:reconciliation-fallback" })); + }); + + it.each(["never-settling", "late-settling"])("reconciliation fallback waits after its cursor advance for %s audit", async (state) => { + const layer = retentionLayer(); + audit.mockImplementation((state === "never-settling" ? () => new Promise(() => undefined) : () => new Promise((resolve) => setTimeout(resolve, 2_100))) as never); + const store = { asyncLayer: layer, consumerId: "consumer", taskCache: new Map(), emitObservedTaskDeleted: vi.fn() } as never; + vi.useFakeTimers(); + let settled = false; + void (new TaskDeletedOutboxConsumer(store) as never).reconcile(0n, { fencingToken: 1n }, "pruned-gap").finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(lifecycle.advanceTaskLifecycleConsumerCursor).toHaveBeenCalledOnce(); + expect(settled).toBe(state === "late-settling"); + }); + + function catchUpStore() { + return { + asyncLayer: retentionLayer(), consumerId: "consumer", taskCache: new Map([["FN-DELETED", { id: "FN-DELETED" }]]), + emitObservedTaskDeleted: vi.fn(), + } as never; + } + + function deletedEvent() { + return { + eventId: "event-1", eventType: "task:deleted", taskId: "FN-DELETED", occurredAt: new Date().toISOString(), seq: 1n, + payload: { taskId: "FN-DELETED", previousColumn: "todo", previousStatus: null, deletedAt: new Date().toISOString(), allowResurrection: false, githubIssueAction: null, closureContext: null, deletedBy: null }, + }; + } + + it.each([ + ["absent", () => undefined, true], + ["synchronous throw", () => { throw new Error("sync"); }, false], + ["rejection", () => Promise.reject(new Error("rejected")), false], + ])("drives catch-up through the production poll batch before a %s audit helper", async (_state, sink, completes) => { + lifecycle.listTaskLifecycleEvents.mockResolvedValueOnce([deletedEvent()]); + audit.mockImplementation(sink as never); + const consumer = new TaskDeletedOutboxConsumer(catchUpStore()); + /* + * FNXC:RunAudit 2026-08-20-07:04: + * The characterization must use poll's production batch path. Set only its lifecycle-owned + * running gate to avoid adding a background scheduler while retaining lease-to-release order. + */ + (consumer as never).running = true; + const operation = consumer.poll(); + if (completes) await expect(operation).resolves.toBe("active"); else await expect(operation).rejects.toThrow(); + expect(lifecycle.acknowledgeTaskLifecycleEvent).toHaveBeenCalledOnce(); + expect(audit).toHaveBeenCalledWith(expect.anything(), expect.objectContaining({ mutationType: "task-deleted-outbox:catch-up" })); + expect(lifecycle.releaseTaskLifecycleLease).toHaveBeenCalledOnce(); + }); + + it("keeps the real catch-up poll and its lease cleanup pending for a never-settling audit", async () => { + lifecycle.listTaskLifecycleEvents.mockResolvedValueOnce([deletedEvent()]); + audit.mockImplementation(() => new Promise(() => undefined)); + const consumer = new TaskDeletedOutboxConsumer(catchUpStore()); + (consumer as never).running = true; + vi.useFakeTimers(); + try { + let settled = false; + void consumer.poll().finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(lifecycle.acknowledgeTaskLifecycleEvent).toHaveBeenCalledOnce(); + expect(settled).toBe(false); + expect(lifecycle.releaseTaskLifecycleLease).not.toHaveBeenCalled(); + } finally { vi.useRealTimers(); } + }); + + it("finishes the real catch-up poll only after a late audit settles beyond the bounded window", async () => { + lifecycle.listTaskLifecycleEvents.mockResolvedValueOnce([deletedEvent()]); + let resolveAudit!: () => void; + audit.mockImplementation(() => new Promise((resolve) => { resolveAudit = resolve; })); + const consumer = new TaskDeletedOutboxConsumer(catchUpStore()); + (consumer as never).running = true; + vi.useFakeTimers(); + try { + const operation = consumer.poll(); + await vi.advanceTimersByTimeAsync(2_100); + expect(lifecycle.acknowledgeTaskLifecycleEvent).toHaveBeenCalledOnce(); + expect(lifecycle.releaseTaskLifecycleLease).not.toHaveBeenCalled(); + resolveAudit(); + await expect(operation).resolves.toBe("active"); + expect(lifecycle.releaseTaskLifecycleLease).toHaveBeenCalledOnce(); + } finally { vi.useRealTimers(); } + }); + + it.each([ + ["absent", () => undefined, true], + ["synchronous throw", () => { throw new Error("sync"); }, false], + ["rejection", () => Promise.reject(new Error("rejected")), false], + ])("drives the real outbox lease-fenced entry through a %s audit helper", async (_state, sink, completes) => { + audit.mockImplementation(sink as never); + const store = { asyncLayer: retentionLayer(), consumerId: "consumer" } as never; + const operation = (new TaskDeletedOutboxConsumer(store) as never).recordLeaseFenced({ fencingToken: 1n }, 1); + if (completes) await expect(operation).resolves.toBeUndefined(); else await expect(operation).rejects.toThrow(); + expect(audit).toHaveBeenCalledWith(store.asyncLayer, expect.objectContaining({ mutationType: "task-deleted-outbox:lease-fenced" })); + }); + + it.each(["never-settling", "late-settling"])("outbox lease-fenced awaits a %s audit helper", async (state) => { + audit.mockImplementation((state === "never-settling" ? () => new Promise(() => undefined) : () => new Promise((resolve) => setTimeout(resolve, 2_100))) as never); + const store = { asyncLayer: retentionLayer(), consumerId: "consumer" } as never; + vi.useFakeTimers(); let settled = false; + void (new TaskDeletedOutboxConsumer(store) as never).recordLeaseFenced({ fencingToken: 1n }, 1).finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(settled).toBe(state === "late-settling"); + }); + + it.each([ + ["absent", () => undefined], + ["synchronous throw", () => { throw new Error("sync"); }], + ["rejection", () => Promise.reject(new Error("reject"))], + ["never-settling", () => new Promise(() => undefined)], + ["late-settling", () => new Promise((resolve) => setTimeout(resolve, 2_100))], + ])("recall capture contains %s injected audit without unhandled rejection", async (_state, injectedAudit) => { + const logger = { warn: vi.fn() }; + const writer = createRecallCaptureWriter({ + layer: {} as never, append: async () => ({ status: "created", record: { id: "recall-1" } }) as never, + audit: injectedAudit as never, logger, + }); + if (_state.includes("settling")) vi.useFakeTimers(); + writer.capture({ origin: "insight", summary: "summary", insightId: "INS-1" }); + if (_state === "never-settling") { + await vi.advanceTimersByTimeAsync(2_100); + expect(logger.warn).not.toHaveBeenCalled(); + } else if (_state === "late-settling") { + await vi.advanceTimersByTimeAsync(2_100); + await writer.flushPendingCaptures(); + expect(logger.warn).not.toHaveBeenCalled(); + } else { + await writer.flushPendingCaptures(); + expect(logger.warn).toHaveBeenCalledTimes(_state === "absent" ? 0 : 1); + } + }); +}); diff --git a/packages/core/src/__tests__/excluded-awaited-run-audit-store-sites.test.ts b/packages/core/src/__tests__/excluded-awaited-run-audit-store-sites.test.ts new file mode 100644 index 0000000000..b1622dfda6 --- /dev/null +++ b/packages/core/src/__tests__/excluded-awaited-run-audit-store-sites.test.ts @@ -0,0 +1,236 @@ +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import type { WorkflowStepResult } from "../types.js"; +import { BUILTIN_CODING_WORKFLOW_IR } from "../workflows/builtin-coding-workflow-ir.js"; +import type { WorkflowIrV2 } from "../workflows/workflow-ir-types.js"; +import { + pgDescribe, + createSharedPgTaskStoreTestHarness, + type SharedPgTaskStoreHarness, +} from "../__test-utils__/pg-test-harness.js"; + +/* + * FNXC:RunAudit 2026-08-20-06:40: + * FN-9178 characterizes current awaited audit behavior through public store entry points; it is + * not a remediation. These fixtures are PG-gated because TaskStore's durable methods require the + * production async layer, while hostile doubles prove the operation's observable ordering. + */ +pgDescribe("FN-9178 awaited store run-audit characterization", () => { + const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_awaited_audit" }); + beforeAll(h.beforeAll); + beforeEach(h.beforeEach); + afterEach(h.afterEach); + afterAll(h.afterAll); + + const failed = (): WorkflowStepResult => ({ + workflowStepId: "review", workflowStepName: "Review", phase: "pre-merge", status: "failed", + completedAt: "2026-08-20T00:00:00.000Z", + }); + const pending = (): WorkflowStepResult => ({ + workflowStepId: "review", workflowStepName: "Review", phase: "pre-merge", status: "pending", + }); + + async function seed(id: string, results: WorkflowStepResult[]) { + const store = h.store(); + await store.createTaskWithReservedId({ description: id, column: "in-review" }, { taskId: id, applyDefaultWorkflowSteps: false }); + await store.updateTask(id, { workflowStepResults: results }); + return store; + } + + it("bypass treats an absent/non-function audit result as a completed write", async () => { + const store = await seed("FN-BYP-ABSENT", [failed()]); + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(() => undefined as never); + const write = vi.spyOn(store as never, "atomicWriteTaskJson"); + write.mockClear(); + await expect(store.bypassFailedPreMergeReviewStep("FN-BYP-ABSENT", { reason: "test", actor: "operator" })).resolves.toBeDefined(); + expect(write).toHaveBeenCalledOnce(); + }); + + it.each([ + ["synchronous throw", () => { throw new Error("sync"); }], + ["rejection", () => Promise.reject(new Error("rejected"))], + ])("bypass rejects and does not persist after a %s audit sink", async (_kind, sink) => { + const id = `FN-BYP-${_kind.replace(/\W/g, "")}`; + const store = await seed(id, [failed()]); + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(sink as never); + const write = vi.spyOn(store as never, "atomicWriteTaskJson"); + write.mockClear(); + await expect(store.bypassFailedPreMergeReviewStep(id, { reason: "test", actor: "operator" })).rejects.toThrow(); + expect(write).not.toHaveBeenCalled(); + }); + + it("bypass remains pending for a never-settling audit sink", async () => { + const store = await seed("FN-BYP-PENDING", [failed()]); + vi.useFakeTimers(); + try { + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(() => new Promise(() => undefined)); + const write = vi.spyOn(store as never, "atomicWriteTaskJson"); + write.mockClear(); + let settled = false; + void store.bypassFailedPreMergeReviewStep("FN-BYP-PENDING", { reason: "test", actor: "operator" }).finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(settled).toBe(false); + expect(write).not.toHaveBeenCalled(); + } finally { vi.useRealTimers(); } + }); + + it("bypass waits for a late-settling audit before writing task.json", async () => { + let resolve!: () => void; + const store = await seed("FN-BYP-LATE", [failed()]); + const audit = vi.spyOn(store, "recordRunAuditEvent").mockImplementation(() => new Promise((done) => { resolve = done; })); + audit.mockClear(); + const write = vi.spyOn(store as never, "atomicWriteTaskJson"); + write.mockClear(); + const operation = store.bypassFailedPreMergeReviewStep("FN-BYP-LATE", { reason: "test", actor: "operator" }); + await vi.waitFor(() => expect(audit).toHaveBeenCalledOnce()); + expect(write).not.toHaveBeenCalled(); + resolve(); + await expect(operation).resolves.toMatchObject({ id: "FN-BYP-LATE" }); + expect(write).toHaveBeenCalledOnce(); + }); + + it.each([ + ["absent", () => undefined, true], + ["synchronous throw", () => { throw new Error("sync"); }, false], + ["rejection", () => Promise.reject(new Error("rejected")), false], + ])("resume-step's persistence follows a %s audit result", async (_state, sink, persists) => { + const id = `FN-RESUME-${_state.replace(/\W/g, "")}`; + const store = await seed(id, [pending()]); + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(sink as never); + const write = vi.spyOn(store as never, "atomicWriteTaskJson"); + write.mockClear(); + const operation = store.resumeWorkflowStep(id, { stepId: "review", reason: "test", actor: "operator" }); + if (persists) await expect(operation).resolves.toBeDefined(); + else await expect(operation).rejects.toThrow(); + expect(write).toHaveBeenCalledTimes(persists ? 1 : 0); + }); + + it.each([ + ["never-settling", () => new Promise(() => undefined)], + ["late-settling", () => new Promise((resolve) => setTimeout(resolve, 2_100))], + ])("resume-step cannot reach persistence while its %s audit is unresolved", async (_state, sink) => { + const id = `FN-RESUME-${_state}`; + const store = await seed(id, [pending()]); + vi.useFakeTimers(); + try { + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(sink as never); + const write = vi.spyOn(store as never, "atomicWriteTaskJson"); + write.mockClear(); + let settled = false; + void store.resumeWorkflowStep(id, { stepId: "review", reason: "test", actor: "operator" }).finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + // The late sink would have settled after the prospective bounded-window; an unbounded await + // either remains pending (never) or only proceeds once its real promise settles (late). + expect(write).not.toHaveBeenCalled(); + expect(settled).toBe(false); + } finally { vi.useRealTimers(); } + }); + + it("preserves WorkflowSwitchRehomeFailedError after all hostile audit outcomes", async () => { + const sourceIr = structuredClone(BUILTIN_CODING_WORKFLOW_IR) as WorkflowIrV2; + sourceIr.columns.push({ id: "audit-source", name: "Audit source", traits: [] }); + for (const [state, sink] of [ + ["absent", () => undefined], + ["sync", () => { throw new Error("sync"); }], + ["reject", () => Promise.reject(new Error("reject"))], + ] as const) { + const store = h.store(); + const source = await store.createWorkflowDefinition({ name: `source ${state}`, ir: sourceIr, layout: {} }); + const target = await store.createWorkflowDefinition({ name: `target ${state}`, ir: structuredClone(BUILTIN_CODING_WORKFLOW_IR) as WorkflowIrV2, layout: {} }); + const task = await store.createTask({ description: `switch ${state}` }); + await store.selectTaskWorkflow(task.id, source.id); + await store.moveTask(task.id, "audit-source", { moveSource: "engine", bypassGuards: true, recoveryRehome: true }); + vi.spyOn(store, "rehomeOccupant").mockResolvedValue({ moved: false, error: "race" } as never); + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(sink as never); + await expect(store.selectTaskWorkflowAndReconcile(task.id, target.id)).rejects.toMatchObject({ name: "WorkflowSwitchRehomeFailedError", committed: true }); + vi.restoreAllMocks(); + } + }); + + it("keeps workflow-switch torn waiting on a never-settling audit promise", async () => { + const store = h.store(); + const sourceIr = structuredClone(BUILTIN_CODING_WORKFLOW_IR) as WorkflowIrV2; + sourceIr.columns.push({ id: "audit-never", name: "Audit", traits: [] }); + const source = await store.createWorkflowDefinition({ name: "source never", ir: sourceIr, layout: {} }); + const target = await store.createWorkflowDefinition({ name: "target never", ir: structuredClone(BUILTIN_CODING_WORKFLOW_IR) as WorkflowIrV2, layout: {} }); + const task = await store.createTask({ description: "switch never" }); + await store.selectTaskWorkflow(task.id, source.id); + await store.moveTask(task.id, "audit-never", { moveSource: "engine", bypassGuards: true, recoveryRehome: true }); + vi.useFakeTimers(); + try { + vi.spyOn(store, "rehomeOccupant").mockResolvedValue({ moved: false } as never); + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(() => new Promise(() => undefined)); + let settled = false; + void store.selectTaskWorkflowAndReconcile(task.id, target.id).finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(settled).toBe(false); + } finally { vi.useRealTimers(); vi.restoreAllMocks(); } + }); + + it("throws WorkflowSwitchRehomeFailedError after a late audit settles beyond the bounded window", async () => { + const store = h.store(); + const sourceIr = structuredClone(BUILTIN_CODING_WORKFLOW_IR) as WorkflowIrV2; + sourceIr.columns.push({ id: "audit-late", name: "Audit", traits: [] }); + const source = await store.createWorkflowDefinition({ name: "source late", ir: sourceIr, layout: {} }); + const target = await store.createWorkflowDefinition({ name: "target late", ir: structuredClone(BUILTIN_CODING_WORKFLOW_IR) as WorkflowIrV2, layout: {} }); + const task = await store.createTask({ description: "switch late" }); + await store.selectTaskWorkflow(task.id, source.id); + await store.moveTask(task.id, "audit-late", { moveSource: "engine", bypassGuards: true, recoveryRehome: true }); + let resolveAudit!: () => void; + let signalAudit!: () => void; + const auditStarted = new Promise((resolve) => { signalAudit = resolve; }); + try { + vi.spyOn(store, "rehomeOccupant").mockResolvedValue({ moved: false } as never); + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(() => new Promise((resolve) => { + resolveAudit = resolve; + signalAudit(); + })); + const operation = store.selectTaskWorkflowAndReconcile(task.id, target.id); + void operation.catch(() => undefined); + await auditStarted; + vi.useFakeTimers(); + await vi.advanceTimersByTimeAsync(2_100); + resolveAudit(); + await expect(operation).rejects.toMatchObject({ name: "WorkflowSwitchRehomeFailedError", committed: true }); + } finally { vi.useRealTimers(); vi.restoreAllMocks(); } + }); + + it.each([ + ["absent", () => undefined, "reconciled"], + ["synchronous throw", () => { throw new Error("sync"); }, "audit-failed: sync"], + ["rejection", () => Promise.reject(new Error("rejected")), "audit-failed: rejected"], + ])("phantom reconciliation exposes %s audit outcome in its public result", async (_state, sink, expected) => { + const store = h.store(); + const task = await store.createTask({ description: `phantom ${_state}` }); + const layer = h.layer(); + const projectId = layer.projectId?.trim() || "__legacy_unscoped__"; + const { rm } = await import("node:fs/promises"); + const { join } = await import("node:path"); + const schema = await import("../postgres/schema/index.js"); + const { and, eq } = await import("drizzle-orm"); + await rm(join(h.rootDir(), ".fusion", "tasks", task.id), { recursive: true, force: true }); + await layer.db.delete(schema.project.tasks).where(and(eq(schema.project.tasks.projectId, projectId), eq(schema.project.tasks.id, task.id))); + await layer.db.insert(schema.project.activityLog).values({ projectId, id: `audit-${task.id}`, timestamp: new Date().toISOString(), type: "task:created", taskId: task.id, details: "orphan" }); + vi.spyOn(store, "recordRunAuditEvent").mockImplementation(sink as never); + const result = await store.reconcilePhantomCommittedReservations(); + if (expected === "reconciled") expect(result.reconciled).toContain(task.id); + else expect(result.skipped).toContainEqual({ id: task.id, reason: expected }); + }); + + it.each(["never", "late"])("phantom reconciliation remains pending when its %s audit does", async (state) => { + const store = h.store(); + const task = await store.createTask({ description: `phantom ${state}` }); + const layer = h.layer(); const projectId = layer.projectId?.trim() || "__legacy_unscoped__"; + const { rm } = await import("node:fs/promises"); const { join } = await import("node:path"); + const schema = await import("../postgres/schema/index.js"); const { and, eq } = await import("drizzle-orm"); + await rm(join(h.rootDir(), ".fusion", "tasks", task.id), { recursive: true, force: true }); + await layer.db.delete(schema.project.tasks).where(and(eq(schema.project.tasks.projectId, projectId), eq(schema.project.tasks.id, task.id))); + await layer.db.insert(schema.project.activityLog).values({ projectId, id: `audit-${task.id}`, timestamp: new Date().toISOString(), type: "task:created", taskId: task.id, details: "orphan" }); + vi.useFakeTimers(); + try { + vi.spyOn(store, "recordRunAuditEvent").mockImplementation((state === "never" ? () => new Promise(() => undefined) : () => new Promise((resolve) => setTimeout(resolve, 2_100))) as never); + let settled = false; void store.reconcilePhantomCommittedReservations().finally(() => { settled = true; }).catch(() => undefined); + await vi.advanceTimersByTimeAsync(2_100); + expect(settled).toBe(false); + } finally { vi.useRealTimers(); } + }); +});