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:
7
.changeset/fn-9177-core-run-audit.md
Normal file
7
.changeset/fn-9177-core-run-audit.md
Normal 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.
|
||||
@@ -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.
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
});
|
||||
});
|
||||
317
packages/core/src/__tests__/core-run-audit-sink-health.test.ts
Normal file
317
packages/core/src/__tests__/core-run-audit-sink-health.test.ts
Normal 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();
|
||||
},
|
||||
);
|
||||
});
|
||||
35
packages/core/src/__tests__/emit-bounded-run-audit.test.ts
Normal file
35
packages/core/src/__tests__/emit-bounded-run-audit.test.ts
Normal 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();
|
||||
});
|
||||
});
|
||||
@@ -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",
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
|
||||
51
packages/core/src/run-audit/emit-bounded-run-audit.ts
Normal file
51
packages/core/src/run-audit/emit-bounded-run-audit.ts
Normal 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(); },
|
||||
);
|
||||
});
|
||||
}
|
||||
@@ -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()}`,
|
||||
|
||||
@@ -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()}`,
|
||||
|
||||
@@ -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()}`,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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()}`,
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -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()}`,
|
||||
|
||||
@@ -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()}`,
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
|
||||
Reference in New Issue
Block a user