FN-9182: Bound reconciliation audit emissions

Keep reconciliation responsive while preserving audit-dependent outcomes.

- add a bounded core audit emitter that reports recorded, absent, failed, and timed-out outcomes
- route workflow-switch and phantom-reservation reconciliation through the outcome-aware seam
- cover hostile sink behavior, caller-visible failures, routing invariants, and operator documentation

Files changed:
 .changeset/fn-9182-bounded-run-audit-outcome.md    |  7 ++
 AGENTS.md                                          |  1 +
 docs/run-audit.md                                  |  4 +-
 .../core-run-audit-emitter-isolation.test.ts       | 10 +--
 .../__tests__/core-run-audit-sink-health.test.ts   | 51 ++++++++++++++-
 .../src/__tests__/emit-bounded-run-audit.test.ts   | 58 ++++++++++++++++-
 .../phantom-reservation-audit-outcome.test.ts      | 76 ++++++++++++++++++++++
 .../core/src/run-audit/emit-bounded-run-audit.ts   | 41 +++++++++---
 .../task-store/async/async-phantom-reservations.ts | 31 +++++----
 .../core/src/task-store/workflow-definitions.ts    | 48 ++++++++------
 10 files changed, 277 insertions(+), 50 deletions(-)

Fusion-Task-Id: FN-9182

Fusion-Task-Lineage: c98e8e5e-2dae-47a0-b969-4156f4beccb3

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-20 00:36:50 -07:00
parent 649b901bbd
commit c8f6afe124
10 changed files with 277 additions and 50 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Keep workflow recovery and reservation cleanup responsive when audit logging stalls.
category: fix
dev: Adds emitBoundedRunAuditWithOutcome for workflow-switch-torn and phantom-reservation reconciliation.

View File

@@ -293,6 +293,7 @@ canonical emitters remain explicit exclusions until their separately scoped hard
--> -->
- FN-9175: New engine run-audit emitters must use `emitBoundedRunAudit` from `packages/engine/src/util/emit-bounded-run-audit.ts`; it absorbs absent, throwing, rejecting, hanging, and late-settling sinks without changing the owning branch, and requires behavioral sink-health coverage. - FN-9175: New engine run-audit emitters must use `emitBoundedRunAudit` from `packages/engine/src/util/emit-bounded-run-audit.ts`; it absorbs absent, throwing, rejecting, hanging, and late-settling sinks without changing the owning branch, and requires behavioral sink-health coverage.
- FNXC:RunAudit 2026-08-20-05:49: FN-9177 requires new core best-effort emitters to use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. It deliberately mirrors the engine seam because core cannot import engine; transactional and deliberately awaited durability writers remain unbounded. - FNXC:RunAudit 2026-08-20-05:49: FN-9177 requires new core best-effort emitters to use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. It deliberately mirrors the engine seam because core cannot import engine; transactional and deliberately awaited durability writers remain unbounded.
- FNXC:RunAudit 2026-08-20-07:14: FN-9182 requires core emitters whose behavior branches on audit success to use `emitBoundedRunAuditWithOutcome`; `emitBoundedRunAudit` remains the default best-effort seam, while transactional and deliberately awaited durability writers remain unbounded.
- FN-9109: `session:cross-runtime-fallback-engaged` records a single retryable-failure handoff from a primary runtime to a deferred CLI runtime. Metadata is ids/outcomes-only (`sessionPurpose`, primary/fallback provider and model IDs, trigger point, failure category, `contextTransferred`); never record error prose or transferred transcript text. - FN-9109: `session:cross-runtime-fallback-engaged` records a single retryable-failure handoff from a primary runtime to a deferred CLI runtime. Metadata is ids/outcomes-only (`sessionPurpose`, primary/fallback provider and model IDs, trigger point, failure category, `contextTransferred`); never record error prose or transferred transcript text.
- FN-8958: `merge:orphan-write-fenced` is emitted once per orphan merge body at its fence's first interaction. Metadata is ids/counts/outcomes-only: `{ taskId, category, interaction, suppressedCount }`; `suppressedCount` is the emit-time count (`1` for `interaction:"suppressed"`, `0` for `interaction:"rejected"`), never a cumulative body total. - FN-8958: `merge:orphan-write-fenced` is emitted once per orphan merge body at its fence's first interaction. Metadata is ids/counts/outcomes-only: `{ taskId, category, interaction, suppressedCount }`; `suppressedCount` is the emit-time count (`1` for `interaction:"suppressed"`, `0` for `interaction:"rejected"`), never a cumulative body total.

