FN-9177: route core run-audit emitters through bounded seam

Make best-effort core telemetry unable to block or alter lifecycle operations.

- Add a bounded core run-audit helper that absorbs missing, throwing, rejecting, hanging, and late-settling sinks.
- Route planner and task-store best-effort emitters through the helper while preserving transactional audit writes.
- Add sink-health, source-routing, and helper regressions plus operator-facing policy documentation and a release changeset.

Files changed:
 .changeset/fn-9177-core-run-audit.md               |   7 +
 AGENTS.md                                          |   5 +-
 docs/run-audit.md                                  |   6 +-
 .../core-run-audit-emitter-isolation.test.ts       |  21 ++
 .../__tests__/core-run-audit-sink-health.test.ts   | 317 +++++++++++++++++++++
 .../src/__tests__/emit-bounded-run-audit.test.ts   |  35 +++
 packages/core/src/planner/planner-intervention.ts  |   9 +-
 .../core/src/planner/planner-overseer-events.ts    |  15 +-
 .../core/src/run-audit/emit-bounded-run-audit.ts   |  51 ++++
 packages/core/src/task-store/audit-ops.ts          |   6 +-
 packages/core/src/task-store/branch-group-ops.ts   |   6 +-
 packages/core/src/task-store/lifecycle-ops.ts      |   4 +-
 packages/core/src/task-store/merge-queue-ops-2.ts  |   4 +-
 packages/core/src/task-store/task-artifacts-ops.ts |   4 +-
 packages/core/src/task-store/task-id-integrity.ts  |  52 ++--
 packages/core/src/task-store/task-mutation-ops.ts  |  13 +-
 packages/core/src/task-store/workflow-integrity.ts |   4 +-
 .../src/task-store/workflow-workitems-ops-2.ts     |   4 +-
 .../core/src/task-store/workflow-workitems-ops.ts  |   4 +-
 19 files changed, 512 insertions(+), 55 deletions(-)

Fusion-Task-Id: FN-9177

Fusion-Task-Lineage: 95082837-1670-4c25-9503-60ad983b9d68

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-19 23:29:37 -07:00
parent 6d51ae84c3
commit f5192a5d15
19 changed files with 512 additions and 55 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Prevent optional core audit sinks from delaying task lifecycle operations.
category: fix
dev: Core best-effort run-audit emitters now use a bounded, non-rejecting seam.

View File

@@ -292,6 +292,7 @@ rotation, and workflow-column boundaries. The bespoke merge-write fence and pack
canonical emitters remain explicit exclusions until their separately scoped hardening work lands.
-->
- 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.
- 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.
@@ -426,4 +427,6 @@ Note: the embedded main-content views Workflows (`_WorkflowEditorView`), Import
FNXC:SettingsNavigation 2026-05-13-08:11:
The modal should be 20% wider than the first section-sidebar layout and use a taller viewport so more settings remain visible without scrolling.
*/
```
```
<!-- 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.

View File

@@ -78,4 +78,8 @@ Adding a new catalogued run-audit event requires updating **both** the typed cat
All engine telemetry must use `emitBoundedRunAudit` from `packages/engine/src/util/emit-bounded-run-audit.ts`. It is best-effort and never load-bearing for lifecycle correctness: absent/non-function, synchronously throwing, rejecting, never-settling, and late-settling sinks are absorbed without altering the owning branch. The seam swallow-logs and bounds each write; it intentionally adds no retry, backoff, or queueing.
This applies to executor, run-auditor, self-healing, merger, PR reconciliation, scheduler, project-engine, plugin, mission-loop, hold-release, goal diagnostics, overseer advisor, mesh-lease, in-process runtime, credential rotation, and workflow-column-boundary emitters. `packages/engine/src/merge/merge-write-fence.ts` retains its bespoke non-`RunAuditEventInput` recorder; `packages/core` emitters, including `recordPlannerIntervention`, remain separately scoped. New engine emitters must ship with a behavioral sink-health regression covering hostile sink states, not only a source-routing assertion.
This applies to executor, run-auditor, self-healing, merger, PR reconciliation, scheduler, project-engine, plugin, mission-loop, hold-release, goal diagnostics, overseer advisor, mesh-lease, in-process runtime, credential rotation, and workflow-column-boundary emitters. `packages/engine/src/merge/merge-write-fence.ts` retains its bespoke non-`RunAuditEventInput` recorder. New engine emitters must ship with a behavioral sink-health regression covering hostile sink states, not only a source-routing assertion.
### 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.

