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) <noreply@runfusion.ai>
This commit is contained in:
@@ -430,3 +430,4 @@ Note: the embedded main-content views Workflows (`_WorkflowEditorView`), Import
|
||||
```
|
||||
<!-- FNXC:RunAudit 2026-08-20-05:49: FN-9177 requires new core best-effort emitters to use the core-owned bounded seam. -->
|
||||
- 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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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<typeof import("../postgres/data-layer.js")>()),
|
||||
recordRunAuditEvent: audit,
|
||||
}));
|
||||
vi.mock("../task-store/async/async-audit.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../task-store/async/async-audit.js")>()),
|
||||
recordRunAuditEvent: asyncAudit,
|
||||
}));
|
||||
vi.mock("../task-store/async/async-persistence.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../task-store/async/async-persistence.js")>()),
|
||||
softDeleteTaskRow: softDelete,
|
||||
readTaskRow,
|
||||
}));
|
||||
vi.mock("../task-store/task-lifecycle-consumer-registry.js", async (importOriginal) => ({
|
||||
...(await importOriginal<typeof import("../task-store/task-lifecycle-consumer-registry.js")>()),
|
||||
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<string, unknown> = {};
|
||||
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<never>(() => 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<void>((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<never>(() => undefined)],
|
||||
["late-settling", () => new Promise<void>((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<never>(() => 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<void>((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<never>(() => undefined) : () => new Promise<void>((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<never>(() => 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<void>((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<never>(() => undefined) : () => new Promise<void>((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<never>(() => undefined)],
|
||||
["late-settling", () => new Promise<void>((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);
|
||||
}
|
||||
});
|
||||
});
|
||||
@@ -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<never>(() => 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<void>((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<never>(() => undefined)],
|
||||
["late-settling", () => new Promise<void>((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<never>(() => 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<void>((resolve) => { signalAudit = resolve; });
|
||||
try {
|
||||
vi.spyOn(store, "rehomeOccupant").mockResolvedValue({ moved: false } as never);
|
||||
vi.spyOn(store, "recordRunAuditEvent").mockImplementation(() => new Promise<void>((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<never>(() => undefined) : () => new Promise<void>((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(); }
|
||||
});
|
||||
});
|
||||
Reference in New Issue
Block a user