View File

@@ -82,10 +82,10 @@ This applies to executor, run-auditor, self-healing, merger, PR reconciliation,
### Core emit-seam policy ### 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. 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. `emitBoundedRunAudit` is the default void seam. `emitBoundedRunAuditWithOutcome` returns `recorded`, `absent`, `failed` (with the original error), or `timed-out` where a forensic throw ordering or caller-visible skipped payload depends on the audit result; workflow-switch torn reconciliation and phantom committed-reservation reconciliation use it. 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 ### 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. 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 and use the bounded outcome seam 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. 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.

View File

@@ -5,14 +5,13 @@ import { describe, expect, it } from "vitest";
/* /*
* FNXC:RunAudit 2026-08-20-07:16: * FNXC:RunAudit 2026-08-20-07:16:
* FN-9178 makes every direct awaited core audit writer a named decision rather than an implicit * 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 * exception. Class A sites are evaluated candidates, C retains its ordering claim, and
* its ordering claim, and transactional/sink writers remain permanent atomicity boundaries. * transactional/sink writers remain permanent atomicity boundaries. FN-9182 migrated class B
* sites to the reporting bounded seam, so they no longer appear in this direct-await inventory.
*/ */
const awaitedClassifications = { const awaitedClassifications = {
"store.ts:task:bypass-review": "C", "store.ts:task:bypass-review": "C",
"store.ts:task:resume-step": "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-creation.ts:intake:resurrection-blocked": "C",
"task-store/task-id-integrity.ts:task: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:catch-up": "A",
@@ -35,7 +34,8 @@ const files = [
"../planner/planner-intervention.ts", "../task-store/audit-ops.ts", "../task-store/branch-group-ops.ts", "../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", "../task-store/merge-queue-ops-2.ts", "../task-store/task-mutation-ops.ts", "../task-store/workflow-integrity.ts",
"../task-store/workflow-workitems-ops.ts", "../task-store/workflow-workitems-ops-2.ts", "../task-store/task-artifacts-ops.ts", "../task-store/workflow-workitems-ops.ts", "../task-store/workflow-workitems-ops-2.ts", "../task-store/task-artifacts-ops.ts",
"../task-store/lifecycle-ops.ts", "../task-store/task-id-integrity.ts", "../task-store/lifecycle-ops.ts", "../task-store/task-id-integrity.ts", "../task-store/workflow-definitions.ts",
"../task-store/async/async-phantom-reservations.ts",
]; ];
const sourceRoot = fileURLToPath(new URL("..", import.meta.url)); const sourceRoot = fileURLToPath(new URL("..", import.meta.url));

View File

@@ -1,6 +1,10 @@
import { afterEach, describe, expect, it, vi } from "vitest"; import { afterEach, describe, expect, it, vi } from "vitest";
const { asyncAuditSink } = vi.hoisted(() => ({ asyncAuditSink: { current: undefined as undefined | ((event: unknown) => unknown) } })); const { asyncAuditSink, workflowReader, resolvedWorkflowIr } = vi.hoisted(() => ({
asyncAuditSink: { current: undefined as undefined | ((event: unknown) => unknown) },
workflowReader: { current: undefined as undefined | (() => unknown) },
resolvedWorkflowIr: { current: undefined as unknown },
}));
vi.mock("../task-store/async/async-audit.js", () => ({ vi.mock("../task-store/async/async-audit.js", () => ({
recordRunAuditEvent: (_layer: unknown, event: unknown) => asyncAuditSink.current?.(event), recordRunAuditEvent: (_layer: unknown, event: unknown) => asyncAuditSink.current?.(event),
})); }));
@@ -11,9 +15,14 @@ vi.mock("../postgres/data-layer.js", async (importOriginal) => {
}); });
vi.mock("../task-store/async/async-persistence.js", () => ({ vi.mock("../task-store/async/async-persistence.js", () => ({
readTaskRow: vi.fn(async () => undefined), readTaskRow: vi.fn(async () => workflowReader.current?.()),
})); }));
vi.mock("../workflows/workflow-ir-resolver.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("../workflows/workflow-ir-resolver.js")>();
return { ...actual, resolveWorkflowIrForTask: vi.fn(async () => resolvedWorkflowIr.current) };
});
vi.mock("../task-store/async/async-transition-pending.js", () => ({ vi.mock("../task-store/async/async-transition-pending.js", () => ({
listTransitionPendingTaskIdsAsync: vi.fn(async () => ["FN-9177"]), listTransitionPendingTaskIdsAsync: vi.fn(async () => ["FN-9177"]),
readTransitionPendingAsync: vi.fn(async () => ({ toColumn: "todo", hooksRemaining: ["default-workflow:postCommit"], startedAt: 1 })), readTransitionPendingAsync: vi.fn(async () => ({ toColumn: "todo", hooksRemaining: ["default-workflow:postCommit"], startedAt: 1 })),
@@ -57,6 +66,8 @@ import { projectMergeRequestToWorkflowWorkItemImpl } from "../task-store/workflo
import { acquireWorkflowWorkItemLeaseImpl } from "../task-store/workflow-workitems-ops-2.js"; import { acquireWorkflowWorkItemLeaseImpl } from "../task-store/workflow-workitems-ops-2.js";
import { markLegacyAutoMergeStampsOnceImpl } from "../task-store/workflow-integrity.js"; import { markLegacyAutoMergeStampsOnceImpl } from "../task-store/workflow-integrity.js";
import { getTraitRegistry } from "../workflows/trait-registry.js"; import { getTraitRegistry } from "../workflows/trait-registry.js";
import { WorkflowSwitchRehomeFailedError } from "../workflows/workflow-reconciliation.js";
import { selectTaskWorkflowAndReconcileImpl } from "../task-store/workflow-definitions.js";
const input = { taskId: "FN-9177", stage: "executor" as const, reason: "test", timestamp: "2026-08-20T00:00:00.000Z" }; const input = { taskId: "FN-9177", stage: "executor" as const, reason: "test", timestamp: "2026-08-20T00:00:00.000Z" };
const façades = [ const façades = [
@@ -258,6 +269,42 @@ describe("core audit emitters tolerate hostile sinks", () => {
}); });
describe("workflow-switch torn reconciliation retains its forensic rejection", () => {
afterEach(() => {
workflowReader.current = undefined;
resolvedWorkflowIr.current = undefined;
});
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"throws the committed torn error for %s audit sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
let readCount = 0;
workflowReader.current = () => (readCount++ === 0 ? undefined : { id: "FN-9182", column: "todo" });
resolvedWorkflowIr.current = { version: 1, nodes: [], edges: [], columns: [{ id: "inbox", name: "Inbox", traits: [] }] };
const owner = {
...hostile.store,
asyncLayer: {},
selectTaskWorkflow: vi.fn(async () => []),
rehomeOccupant: vi.fn(async () => ({ moved: false })),
};
const operation = selectTaskWorkflowAndReconcileImpl(owner as never, "FN-9182", "WF-9182");
const rejection = expect(operation).rejects.toMatchObject({
name: "WorkflowSwitchRehomeFailedError",
taskId: "FN-9182",
workflowId: "WF-9182",
fromColumn: "todo",
intendedColumn: "inbox",
committed: true,
} satisfies Partial<WorkflowSwitchRehomeFailedError>);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
await rejection;
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
},
);
});
describe("remaining production owners retain hostile-sink isolation", () => { describe("remaining production owners retain hostile-sink isolation", () => {
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)( it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps plugin trait degradation and merge-request projection non-blocking for %s sinks", async (mode) => { "keeps plugin trait degradation and merge-request projection non-blocking for %s sinks", async (mode) => {

View File

@@ -1,5 +1,9 @@
import { afterEach, describe, expect, it, vi } from "vitest"; import { afterEach, describe, expect, it, vi } from "vitest";
import { CORE_RUN_AUDIT_EMIT_TIMEOUT_MS, emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js"; import {
CORE_RUN_AUDIT_EMIT_TIMEOUT_MS,
emitBoundedRunAudit,
emitBoundedRunAuditWithOutcome,
} from "../run-audit/emit-bounded-run-audit.js";
const event = { taskId: "FN-1", agentId: "system", runId: "run", domain: "database" as const, mutationType: "test:audit", target: "FN-1", metadata: {} }; const event = { taskId: "FN-1", agentId: "system", runId: "run", domain: "database" as const, mutationType: "test:audit", target: "FN-1", metadata: {} };
@@ -33,3 +37,55 @@ describe("emitBoundedRunAudit", () => {
await expect(late).resolves.toBeUndefined(); await expect(late).resolves.toBeUndefined();
}); });
}); });
describe("emitBoundedRunAuditWithOutcome", () => {
it("classifies absent, successful, and failed sinks with original errors", async () => {
const log = { warn: vi.fn() };
const thrown = new Error("thrown");
const rejected = new Error("rejected");
await expect(emitBoundedRunAuditWithOutcome(undefined, event, { log })).resolves.toEqual({ outcome: "absent" });
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: 1 as never }, event, { log })).resolves.toEqual({ outcome: "absent" });
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => undefined }, event, { log })).resolves.toEqual({ outcome: "recorded" });
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => { throw thrown; } }, event, { log })).resolves.toEqual({ outcome: "failed", error: thrown });
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => Promise.reject(rejected) }, event, { log })).resolves.toEqual({ outcome: "failed", error: rejected });
});
it("bounds never and late-settling sinks with explicit outcomes", async () => {
vi.useFakeTimers();
const log = { warn: vi.fn() };
const never = emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => new Promise(() => undefined) }, event, { log });
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
await expect(never).resolves.toEqual({ outcome: "timed-out" });
let resolveLate!: () => void;
const lateResolve = emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => new Promise<void>((resolve) => { resolveLate = resolve; }) }, event, { timeoutMs: 1, log });
await vi.advanceTimersByTimeAsync(1);
resolveLate();
await expect(lateResolve).resolves.toEqual({ outcome: "timed-out" });
let rejectLate!: (error: Error) => void;
const lateReject = emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => new Promise<void>((_resolve, reject) => { rejectLate = reject; }) }, event, { timeoutMs: 1, log });
await vi.advanceTimersByTimeAsync(1);
rejectLate(new Error("late rejection"));
await expect(lateReject).resolves.toEqual({ outcome: "timed-out" });
});
it("still invokes the void seam synchronously and resolves undefined for every sink state", async () => {
vi.useFakeTimers();
const modes = [
undefined,
() => undefined,
() => { throw new Error("throw"); },
() => Promise.reject(new Error("reject")),
() => new Promise(() => undefined),
];
for (const recordRunAuditEvent of modes) {
const sink = vi.fn(recordRunAuditEvent);
const promise = emitBoundedRunAudit(recordRunAuditEvent === undefined ? undefined : { recordRunAuditEvent: sink }, event, { log: { warn: vi.fn() } });
if (recordRunAuditEvent !== undefined) expect(sink).toHaveBeenCalledTimes(1);
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
await expect(promise).resolves.toBeUndefined();
}
});
});