View File

@@ -0,0 +1,21 @@
import { readFileSync } from "node:fs";
import { fileURLToPath } from "node:url";
import { describe, expect, it } from "vitest";
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",
];
describe("core run-audit emitter isolation", () => {
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");
expect(source, relative).toContain("emitBoundedRunAudit");
// Interface fields and host adapters may retain this identifier; only direct void/await calls are forbidden.
expect(source.match(/(?:await|void)\s+(?:\([^)]*\)\.)?recordRunAuditEvent\??\s*\(/g), relative).toBeNull();
}
});
});

View File

@@ -0,0 +1,317 @@
import { afterEach, describe, expect, it, vi } from "vitest";
const { asyncAuditSink } = vi.hoisted(() => ({ asyncAuditSink: { current: undefined as undefined | ((event: unknown) => unknown) } }));
vi.mock("../task-store/async/async-audit.js", () => ({
recordRunAuditEvent: (_layer: unknown, event: unknown) => asyncAuditSink.current?.(event),
}));
vi.mock("../postgres/data-layer.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("../postgres/data-layer.js")>();
return { ...actual, recordRunAuditEvent: (_layer: unknown, event: unknown) => asyncAuditSink.current?.(event) };
});
vi.mock("../task-store/async/async-persistence.js", () => ({
readTaskRow: vi.fn(async () => undefined),
}));
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 })),
clearTransitionPendingAsync: vi.fn(async () => undefined),
writeTransitionPendingAsync: vi.fn(async () => undefined),
}));
vi.mock("../task-store/async/async-workflow-workitems.js", async (importOriginal) => {
const actual = await importOriginal<typeof import("../task-store/async/async-workflow-workitems.js")>();
return {
...actual,
recordCompletionHandoff: vi.fn(async (_db: unknown, taskId: string, source: string, acceptedAt?: string) => ({
taskId, source, acceptedAt: acceptedAt ?? "2026-08-20T00:00:00.000Z",
})),
getCompletionHandoffMarker: vi.fn(async () => ({
taskId: "FN-9177", source: "test", acceptedAt: "2026-08-20T00:00:00.000Z",
})),
clearCompletionHandoffMarker: vi.fn(async () => undefined),
getWorkflowWorkItem: vi.fn(async () => ({ id: "WI-9177", taskId: "FN-9177", runId: "run-9177", state: "running", leaseOwner: "worker" })),
withTaskWorkflowSerialization: vi.fn(async (_tx: unknown, _projectId: string, _taskId: string, work: () => unknown) => work()),
};
});
import {
emitOverseerConfirmation,
emitOverseerEscalation,
emitOverseerObservation,
emitOverseerRecoveryAttempt,
emitOverseerRetry,
emitOverseerSteering,
} from "../planner/planner-overseer-events.js";
import { CORE_RUN_AUDIT_EMIT_TIMEOUT_MS } from "../run-audit/emit-bounded-run-audit.js";
import { runPluginColumnTransitionHooksImpl } from "../task-store/audit-ops.js";
import { rehomeOccupantImpl } from "../task-store/branch-group-ops.js";
import { applyPrMergedTransitionImpl } from "../task-store/merge-queue-ops-2.js";
import { insertRunAuditEventRowImpl, recordDependencyCycleRejectedAuditImpl } from "../task-store/task-id-integrity.js";
import { setCompletionHandoffAcceptedMarkerImpl, reconcileLegacyAutoMergeStampsImpl } from "../task-store/task-mutation-ops.js";
import { clearCompletionHandoffAcceptedMarkerImpl } from "../task-store/task-artifacts-ops.js";
import { recoverStaleTransitionPendingImpl } from "../task-store/lifecycle-ops.js";
import { projectMergeRequestToWorkflowWorkItemImpl } from "../task-store/workflow-workitems-ops.js";
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";
const input = { taskId: "FN-9177", stage: "executor" as const, reason: "test", timestamp: "2026-08-20T00:00:00.000Z" };
const façades = [
emitOverseerObservation,
emitOverseerSteering,
emitOverseerRecoveryAttempt,
emitOverseerRetry,
emitOverseerConfirmation,
emitOverseerEscalation,
] as const;
type SinkMode = "absent" | "throw" | "reject" | "never" | "late-resolve" | "late-reject";
function hostileSink(mode: SinkMode) {
let settle: (() => void) | undefined;
const sink = vi.fn(() => {
if (mode === "throw") throw new Error("audit throw");
if (mode === "reject") return Promise.reject(new Error("audit rejection"));
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 rejection")); });
return undefined;
});
return { store: mode === "absent" ? {} : { recordRunAuditEvent: sink }, sink, settle: () => settle?.() };
}
async function settleBounded(promise: Promise<void>, mode: SinkMode, settle: () => void): Promise<void> {
if (mode === "never" || mode.startsWith("late-")) {
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
settle();
}
await expect(promise).resolves.toBeUndefined();
}
afterEach(() => vi.useRealTimers());
describe("core audit emitters tolerate hostile sinks", () => {
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps every planner façade bounded for %s sinks", async (mode) => {
vi.useFakeTimers();
for (const emit of façades) {
const hostile = hostileSink(mode);
const promise = emit({ ...input, store: hostile.store as never });
// The planner's immediate timeline users observe this call before awaiting.
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
await settleBounded(promise, mode, hostile.settle);
}
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps synchronous store audit helpers synchronous for %s sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
asyncAuditSink.current = mode === "absent" ? undefined : hostile.sink;
const owner = { backendMode: true, asyncLayer: {} };
expect(() => insertRunAuditEventRowImpl(owner as never, {
taskId: "FN-9177", domain: "database", mutationType: "test:sync-helper", target: "FN-9177",
})).not.toThrow();
expect(() => recordDependencyCycleRejectedAuditImpl(owner as never, "FN-9177", ["FN-9177"], "updateTask")).not.toThrow();
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(2);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
asyncAuditSink.current = undefined;
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"does not alter workflow reconciliation for %s audit sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const task = { id: "FN-9177", column: "todo" };
const owner = {
...hostile.store,
backendMode: false,
readTaskFromDb: () => task,
moveTask: vi.fn(async () => task),
};
await expect(rehomeOccupantImpl(owner as never, task.id, "todo", "workflow-edit-rehome", { fixture: true })).resolves.toEqual({ moved: true });
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
// Reconciliation is fire-and-forget: its result is available before a hostile audit settles.
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"does not alter the merged-PR completion owner for %s audit sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const task = { id: "FN-9177", column: "in-review", prInfo: { status: "merged", number: 17 }, dependencies: [], steps: [] };
const owner = {
...hostile.store,
getTask: vi.fn(async () => task),
getTaskWorkflowSelection: () => undefined,
getTaskWorkflowSelectionAsync: async () => undefined,
moveTask: vi.fn(async () => ({ ...task, column: "done" })),
emit: vi.fn(),
};
await expect(applyPrMergedTransitionImpl(owner as never, task.id, { agentId: "merger", runId: "run-9177" })).resolves.toEqual({ moved: true });
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps completion handoff marker ownership non-blocking for %s audit sinks", async (mode) => {
vi.useFakeTimers();
// The marker helper is driven through its real exported owner. Persistence is stubbed only to
// isolate the post-commit optional audit dependency from the marker's durable write contract.
const hostile = hostileSink(mode);
const marker = { taskId: "FN-9177", acceptedAt: "2026-08-20T00:00:00.000Z", source: "test" };
const owner = { ...hostile.store, asyncLayer: { db: {} } };
await expect(setCompletionHandoffAcceptedMarkerImpl(owner as never, marker.taskId, { source: marker.source, acceptedAt: marker.acceptedAt })).resolves.toMatchObject({ source: marker.source });
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps completion handoff clearing non-blocking for %s audit sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const owner = { ...hostile.store, asyncLayer: { db: {} } };
await expect(clearCompletionHandoffAcceptedMarkerImpl(owner as never, "FN-9177")).resolves.toBeUndefined();
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps legacy stamp reconciliation results unchanged for %s audit sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const task = { id: "FN-9177", column: "in-review", autoMerge: true, autoMergeProvenance: "legacy-stamp" };
const owner = {
...hostile.store,
listLegacyAutoMergeStampCandidates: vi.fn(async () => [task]),
listWorkflowDefinitions: vi.fn(async () => []),
getTask: vi.fn(async () => ({ ...task })),
isLegacyAutoMergeStampCandidate: vi.fn(() => true),
atomicWriteTaskJson: vi.fn(async () => undefined),
taskDir: vi.fn(() => "/fixture/FN-9177"),
isWatching: false,
emitTaskLifecycleEventSafely: vi.fn(),
};
await expect(reconcileLegacyAutoMergeStampsImpl(owner as never, { apply: true })).resolves.toEqual([
{ taskId: task.id, column: task.column, cleared: true },
]);
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps work-item lease acquisition non-blocking for %s audit sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const tx = { update: vi.fn(() => ({ set: vi.fn(() => ({ where: vi.fn(async () => undefined) })) })) };
const owner = {
...hostile.store,
asyncLayer: { projectId: "project-9177", transactionImmediate: async (work: (tx: typeof tx) => unknown) => work(tx) },
};
await expect(acquireWorkflowWorkItemLeaseImpl(owner as never, "WI-9177", "worker", { leaseDurationMs: 1, now: "2026-08-20T00:00:00.000Z" })).resolves.toMatchObject({ id: "WI-9177" });
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps legacy stamp marking non-blocking for %s audit sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const task = { id: "FN-9177", column: "in-review", autoMerge: true };
const prepare = vi.fn(() => ({ get: vi.fn(() => undefined), run: vi.fn() }));
const owner = {
...hostile.store,
db: { prepare, bumpLastModified: vi.fn() },
listLegacyAutoMergeStampCandidates: vi.fn(async () => [task]),
listWorkflowDefinitions: vi.fn(async () => []),
getTask: vi.fn(async () => ({ ...task })),
isLegacyAutoMergeStampCandidate: vi.fn(() => true),
atomicWriteTaskJson: vi.fn(async () => undefined),
taskDir: vi.fn(() => "/fixture/FN-9177"),
isWatching: false,
emitTaskLifecycleEventSafely: vi.fn(),
};
await expect(markLegacyAutoMergeStampsOnceImpl(owner as never)).resolves.toBeUndefined();
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
});
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) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const registry = getTraitRegistry();
// The exported hook runner is exercised with a real plugin descriptor whose missing
// implementation takes the production "no-impl" degradation audit branch.
try { registry.registerTrait({ id: "plugin:fn-9177-audit", name: "audit", flags: {}, hooks: { onEnter: true } }); } catch { /* singleton may retain the fixture trait */ }
const hookOwner = { ...hostile.store, asyncLayer: { db: {} }, rowToTask: (value: unknown) => value };
await expect(runPluginColumnTransitionHooksImpl(hookOwner as never, "FN-9177", {
version: 1, nodes: [], edges: [], columns: [{ id: "todo", name: "Todo", traits: [{ trait: "plugin:fn-9177-audit" }] }],
} as never, "other", "todo")).resolves.toBeUndefined();
asyncAuditSink.current = mode === "absent" ? undefined : hostile.sink;
const projectionOwner = {
asyncLayer: {}, getMergeRequestRecordAsync: vi.fn(async () => ({ state: "pending", attemptCount: 1, updatedAt: "2026-08-20T00:00:00.000Z" })),
workflowStateForMergeRequestState: vi.fn(() => "runnable"),
upsertWorkflowWorkItem: vi.fn(async () => ({ id: "WI-9177", runId: "run-9177", state: "runnable", kind: "merge" })),
cancelActiveWorkflowWorkItemsForTask: vi.fn(async () => undefined),
};
await expect(projectMergeRequestToWorkflowWorkItemImpl(projectionOwner as never, "FN-9177")).resolves.toMatchObject({ id: "WI-9177" });
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalled();
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
asyncAuditSink.current = undefined;
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps transition-pending recovery results unchanged for %s sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const owner = {
...hostile.store, backendMode: true, asyncLayer: { db: {} },
withTaskLock: async (_id: string, work: () => unknown) => work(),
};
await expect(recoverStaleTransitionPendingImpl(owner as never)).resolves.toEqual({ scanned: 1, recovered: 1, degradedHooks: 0 });
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
"keeps the already-rehomed reconciliation owner non-blocking for %s sinks", async (mode) => {
vi.useFakeTimers();
const hostile = hostileSink(mode);
const task = { id: "FN-9177", column: "todo" };
const owner = { ...hostile.store, backendMode: false, readTaskFromDb: () => task, moveTask: vi.fn() };
await expect(rehomeOccupantImpl(owner as never, task.id, "todo", "workflow-edit-rehome", { fixture: true })).resolves.toEqual({ moved: true });
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
hostile.settle();
},
);
});

