From c8f6afe1246f0a6903708bd0d35d686c5892630b Mon Sep 17 00:00:00 2001 From: gsxdsm Date: Thu, 20 Aug 2026 00:36:50 -0700 Subject: [PATCH] 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) --- .../fn-9182-bounded-run-audit-outcome.md | 7 ++ AGENTS.md | 1 + docs/run-audit.md | 4 +- .../core-run-audit-emitter-isolation.test.ts | 10 +-- .../core-run-audit-sink-health.test.ts | 51 ++++++++++++- .../__tests__/emit-bounded-run-audit.test.ts | 58 +++++++++++++- .../phantom-reservation-audit-outcome.test.ts | 76 +++++++++++++++++++ .../src/run-audit/emit-bounded-run-audit.ts | 41 ++++++++-- .../async/async-phantom-reservations.ts | 31 +++++--- .../src/task-store/workflow-definitions.ts | 48 +++++++----- 10 files changed, 277 insertions(+), 50 deletions(-) create mode 100644 .changeset/fn-9182-bounded-run-audit-outcome.md create mode 100644 packages/core/src/__tests__/phantom-reservation-audit-outcome.test.ts diff --git a/.changeset/fn-9182-bounded-run-audit-outcome.md b/.changeset/fn-9182-bounded-run-audit-outcome.md new file mode 100644 index 0000000000..2cfb4f4f59 --- /dev/null +++ b/.changeset/fn-9182-bounded-run-audit-outcome.md @@ -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. diff --git a/AGENTS.md b/AGENTS.md index 31f069a69c..e1f6c26143 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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. - 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-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. diff --git a/docs/run-audit.md b/docs/run-audit.md index 65682638ab..2bcbb263c9 100644 --- a/docs/run-audit.md +++ b/docs/run-audit.md @@ -82,10 +82,10 @@ 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. +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 -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. 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 22ffbf352e..9641cea8f4 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 @@ -5,14 +5,13 @@ 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. + * exception. Class A sites are evaluated candidates, C retains its ordering claim, and + * 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 = { "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", @@ -35,7 +34,8 @@ 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", "../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)); diff --git a/packages/core/src/__tests__/core-run-audit-sink-health.test.ts b/packages/core/src/__tests__/core-run-audit-sink-health.test.ts index 7b6dc8d01f..6197ea12c2 100644 --- a/packages/core/src/__tests__/core-run-audit-sink-health.test.ts +++ b/packages/core/src/__tests__/core-run-audit-sink-health.test.ts @@ -1,6 +1,10 @@ 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", () => ({ 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", () => ({ - readTaskRow: vi.fn(async () => undefined), + readTaskRow: vi.fn(async () => workflowReader.current?.()), })); +vi.mock("../workflows/workflow-ir-resolver.js", async (importOriginal) => { + const actual = await importOriginal(); + return { ...actual, resolveWorkflowIrForTask: vi.fn(async () => resolvedWorkflowIr.current) }; +}); + vi.mock("../task-store/async/async-transition-pending.js", () => ({ listTransitionPendingTaskIdsAsync: vi.fn(async () => ["FN-9177"]), 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 { markLegacyAutoMergeStampsOnceImpl } from "../task-store/workflow-integrity.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 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); + 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", () => { 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) => { diff --git a/packages/core/src/__tests__/emit-bounded-run-audit.test.ts b/packages/core/src/__tests__/emit-bounded-run-audit.test.ts index a89a516376..17fc9ff915 100644 --- a/packages/core/src/__tests__/emit-bounded-run-audit.test.ts +++ b/packages/core/src/__tests__/emit-bounded-run-audit.test.ts @@ -1,5 +1,9 @@ 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: {} }; @@ -33,3 +37,55 @@ describe("emitBoundedRunAudit", () => { 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((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((_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(); + } + }); +}); diff --git a/packages/core/src/__tests__/phantom-reservation-audit-outcome.test.ts b/packages/core/src/__tests__/phantom-reservation-audit-outcome.test.ts new file mode 100644 index 0000000000..48b85ca7a5 --- /dev/null +++ b/packages/core/src/__tests__/phantom-reservation-audit-outcome.test.ts @@ -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((resolve) => { settle = resolve; }); + if (mode === "late-reject") return new Promise((_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); + }); +}); diff --git a/packages/core/src/run-audit/emit-bounded-run-audit.ts b/packages/core/src/run-audit/emit-bounded-run-audit.ts index 6ddcea45e5..c4b0a4c224 100644 --- a/packages/core/src/run-audit/emit-bounded-run-audit.ts +++ b/packages/core/src/run-audit/emit-bounded-run-audit.ts @@ -9,6 +9,12 @@ export type RunAuditSinkHost = { 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 }; 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, * 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, event: RunAuditEvent, options: { timeoutMs?: number; log?: RunAuditLogger } = {}, ): Promise { + 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 { const log = options.log ?? defaultLog; const sink = host?.recordRunAuditEvent; - if (typeof sink !== "function") return; + if (typeof sink !== "function") return { outcome: "absent" }; let sinkPromise: Promise; try { + // Invoke before the first await: synchronous timeline consumers observe every audit attempt. sinkPromise = Promise.resolve(sink.call(host, event as RunAuditEventInput)); - } catch { + } catch (error) { log.warn(`[run-audit] failed to record ${event.mutationType}`); - return; + return { outcome: "failed", error }; } void sinkPromise.catch(() => undefined); - await new Promise((resolve) => { + return new Promise((resolve) => { const timer = setTimeout(() => { log.warn(`[run-audit] timed out recording ${event.mutationType}`); - resolve(); + resolve({ outcome: "timed-out" }); }, options.timeoutMs ?? CORE_RUN_AUDIT_EMIT_TIMEOUT_MS); timer.unref?.(); void sinkPromise.then( - () => { clearTimeout(timer); resolve(); }, - () => { clearTimeout(timer); log.warn(`[run-audit] failed to record ${event.mutationType}`); resolve(); }, + () => { clearTimeout(timer); resolve({ outcome: "recorded" }); }, + (error) => { + clearTimeout(timer); + log.warn(`[run-audit] failed to record ${event.mutationType}`); + resolve({ outcome: "failed", error }); + }, ); }); } diff --git a/packages/core/src/task-store/async/async-phantom-reservations.ts b/packages/core/src/task-store/async/async-phantom-reservations.ts index 77b47f6b4f..d52f138cd5 100644 --- a/packages/core/src/task-store/async/async-phantom-reservations.ts +++ b/packages/core/src/task-store/async/async-phantom-reservations.ts @@ -4,6 +4,7 @@ import { and, asc, eq, inArray, isNull } from "drizzle-orm"; import { alias } from "drizzle-orm/pg-core"; import * as schema from "../../postgres/schema/index.js"; import type { TaskStore } from "../../store.js"; +import { emitBoundedRunAuditWithOutcome } from "../../run-audit/emit-bounded-run-audit.js"; export interface PhantomReservationReconcileResult { reconciled: string[]; @@ -144,24 +145,30 @@ export async function reconcilePhantomCommittedReservationsAsync( continue; } if (pruned.prunedActivityLog > 0 || pruned.prunedAgents > 0) { - try { - await store.recordRunAuditEvent({ - agentId: "self-healing", - runId: `phantom-reservation:${taskId}`, - taskId, - domain: "database", - mutationType: "task:reconcile-phantom-committed-reservation", - target: taskId, - metadata: { reservationStatus: "committed", ...pruned }, - }); - } catch (error) { + const auditResult = await emitBoundedRunAuditWithOutcome(store, { + agentId: "self-healing", + runId: `phantom-reservation:${taskId}`, + taskId, + domain: "database", + mutationType: "task:reconcile-phantom-committed-reservation", + target: taskId, + metadata: { reservationStatus: "committed", ...pruned }, + }); + if (auditResult.outcome === "failed" || auditResult.outcome === "timed-out") { /* 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. + + 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({ 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; } diff --git a/packages/core/src/task-store/workflow-definitions.ts b/packages/core/src/task-store/workflow-definitions.ts index 31ae09e43d..bf89680364 100644 --- a/packages/core/src/task-store/workflow-definitions.ts +++ b/packages/core/src/task-store/workflow-definitions.ts @@ -46,6 +46,10 @@ import { resolveWorkflowIrForTask } from "../workflows/workflow-ir-resolver.js"; import { acquireTaskAdvisoryXactLock } from "./task-advisory-lock.js"; import { resolveProjectColumnsForRoles, REVIEW_ROLES } from "../project-lane-vocabulary.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, 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 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 - could silently be absent. The write is awaited and its own failure is swallowed - so it can never mask the rejection the caller actually needs to see. + could silently be absent. The bounded outcome is still awaited before throwing, + but a hostile sink cannot stall the switch forever. NO ERROR PROSE (PR #2512 review — CodeRabbit). `outcome.error` is a propagated `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 prose. The bounded outcome code goes here; the human-readable reason travels on 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 { - await store.recordRunAuditEvent({ - taskId, - agentId: "system", - runId: `workflow-switch-torn-${taskId}`, - domain: "database", - mutationType: "task:workflow-switch-torn", - target: taskId, - metadata: { - workflowId, - fromColumn, - intendedColumn: decision.targetColumn, - selectionCommitted: true, - outcome: "rehome-rejected", - }, - }); - } catch { - // An audit-write failure must not replace the rejection being reported. + const auditResult = await emitBoundedRunAuditWithOutcome(store, { + taskId, + agentId: "system", + runId: `workflow-switch-torn-${taskId}`, + domain: "database", + mutationType: "task:workflow-switch-torn", + target: taskId, + metadata: { + workflowId, + fromColumn, + intendedColumn: decision.targetColumn, + selectionCommitted: true, + outcome: "rehome-rejected", + }, + }); + if (auditResult.outcome !== "recorded") { + workflowDefinitionLog.warn(`[workflow-switch] audit ${auditResult.outcome} task=${taskId} workflow=${workflowId}`); } throw new WorkflowSwitchRehomeFailedError({ taskId,