View File

@@ -0,0 +1,76 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import { CORE_RUN_AUDIT_EMIT_TIMEOUT_MS } from "../run-audit/emit-bounded-run-audit.js";
import { reconcilePhantomCommittedReservationsAsync } from "../task-store/async/async-phantom-reservations.js";
type SinkMode = "absent" | "recorded" | "throw" | "reject" | "never" | "late-resolve" | "late-reject";
function createFixture(mode: SinkMode) {
let settle: (() => void) | undefined;
const sink = vi.fn(() => {
if (mode === "throw") throw new Error("forced audit failure");
if (mode === "reject") return Promise.reject(new Error("forced audit failure"));
if (mode === "never") return new Promise(() => undefined);
if (mode === "late-resolve") return new Promise<void>((resolve) => { settle = resolve; });
if (mode === "late-reject") return new Promise<void>((_resolve, reject) => { settle = () => reject(new Error("late audit failure")); });
return undefined;
});
const ids = ["FN-9182-a", "FN-9182-b"];
const reservations = ids.map((taskId) => ({ taskId, liveId: null, projectArchiveId: null, coldArchiveId: null }));
const query = {
from: vi.fn(), leftJoin: vi.fn(), where: vi.fn(), orderBy: vi.fn(async () => reservations),
};
query.from.mockReturnValue(query); query.leftJoin.mockReturnValue(query); query.where.mockReturnValue(query);
const txQuery = { from: vi.fn(), leftJoin: vi.fn(), where: vi.fn(async () => ids.map((taskId) => ({ taskId }))) };
txQuery.from.mockReturnValue(txQuery); txQuery.leftJoin.mockReturnValue(txQuery);
let deleteCount = 0;
const tx = {
select: vi.fn(() => txQuery),
delete: vi.fn(() => ({ where: vi.fn(() => ({ returning: vi.fn(async () => {
deleteCount += 1;
return deleteCount === 1 ? ids.map((taskId) => ({ taskId })) : [];
}) })) })),
};
const layer = {
projectId: "project-9182",
db: { select: vi.fn(() => query) },
transactionImmediate: async (work: (value: typeof tx) => unknown) => work(tx),
};
return {
store: {
getAsyncLayer: () => layer,
taskDir: (id: string) => `/definitely-missing-fn-9182/${id}`,
...(mode === "absent" ? {} : { recordRunAuditEvent: sink }),
},
sink,
settle: () => settle?.(),
ids,
};
}
afterEach(() => vi.useRealTimers());
describe("phantom reservation audit outcomes", () => {
it.each(["absent", "recorded", "throw", "reject"] as const)("preserves reconciliation payloads for %s sinks", async (mode) => {
const fixture = createFixture(mode);
const result = await reconcilePhantomCommittedReservationsAsync(fixture.store as never);
if (mode === "throw" || mode === "reject") {
expect(result).toEqual({ reconciled: [], skipped: fixture.ids.map((id) => ({ id, reason: "audit-failed: forced audit failure" })) });
} else {
expect(result).toEqual({ reconciled: fixture.ids, skipped: [] });
}
if (mode !== "absent") expect(fixture.sink).toHaveBeenCalledTimes(2);
});
it.each(["never", "late-resolve", "late-reject"] as const)("bounds %s audit sinks as per-reservation skipped results", async (mode) => {
vi.useFakeTimers();
const fixture = createFixture(mode);
const operation = reconcilePhantomCommittedReservationsAsync(fixture.store as never);
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS * fixture.ids.length);
fixture.settle();
await expect(operation).resolves.toEqual({
reconciled: [],
skipped: fixture.ids.map((id) => ({ id, reason: "audit-failed: timed-out" })),
});
expect(fixture.sink).toHaveBeenCalledTimes(2);
});
});