View File

@@ -0,0 +1,35 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import { CORE_RUN_AUDIT_EMIT_TIMEOUT_MS, emitBoundedRunAudit } 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: {} };
afterEach(() => vi.useRealTimers());
describe("emitBoundedRunAudit", () => {
it("calls sinks synchronously and resolves healthy, absent, throwing, and rejecting sinks", async () => {
const log = { warn: vi.fn() };
const sink = vi.fn(() => undefined);
const promise = emitBoundedRunAudit({ recordRunAuditEvent: sink }, event, { log });
expect(sink).toHaveBeenCalledWith(event);
await expect(promise).resolves.toBeUndefined();
await expect(emitBoundedRunAudit(undefined, event, { log })).resolves.toBeUndefined();
await expect(emitBoundedRunAudit({ recordRunAuditEvent: () => { throw new Error("no"); } }, event, { log })).resolves.toBeUndefined();
await expect(emitBoundedRunAudit({ recordRunAuditEvent: () => Promise.reject(new Error("no")) }, event, { log })).resolves.toBeUndefined();
expect(log.warn).toHaveBeenCalledWith("[run-audit] failed to record test:audit");
});
it("bounds never-settling and late-settling sinks without unhandled rejections", async () => {
vi.useFakeTimers();
const log = { warn: vi.fn() };
const never = emitBoundedRunAudit({ recordRunAuditEvent: () => new Promise(() => undefined) }, event, { log });
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
await expect(never).resolves.toBeUndefined();
expect(log.warn).toHaveBeenCalledWith("[run-audit] timed out recording test:audit");
let reject!: (error: Error) => void;
const late = emitBoundedRunAudit({ recordRunAuditEvent: () => new Promise<unknown>((_, fail) => { reject = fail; }) }, event, { timeoutMs: 1, log });
await vi.advanceTimersByTimeAsync(1);
reject(new Error("late"));
await expect(late).resolves.toBeUndefined();
});
});