View File

@@ -9,6 +9,12 @@ export type RunAuditSinkHost = {
export type RunAuditLogger = { warn: (message: string) => void }; export type RunAuditLogger = { warn: (message: string) => void };
export type BoundedRunAuditResult =
| { outcome: "recorded" }
| { outcome: "absent" }
| { outcome: "failed"; error: unknown }
| { outcome: "timed-out" };
type RunAuditEvent = RunAuditEventInput | { mutationType: string; [key: string]: unknown }; type RunAuditEvent = RunAuditEventInput | { mutationType: string; [key: string]: unknown };
const defaultLog = createLogger("run-audit"); const defaultLog = createLogger("run-audit");
@@ -19,33 +25,52 @@ const defaultLog = createLogger("run-audit");
* without creating its documented dependency cycle. Best-effort telemetry must not block, reject, * without creating its documented dependency cycle. Best-effort telemetry must not block, reject,
* or otherwise alter the lifecycle operation which emitted it; no retry, queue, or backoff belongs here. * or otherwise alter the lifecycle operation which emitted it; no retry, queue, or backoff belongs here.
*/ */
export async function emitBoundedRunAudit( export function emitBoundedRunAudit(
host: RunAuditSinkHost, host: RunAuditSinkHost,
event: RunAuditEvent, event: RunAuditEvent,
options: { timeoutMs?: number; log?: RunAuditLogger } = {}, options: { timeoutMs?: number; log?: RunAuditLogger } = {},
): Promise<void> { ): Promise<void> {
return emitBoundedRunAuditWithOutcome(host, event, options).then(() => undefined);
}
/**
* FNXC:RunAudit 2026-08-20-07:14:
* FN-9182 supplies an explicit result for callers whose forensic throw ordering or caller-visible
* skipped payload derives from whether audit recording landed. This remains best-effort telemetry,
* not a durability guarantee: it bounds hostile sinks while preserving their original failure value.
*/
export async function emitBoundedRunAuditWithOutcome(
host: RunAuditSinkHost,
event: RunAuditEvent,
options: { timeoutMs?: number; log?: RunAuditLogger } = {},
): Promise<BoundedRunAuditResult> {
const log = options.log ?? defaultLog; const log = options.log ?? defaultLog;
const sink = host?.recordRunAuditEvent; const sink = host?.recordRunAuditEvent;
if (typeof sink !== "function") return; if (typeof sink !== "function") return { outcome: "absent" };
let sinkPromise: Promise<unknown>; let sinkPromise: Promise<unknown>;
try { try {
// Invoke before the first await: synchronous timeline consumers observe every audit attempt.
sinkPromise = Promise.resolve(sink.call(host, event as RunAuditEventInput)); sinkPromise = Promise.resolve(sink.call(host, event as RunAuditEventInput));
} catch { } catch (error) {
log.warn(`[run-audit] failed to record ${event.mutationType}`); log.warn(`[run-audit] failed to record ${event.mutationType}`);
return; return { outcome: "failed", error };
} }
void sinkPromise.catch(() => undefined); void sinkPromise.catch(() => undefined);
await new Promise<void>((resolve) => { return new Promise<BoundedRunAuditResult>((resolve) => {
const timer = setTimeout(() => { const timer = setTimeout(() => {
log.warn(`[run-audit] timed out recording ${event.mutationType}`); log.warn(`[run-audit] timed out recording ${event.mutationType}`);
resolve(); resolve({ outcome: "timed-out" });
}, options.timeoutMs ?? CORE_RUN_AUDIT_EMIT_TIMEOUT_MS); }, options.timeoutMs ?? CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
timer.unref?.(); timer.unref?.();
void sinkPromise.then( void sinkPromise.then(
() => { clearTimeout(timer); resolve(); }, () => { clearTimeout(timer); resolve({ outcome: "recorded" }); },
() => { clearTimeout(timer); log.warn(`[run-audit] failed to record ${event.mutationType}`); resolve(); }, (error) => {
clearTimeout(timer);
log.warn(`[run-audit] failed to record ${event.mutationType}`);
resolve({ outcome: "failed", error });
},
); );
}); });
} }

View File

@@ -4,6 +4,7 @@ import { and, asc, eq, inArray, isNull } from "drizzle-orm";
import { alias } from "drizzle-orm/pg-core"; import { alias } from "drizzle-orm/pg-core";
import * as schema from "../../postgres/schema/index.js"; import * as schema from "../../postgres/schema/index.js";
import type { TaskStore } from "../../store.js"; import type { TaskStore } from "../../store.js";
import { emitBoundedRunAuditWithOutcome } from "../../run-audit/emit-bounded-run-audit.js";
export interface PhantomReservationReconcileResult { export interface PhantomReservationReconcileResult {
reconciled: string[]; reconciled: string[];
@@ -144,24 +145,30 @@ export async function reconcilePhantomCommittedReservationsAsync(
continue; continue;
} }
if (pruned.prunedActivityLog > 0 || pruned.prunedAgents > 0) { if (pruned.prunedActivityLog > 0 || pruned.prunedAgents > 0) {
try { const auditResult = await emitBoundedRunAuditWithOutcome(store, {
await store.recordRunAuditEvent({ agentId: "self-healing",
agentId: "self-healing", runId: `phantom-reservation:${taskId}`,
runId: `phantom-reservation:${taskId}`, taskId,
taskId, domain: "database",
domain: "database", mutationType: "task:reconcile-phantom-committed-reservation",
mutationType: "task:reconcile-phantom-committed-reservation", target: taskId,
target: taskId, metadata: { reservationStatus: "committed", ...pruned },
metadata: { reservationStatus: "committed", ...pruned }, });
}); if (auditResult.outcome === "failed" || auditResult.outcome === "timed-out") {
} catch (error) {
/* /*
FNXC:PostgresReservationRecovery 2026-07-14-21:55: FNXC:PostgresReservationRecovery 2026-07-14-21:55:
Audit emission is isolated per reconciled reservation. One failed audit must not relabel earlier successful reconciliations as skipped or prevent later IDs from completing their own bookkeeping. Audit emission is isolated per reconciled reservation. One failed audit must not relabel earlier successful reconciliations as skipped or prevent later IDs from completing their own bookkeeping.
FNXC:RunAudit 2026-08-20-07:14:
FN-9182 bounds this awaited audit and classifies its outcome. Preserve the
audit-failed prefix because callers already expose it per reservation; timeout
is a new bounded failure state rather than an indefinitely stalled sweep.
*/ */
result.skipped.push({ result.skipped.push({
id: taskId, id: taskId,
reason: `audit-failed: ${error instanceof Error ? error.message : String(error)}`, reason: auditResult.outcome === "timed-out"
? "audit-failed: timed-out"
: `audit-failed: ${auditResult.error instanceof Error ? auditResult.error.message : String(auditResult.error)}`,
}); });
continue; continue;
} }

View File

@@ -46,6 +46,10 @@ import { resolveWorkflowIrForTask } from "../workflows/workflow-ir-resolver.js";
import { acquireTaskAdvisoryXactLock } from "./task-advisory-lock.js"; import { acquireTaskAdvisoryXactLock } from "./task-advisory-lock.js";
import { resolveProjectColumnsForRoles, REVIEW_ROLES } from "../project-lane-vocabulary.js"; import { resolveProjectColumnsForRoles, REVIEW_ROLES } from "../project-lane-vocabulary.js";
import type { InReviewDurationLanes } from "./async/async-audit.js"; import type { InReviewDurationLanes } from "./async/async-audit.js";
import { createLogger } from "../process/logger.js";
import { emitBoundedRunAuditWithOutcome } from "../run-audit/emit-bounded-run-audit.js";
const workflowDefinitionLog = createLogger("workflow-definitions");
export async function getAgentLogsByTimeRangeImpl(store: TaskStore, export async function getAgentLogsByTimeRangeImpl(store: TaskStore,
taskId: string, taskId: string,
@@ -849,33 +853,37 @@ export async function selectTaskWorkflowAndReconcileImpl(store: TaskStore,
AWAITED, not fire-and-forget (PR #2512 review — CodeRabbit). The whole claim of AWAITED, not fire-and-forget (PR #2512 review — CodeRabbit). The whole claim of
this branch is that the divergence is a fact on disk; `void`-ing the write and this branch is that the divergence is a fact on disk; `void`-ing the write and
throwing on the next line meant the one artifact self-healing is meant to find throwing on the next line meant the one artifact self-healing is meant to find
could silently be absent. The write is awaited and its own failure is swallowed could silently be absent. The bounded outcome is still awaited before throwing,
so it can never mask the rejection the caller actually needs to see. but a hostile sink cannot stall the switch forever.
NO ERROR PROSE (PR #2512 review — CodeRabbit). `outcome.error` is a propagated NO ERROR PROSE (PR #2512 review — CodeRabbit). `outcome.error` is a propagated
`err.message`; persisting it would contradict this file's own "ids/columns/ `err.message`; persisting it would contradict this file's own "ids/columns/
outcomes-only" claim and the project rule that run-audit never stores error outcomes-only" claim and the project rule that run-audit never stores error
prose. The bounded outcome code goes here; the human-readable reason travels on prose. The bounded outcome code goes here; the human-readable reason travels on
the thrown error, which is not persisted. the thrown error, which is not persisted.
FNXC:RunAudit 2026-08-20-07:14:
FN-9182 bounds the awaited forensic emission and reports its outcome without
changing this branch's typed rejection. A timeout is observable in logs rather
than indefinitely hiding a torn workflow switch from its caller.
*/ */
try { const auditResult = await emitBoundedRunAuditWithOutcome(store, {
await store.recordRunAuditEvent({ taskId,
taskId, agentId: "system",
agentId: "system", runId: `workflow-switch-torn-${taskId}`,
runId: `workflow-switch-torn-${taskId}`, domain: "database",
domain: "database", mutationType: "task:workflow-switch-torn",
mutationType: "task:workflow-switch-torn", target: taskId,
target: taskId, metadata: {
metadata: { workflowId,
workflowId, fromColumn,
fromColumn, intendedColumn: decision.targetColumn,
intendedColumn: decision.targetColumn, selectionCommitted: true,
selectionCommitted: true, outcome: "rehome-rejected",
outcome: "rehome-rejected", },
}, });
}); if (auditResult.outcome !== "recorded") {
} catch { workflowDefinitionLog.warn(`[workflow-switch] audit ${auditResult.outcome} task=${taskId} workflow=${workflowId}`);
// An audit-write failure must not replace the rejection being reported.
} }
throw new WorkflowSwitchRehomeFailedError({ throw new WorkflowSwitchRehomeFailedError({
taskId, taskId,