View File

@@ -9,6 +9,7 @@ import type {
RunAuditEventInput,
} from "../types.js";
import { OVERSEER_INTERVENTION_MUTATION } from "../types.js";
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/**
* FNXC:PlannerOversight 2026-07-04-18:00:
@@ -23,7 +24,7 @@ import { OVERSEER_INTERVENTION_MUTATION } from "../types.js";
/** Minimal store seam this module depends on (satisfied by `TaskStore`). */
export interface PlannerInterventionStore {
recordRunAuditEvent(input: RunAuditEventInput): RunAuditEvent | Promise<RunAuditEvent>;
recordRunAuditEvent(input: RunAuditEventInput): unknown;
}
/** Read-capable seam used only by timeline queries. */
@@ -84,11 +85,11 @@ const KNOWN_SOURCE_LINK_KINDS: readonly PlannerInterventionSourceLink["kind"][]
"url",
];
/** Records one planner-intervention timeline entry as a run-audit event under `overseer:intervention`. Non-throwing on optional-field absence. */
/** Records one bounded, best-effort, non-rejecting planner-intervention audit event under `overseer:intervention`. */
export function recordPlannerIntervention(
store: PlannerInterventionStore,
input: RecordPlannerInterventionInput,
): RunAuditEvent | Promise<RunAuditEvent> {
): Promise<void> {
const metadata: Record<string, unknown> = {
stage: input.stage,
reason: input.reason,
@@ -109,7 +110,7 @@ export function recordPlannerIntervention(
metadata.advisorSlug = input.advisorSlug.trim();
}
return store.recordRunAuditEvent({
return emitBoundedRunAudit(store, {
timestamp: input.timestamp,
taskId: input.taskId,
agentId: input.agentId ?? "overseer",

View File

@@ -3,7 +3,6 @@ import type {
PlannerInterventionOutcome,
PlannerInterventionSourceLink,
PlannerOversightStage,
RunAuditEvent,
} from "../types.js";
import { type PlannerInterventionStore, recordPlannerIntervention } from "./planner-intervention.js";
@@ -64,7 +63,7 @@ function normalizeAndRecord(
input: OverseerEventInput,
action: PlannerInterventionAction,
defaultOutcome: PlannerInterventionOutcome,
): RunAuditEvent | Promise<RunAuditEvent> {
): Promise<void> {
return recordPlannerIntervention(input.store, {
taskId: input.taskId,
runId: input.runId,
@@ -88,7 +87,7 @@ function normalizeAndRecord(
* action taken). Default outcome: `"succeeded"` (the observation itself always
* "succeeds"; attempt fields are typically omitted for this category).
*/
export function emitOverseerObservation(input: OverseerEventInput): RunAuditEvent | Promise<RunAuditEvent> {
export function emitOverseerObservation(input: OverseerEventInput): Promise<void> {
return normalizeAndRecord(input, "observe", "succeeded");
}
@@ -97,7 +96,7 @@ export function emitOverseerObservation(input: OverseerEventInput): RunAuditEven
* Default outcome: `"pending"` (guidance has been injected; whether it lands
* successfully is determined by a later observation/retry).
*/
export function emitOverseerSteering(input: OverseerEventInput): RunAuditEvent | Promise<RunAuditEvent> {
export function emitOverseerSteering(input: OverseerEventInput): Promise<void> {
return normalizeAndRecord(input, "inject-guidance", "pending");
}
@@ -106,7 +105,7 @@ export function emitOverseerSteering(input: OverseerEventInput): RunAuditEvent |
* Default outcome: `"pending"`. Callers should supply `attemptCount`/`attemptLimit`
* so the timeline can render bounded-recovery progress.
*/
export function emitOverseerRecoveryAttempt(input: OverseerEventInput): RunAuditEvent | Promise<RunAuditEvent> {
export function emitOverseerRecoveryAttempt(input: OverseerEventInput): Promise<void> {
return normalizeAndRecord(input, "request-fix", "pending");
}
@@ -115,7 +114,7 @@ export function emitOverseerRecoveryAttempt(input: OverseerEventInput): RunAudit
* Callers should supply `attemptCount`/`attemptLimit` so the timeline can
* render bounded-retry progress.
*/
export function emitOverseerRetry(input: OverseerEventInput): RunAuditEvent | Promise<RunAuditEvent> {
export function emitOverseerRetry(input: OverseerEventInput): Promise<void> {
return normalizeAndRecord(input, "retry", "pending");
}
@@ -123,7 +122,7 @@ export function emitOverseerRetry(input: OverseerEventInput): RunAuditEvent | Pr
* Records a merge/PR confirmation request raised to a human. Default outcome:
* `"awaiting-confirmation"`.
*/
export function emitOverseerConfirmation(input: OverseerEventInput): RunAuditEvent | Promise<RunAuditEvent> {
export function emitOverseerConfirmation(input: OverseerEventInput): Promise<void> {
return normalizeAndRecord(input, "request-confirmation", "awaiting-confirmation");
}
@@ -133,6 +132,6 @@ export function emitOverseerConfirmation(input: OverseerEventInput): RunAuditEve
* example a caller may escalate with outcome `"skipped"` when escalation is
* itself bypassed by a human-control guard.
*/
export function emitOverseerEscalation(input: OverseerEventInput): RunAuditEvent | Promise<RunAuditEvent> {
export function emitOverseerEscalation(input: OverseerEventInput): Promise<void> {
return normalizeAndRecord(input, "escalate", "failed");
}

View File

@@ -0,0 +1,51 @@
import type { RunAuditEventInput } from "../types.js";
import { createLogger } from "../process/logger.js";
export const CORE_RUN_AUDIT_EMIT_TIMEOUT_MS = 2_000;
export type RunAuditSinkHost = {
recordRunAuditEvent?: (input: RunAuditEventInput) => unknown;
} | null | undefined;
export type RunAuditLogger = { warn: (message: string) => void };
type RunAuditEvent = RunAuditEventInput | { mutationType: string; [key: string]: unknown };
const defaultLog = createLogger("run-audit");
/**
* FNXC:RunAudit 2026-08-20-05:49:
* FN-9177 keeps this bounded optional-audit seam in core deliberately: core cannot import engine
* 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(
host: RunAuditSinkHost,
event: RunAuditEvent,
options: { timeoutMs?: number; log?: RunAuditLogger } = {},
): Promise<void> {
const log = options.log ?? defaultLog;
const sink = host?.recordRunAuditEvent;
if (typeof sink !== "function") return;
let sinkPromise: Promise<unknown>;
try {
sinkPromise = Promise.resolve(sink.call(host, event as RunAuditEventInput));
} catch {
log.warn(`[run-audit] failed to record ${event.mutationType}`);
return;
}
void sinkPromise.catch(() => undefined);
await new Promise<void>((resolve) => {
const timer = setTimeout(() => {
log.warn(`[run-audit] timed out recording ${event.mutationType}`);
resolve();
}, 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(); },
);
});
}

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
/**
* audit-ops operations.
*
@@ -79,7 +81,7 @@ export async function runPluginColumnTransitionHooksImpl(store: TaskStore, taskI
const resolved = registry.resolveTraitHook(traitId, hookKind);
if (resolved.warning) {
// Degraded (no impl / force-disabled) → passive no-op, audit the warning.
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId,
agentId: "system",
runId: `plugin-trait-hook-${traitId}-${taskId}-${Date.now()}`,
@@ -93,7 +95,7 @@ export async function runPluginColumnTransitionHooksImpl(store: TaskStore, taskI
await resolved.impl({ task: taskDetail, context: { fromColumn, toColumn, hookKind } });
} catch (err) {
// A throwing plugin hook DEGRADES — audited, never wedges the lock.
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId,
agentId: "system",
runId: `plugin-trait-hook-${traitId}-${taskId}-${Date.now()}`,

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
/**
* branch-group-ops operations.
*
@@ -429,7 +431,7 @@ export async function rehomeOccupantImpl(store: TaskStore, taskId: string, targe
if (fromColumn === targetColumn) {
// Already in the target column — nothing to move, but still record the
// reconciliation decision for audit traceability.
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId,
agentId: "system",
runId: `workflow-reconcile-${reason}-${taskId}-${Date.now()}`,
@@ -461,7 +463,7 @@ export async function rehomeOccupantImpl(store: TaskStore, taskId: string, targe
} catch (err) {
error = err instanceof Error ? err.message : String(err);
}
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId,
agentId: "system",
runId: `workflow-reconcile-${reason}-${taskId}-${Date.now()}`,

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
/**
* lifecycle-ops operations.
*
@@ -997,7 +999,7 @@ export async function recoverStaleTransitionPendingImpl(store: TaskStore): Promi
// best-effort; a later sweep retries.
}
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId: id,
agentId: "system",
runId: `transition-pending-recovery-${id}-${Date.now()}`,

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
/**
* merge-queue-ops-2 operations.
*
@@ -205,7 +207,7 @@ export async function applyPrMergedTransitionImpl(store: TaskStore, taskId: stri
} satisfies MergeResult);
if (ctx?.agentId && ctx?.runId) {
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId,
agentId: ctx.agentId,
runId: ctx.runId,

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
/**
* FNXC:CodeOrganization 2026-07-20-14:00:
* Domain rename from remaining-ops-7: merge-queue peeks, archive/done transitions,
@@ -62,7 +64,7 @@ export async function clearCompletionHandoffAcceptedMarkerImpl(store: TaskStore,
const existing = await getCompletionHandoffMarkerAsync(layer.db, taskId);
if (!existing) return;
await clearCompletionHandoffMarkerAsync(layer.db, taskId);
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId,
agentId: "system",
runId: `completion-handoff-clear:${taskId}:${Date.now()}`,

View File

@@ -10,6 +10,8 @@
*/
import { TaskStore } from "../store.js";
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so synchronous store helpers remain non-blocking. */
import { randomUUID } from "node:crypto";
import { and, eq } from "drizzle-orm";
import { ArchiveDatabase } from "../db/archive-db.js";
@@ -577,18 +579,20 @@ export function insertRunAuditEventRowImpl(store: TaskStore, input: Omit<RunAudi
const eventId = randomUUID();
const agentId = input.agentId ?? "store";
const runId = input.runId ?? `store:${input.mutationType}:${input.taskId ?? input.target}:${eventId}`;
void recordRunAuditEventAsync(store.asyncLayer, {
timestamp: input.timestamp,
taskId: input.taskId,
agentId,
runId,
domain: input.domain,
mutationType: input.mutationType,
target: input.target,
metadata: input.metadata as Record<string, unknown> | undefined,
}).catch((err) => {
storeLog.warn(`[run-audit-event-failed] ${input.mutationType}:${input.taskId ?? input.target}`, { error: getErrorMessage(err) });
});
void emitBoundedRunAudit(
{ recordRunAuditEvent: (event) => recordRunAuditEventAsync(store.asyncLayer!, event) },
{
timestamp: input.timestamp,
taskId: input.taskId,
agentId,
runId,
domain: input.domain,
mutationType: input.mutationType,
target: input.target,
metadata: input.metadata as Record<string, unknown> | undefined,
},
{ log: { warn: (detail) => storeLog.warn(`[run-audit-event-failed] ${input.mutationType}:${input.taskId ?? input.target}`, { error: detail }) } },
);
return;
}
const eventId = randomUUID();
@@ -836,17 +840,19 @@ export function recordDependencyCycleRejectedAuditImpl(store: TaskStore,
*/
if (store.backendMode && store.asyncLayer) {
const mutationType = source === "replication" ? "task:dependency-cycle-rejected-replication" : "task:dependency-cycle-rejected";
void recordRunAuditEventAsync(store.asyncLayer, {
taskId,
agentId: "store",
runId: `store:${mutationType}:${taskId}`,
domain: "database",
mutationType,
target: taskId,
metadata: { taskId, cyclePath, source } as Record<string, unknown>,
}).catch((err) => {
storeLog.warn(`[dependency-cycle-rejected-audit-failed] ${taskId}`, { error: getErrorMessage(err) });
});
void emitBoundedRunAudit(
{ recordRunAuditEvent: (event) => recordRunAuditEventAsync(store.asyncLayer!, event) },
{
taskId,
agentId: "store",
runId: `store:${mutationType}:${taskId}`,
domain: "database",
mutationType,
target: taskId,
metadata: { taskId, cyclePath, source } as Record<string, unknown>,
},
{ log: { warn: (detail) => storeLog.warn(`[dependency-cycle-rejected-audit-failed] ${taskId}`, { error: detail }) } },
);
return;
}
store.insertRunAuditEventRow({

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
import { createLogger } from "../process/logger.js";
import { resolveLegacyStampReviewColumns } from "./task-store-helpers.js";
@@ -843,7 +845,7 @@ export async function setCompletionHandoffAcceptedMarkerImpl(store: TaskStore, t
await recordCompletionHandoffAsync(layer.db, taskId, opts.source, opts.acceptedAt);
const marker = await getCompletionHandoffMarkerAsync(layer.db, taskId);
if (!marker) throw new Error(`Failed to set completion handoff marker for ${taskId}`);
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId,
agentId: "system",
runId: `completion-handoff:${taskId}:${Date.now()}`,
@@ -851,12 +853,7 @@ export async function setCompletionHandoffAcceptedMarkerImpl(store: TaskStore, t
mutationType: "task:completion-handoff-accepted",
target: taskId,
metadata: { taskId, acceptedAt: marker.acceptedAt, source: marker.source },
}).catch((err) => {
storeLog.warn("completion-handoff audit write failed", {
taskId,
error: err instanceof Error ? err.message : String(err),
});
});
}, { log: { warn: (detail) => storeLog.warn("completion-handoff audit write failed", { taskId, detail }) } });
return marker as CompletionHandoffMarker;
}
@@ -892,7 +889,7 @@ export async function reconcileLegacyAutoMergeStampsImpl(store: TaskStore, optio
if (store.isWatching) store.taskCache.set(current.id, { ...current });
store.emitTaskLifecycleEventSafely("task:updated", [current]);
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId: current.id,
agentId: "system",
runId: `legacy-auto-merge-stamp-clear-${current.id}-${Date.now()}`,

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
import { createLogger } from "../process/logger.js";
import { resolveLegacyStampReviewColumns } from "./task-store-helpers.js";
@@ -53,7 +55,7 @@ export async function markLegacyAutoMergeStampsOnceImpl(store: TaskStore): Promi
store.emitTaskLifecycleEventSafely("task:updated", [current]);
markedTaskIds.push(current.id);
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId: current.id,
agentId: "system",
runId: `legacy-auto-merge-stamp-mark-${current.id}-${Date.now()}`,

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
/**
* workflow-workitems-ops-2 operations.
*
@@ -101,7 +103,7 @@ export async function acquireWorkflowWorkItemLeaseImpl(store: TaskStore, id: str
});
if (!updated) return null;
// Record the audit event (fire-and-forget).
void store.recordRunAuditEvent({
void emitBoundedRunAudit(store, {
taskId: updated.taskId,
agentId: "system",
runId: updated.runId,

View File

@@ -1,3 +1,5 @@
import { emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
/* FNXC:RunAudit 2026-08-20-05:49: FN-9177 bounds optional audit telemetry so a hostile sink cannot alter this lifecycle path. */
/**
* workflow-workitems-ops operations.
*
@@ -53,7 +55,7 @@ export async function projectMergeRequestToWorkflowWorkItemImpl(store: TaskStore
now: opts.now ?? record.updatedAt,
lastError: "superseded-by-merge-request-projection",
});
void recordRunAuditEventAsync(layer, {
void emitBoundedRunAudit({ recordRunAuditEvent: (input) => recordRunAuditEventAsync(layer, input) }, {
taskId,
agentId: "system",
runId: item.runId,