FN-9182: Bound reconciliation audit emissions
Keep reconciliation responsive while preserving audit-dependent outcomes. - add a bounded core audit emitter that reports recorded, absent, failed, and timed-out outcomes - route workflow-switch and phantom-reservation reconciliation through the outcome-aware seam - cover hostile sink behavior, caller-visible failures, routing invariants, and operator documentation Files changed: .changeset/fn-9182-bounded-run-audit-outcome.md | 7 ++ AGENTS.md | 1 + docs/run-audit.md | 4 +- .../core-run-audit-emitter-isolation.test.ts | 10 +-- .../__tests__/core-run-audit-sink-health.test.ts | 51 ++++++++++++++- .../src/__tests__/emit-bounded-run-audit.test.ts | 58 ++++++++++++++++- .../phantom-reservation-audit-outcome.test.ts | 76 ++++++++++++++++++++++ .../core/src/run-audit/emit-bounded-run-audit.ts | 41 +++++++++--- .../task-store/async/async-phantom-reservations.ts | 31 +++++---- .../core/src/task-store/workflow-definitions.ts | 48 ++++++++------ 10 files changed, 277 insertions(+), 50 deletions(-) Fusion-Task-Id: FN-9182 Fusion-Task-Lineage: c98e8e5e-2dae-47a0-b969-4156f4beccb3 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-9182-bounded-run-audit-outcome.md
Normal file
7
.changeset/fn-9182-bounded-run-audit-outcome.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": patch
|
||||||
|
---
|
||||||
|
|
||||||
|
summary: Keep workflow recovery and reservation cleanup responsive when audit logging stalls.
|
||||||
|
category: fix
|
||||||
|
dev: Adds emitBoundedRunAuditWithOutcome for workflow-switch-torn and phantom-reservation reconciliation.
|
||||||
@@ -293,6 +293,7 @@ canonical emitters remain explicit exclusions until their separately scoped hard
|
|||||||
-->
|
-->
|
||||||
- FN-9175: New engine run-audit emitters must use `emitBoundedRunAudit` from `packages/engine/src/util/emit-bounded-run-audit.ts`; it absorbs absent, throwing, rejecting, hanging, and late-settling sinks without changing the owning branch, and requires behavioral sink-health coverage.
|
- FN-9175: New engine run-audit emitters must use `emitBoundedRunAudit` from `packages/engine/src/util/emit-bounded-run-audit.ts`; it absorbs absent, throwing, rejecting, hanging, and late-settling sinks without changing the owning branch, and requires behavioral sink-health coverage.
|
||||||
- FNXC:RunAudit 2026-08-20-05:49: FN-9177 requires new core best-effort emitters to use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. It deliberately mirrors the engine seam because core cannot import engine; transactional and deliberately awaited durability writers remain unbounded.
|
- FNXC:RunAudit 2026-08-20-05:49: FN-9177 requires new core best-effort emitters to use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. It deliberately mirrors the engine seam because core cannot import engine; transactional and deliberately awaited durability writers remain unbounded.
|
||||||
|
- FNXC:RunAudit 2026-08-20-07:14: FN-9182 requires core emitters whose behavior branches on audit success to use `emitBoundedRunAuditWithOutcome`; `emitBoundedRunAudit` remains the default best-effort seam, while transactional and deliberately awaited durability writers remain unbounded.
|
||||||
- FN-9109: `session:cross-runtime-fallback-engaged` records a single retryable-failure handoff from a primary runtime to a deferred CLI runtime. Metadata is ids/outcomes-only (`sessionPurpose`, primary/fallback provider and model IDs, trigger point, failure category, `contextTransferred`); never record error prose or transferred transcript text.
|
- FN-9109: `session:cross-runtime-fallback-engaged` records a single retryable-failure handoff from a primary runtime to a deferred CLI runtime. Metadata is ids/outcomes-only (`sessionPurpose`, primary/fallback provider and model IDs, trigger point, failure category, `contextTransferred`); never record error prose or transferred transcript text.
|
||||||
|
|
||||||
- FN-8958: `merge:orphan-write-fenced` is emitted once per orphan merge body at its fence's first interaction. Metadata is ids/counts/outcomes-only: `{ taskId, category, interaction, suppressedCount }`; `suppressedCount` is the emit-time count (`1` for `interaction:"suppressed"`, `0` for `interaction:"rejected"`), never a cumulative body total.
|
- FN-8958: `merge:orphan-write-fenced` is emitted once per orphan merge body at its fence's first interaction. Metadata is ids/counts/outcomes-only: `{ taskId, category, interaction, suppressedCount }`; `suppressedCount` is the emit-time count (`1` for `interaction:"suppressed"`, `0` for `interaction:"rejected"`), never a cumulative body total.
|
||||||
|
|||||||
@@ -82,10 +82,10 @@ This applies to executor, run-auditor, self-healing, merger, PR reconciliation,
|
|||||||
|
|
||||||
### Core emit-seam policy
|
### Core emit-seam policy
|
||||||
|
|
||||||
Core best-effort emitters use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. This is a deliberate copy of the engine seam because `@fusion/core` cannot import `@fusion/engine`; it synchronously invokes a valid sink, then absorbs throws, rejection, timeout, and late settlement without making telemetry lifecycle-load-bearing. Transactional writers and explicitly awaited durability/ordering writers remain unbounded. `packages/core/src/__tests__/core-run-audit-sink-health.test.ts` and `core-run-audit-emitter-isolation.test.ts` respectively enforce hostile-sink behavior and source routing.
|
Core best-effort emitters use `packages/core/src/run-audit/emit-bounded-run-audit.ts`. This is a deliberate copy of the engine seam because `@fusion/core` cannot import `@fusion/engine`; it synchronously invokes a valid sink, then absorbs throws, rejection, timeout, and late settlement without making telemetry lifecycle-load-bearing. `emitBoundedRunAudit` is the default void seam. `emitBoundedRunAuditWithOutcome` returns `recorded`, `absent`, `failed` (with the original error), or `timed-out` where a forensic throw ordering or caller-visible skipped payload depends on the audit result; workflow-switch torn reconciliation and phantom committed-reservation reconciliation use it. Transactional writers and explicitly awaited durability/ordering writers remain unbounded. `packages/core/src/__tests__/core-run-audit-sink-health.test.ts` and `core-run-audit-emitter-isolation.test.ts` respectively enforce hostile-sink behavior and source routing.
|
||||||
|
|
||||||
### Awaited core exclusion decision
|
### Awaited core exclusion decision
|
||||||
|
|
||||||
FN-9178 classified the remaining awaited sites with hostile-sink characterization tests. `task-deleted-outbox:catch-up`, `:reconciliation-fallback`, `:lease-fenced`, `:retention-pruned`, and recall-capture audit events are class A (bound-safe) candidates. `task:workflow-switch-torn` and `task:reconcile-phantom-committed-reservation` are class B: a future bounded seam must return `recorded`, `failed`, `timed-out`, or `absent`, because their throw/result payload depends on audit outcome. `task:bypass-review`, `task:resume-step`, and both resurrection-blocked records are class C and intentionally unbounded: they claim persistence before return, destructive cleanup, or a forensic throw.
|
FN-9178 classified the remaining awaited sites with hostile-sink characterization tests. `task-deleted-outbox:catch-up`, `:reconciliation-fallback`, `:lease-fenced`, `:retention-pruned`, and recall-capture audit events are class A (bound-safe) candidates. `task:workflow-switch-torn` and `task:reconcile-phantom-committed-reservation` are class B and use the bounded outcome seam because their throw/result payload depends on audit outcome. `task:bypass-review`, `task:resume-step`, and both resurrection-blocked records are class C and intentionally unbounded: they claim persistence before return, destructive cleanup, or a forensic throw.
|
||||||
|
|
||||||
All `recordRunAuditEventWithinTransaction(tx, ...)` calls and the `recordRunAuditEventBackend(tx, ...)` transactional call are permanently out of scope. Their audit row shares a transaction with the mutation it describes; bounding would split that atomicity. The full matrix and evidence pointers are in the FN-9178 `decision` task document; `excluded-awaited-run-audit-store-sites.test.ts`, `excluded-awaited-run-audit-layer-sites.test.ts`, and the core routing ratchet pin this boundary.
|
All `recordRunAuditEventWithinTransaction(tx, ...)` calls and the `recordRunAuditEventBackend(tx, ...)` transactional call are permanently out of scope. Their audit row shares a transaction with the mutation it describes; bounding would split that atomicity. The full matrix and evidence pointers are in the FN-9178 `decision` task document; `excluded-awaited-run-audit-store-sites.test.ts`, `excluded-awaited-run-audit-layer-sites.test.ts`, and the core routing ratchet pin this boundary.
|
||||||
|
|||||||
@@ -5,14 +5,13 @@ import { describe, expect, it } from "vitest";
|
|||||||
/*
|
/*
|
||||||
* FNXC:RunAudit 2026-08-20-07:16:
|
* FNXC:RunAudit 2026-08-20-07:16:
|
||||||
* FN-9178 makes every direct awaited core audit writer a named decision rather than an implicit
|
* FN-9178 makes every direct awaited core audit writer a named decision rather than an implicit
|
||||||
* exception. Class A sites are evaluated candidates, B needs a reporting bounded seam, C retains
|
* exception. Class A sites are evaluated candidates, C retains its ordering claim, and
|
||||||
* its ordering claim, and transactional/sink writers remain permanent atomicity boundaries.
|
* transactional/sink writers remain permanent atomicity boundaries. FN-9182 migrated class B
|
||||||
|
* sites to the reporting bounded seam, so they no longer appear in this direct-await inventory.
|
||||||
*/
|
*/
|
||||||
const awaitedClassifications = {
|
const awaitedClassifications = {
|
||||||
"store.ts:task:bypass-review": "C",
|
"store.ts:task:bypass-review": "C",
|
||||||
"store.ts:task:resume-step": "C",
|
"store.ts:task:resume-step": "C",
|
||||||
"task-store/workflow-definitions.ts:task:workflow-switch-torn": "B",
|
|
||||||
"task-store/async/async-phantom-reservations.ts:task:reconcile-phantom-committed-reservation": "B",
|
|
||||||
"task-store/task-creation.ts:intake:resurrection-blocked": "C",
|
"task-store/task-creation.ts:intake:resurrection-blocked": "C",
|
||||||
"task-store/task-id-integrity.ts:task:resurrection-blocked": "C",
|
"task-store/task-id-integrity.ts:task:resurrection-blocked": "C",
|
||||||
"task-store/task-deleted-outbox-consumer.ts:task-deleted-outbox:catch-up": "A",
|
"task-store/task-deleted-outbox-consumer.ts:task-deleted-outbox:catch-up": "A",
|
||||||
@@ -35,7 +34,8 @@ const files = [
|
|||||||
"../planner/planner-intervention.ts", "../task-store/audit-ops.ts", "../task-store/branch-group-ops.ts",
|
"../planner/planner-intervention.ts", "../task-store/audit-ops.ts", "../task-store/branch-group-ops.ts",
|
||||||
"../task-store/merge-queue-ops-2.ts", "../task-store/task-mutation-ops.ts", "../task-store/workflow-integrity.ts",
|
"../task-store/merge-queue-ops-2.ts", "../task-store/task-mutation-ops.ts", "../task-store/workflow-integrity.ts",
|
||||||
"../task-store/workflow-workitems-ops.ts", "../task-store/workflow-workitems-ops-2.ts", "../task-store/task-artifacts-ops.ts",
|
"../task-store/workflow-workitems-ops.ts", "../task-store/workflow-workitems-ops-2.ts", "../task-store/task-artifacts-ops.ts",
|
||||||
"../task-store/lifecycle-ops.ts", "../task-store/task-id-integrity.ts",
|
"../task-store/lifecycle-ops.ts", "../task-store/task-id-integrity.ts", "../task-store/workflow-definitions.ts",
|
||||||
|
"../task-store/async/async-phantom-reservations.ts",
|
||||||
];
|
];
|
||||||
|
|
||||||
const sourceRoot = fileURLToPath(new URL("..", import.meta.url));
|
const sourceRoot = fileURLToPath(new URL("..", import.meta.url));
|
||||||
|
|||||||
@@ -1,6 +1,10 @@
|
|||||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||||
|
|
||||||
const { asyncAuditSink } = vi.hoisted(() => ({ asyncAuditSink: { current: undefined as undefined | ((event: unknown) => unknown) } }));
|
const { asyncAuditSink, workflowReader, resolvedWorkflowIr } = vi.hoisted(() => ({
|
||||||
|
asyncAuditSink: { current: undefined as undefined | ((event: unknown) => unknown) },
|
||||||
|
workflowReader: { current: undefined as undefined | (() => unknown) },
|
||||||
|
resolvedWorkflowIr: { current: undefined as unknown },
|
||||||
|
}));
|
||||||
vi.mock("../task-store/async/async-audit.js", () => ({
|
vi.mock("../task-store/async/async-audit.js", () => ({
|
||||||
recordRunAuditEvent: (_layer: unknown, event: unknown) => asyncAuditSink.current?.(event),
|
recordRunAuditEvent: (_layer: unknown, event: unknown) => asyncAuditSink.current?.(event),
|
||||||
}));
|
}));
|
||||||
@@ -11,9 +15,14 @@ vi.mock("../postgres/data-layer.js", async (importOriginal) => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
vi.mock("../task-store/async/async-persistence.js", () => ({
|
vi.mock("../task-store/async/async-persistence.js", () => ({
|
||||||
readTaskRow: vi.fn(async () => undefined),
|
readTaskRow: vi.fn(async () => workflowReader.current?.()),
|
||||||
}));
|
}));
|
||||||
|
|
||||||
|
vi.mock("../workflows/workflow-ir-resolver.js", async (importOriginal) => {
|
||||||
|
const actual = await importOriginal<typeof import("../workflows/workflow-ir-resolver.js")>();
|
||||||
|
return { ...actual, resolveWorkflowIrForTask: vi.fn(async () => resolvedWorkflowIr.current) };
|
||||||
|
});
|
||||||
|
|
||||||
vi.mock("../task-store/async/async-transition-pending.js", () => ({
|
vi.mock("../task-store/async/async-transition-pending.js", () => ({
|
||||||
listTransitionPendingTaskIdsAsync: vi.fn(async () => ["FN-9177"]),
|
listTransitionPendingTaskIdsAsync: vi.fn(async () => ["FN-9177"]),
|
||||||
readTransitionPendingAsync: vi.fn(async () => ({ toColumn: "todo", hooksRemaining: ["default-workflow:postCommit"], startedAt: 1 })),
|
readTransitionPendingAsync: vi.fn(async () => ({ toColumn: "todo", hooksRemaining: ["default-workflow:postCommit"], startedAt: 1 })),
|
||||||
@@ -57,6 +66,8 @@ import { projectMergeRequestToWorkflowWorkItemImpl } from "../task-store/workflo
|
|||||||
import { acquireWorkflowWorkItemLeaseImpl } from "../task-store/workflow-workitems-ops-2.js";
|
import { acquireWorkflowWorkItemLeaseImpl } from "../task-store/workflow-workitems-ops-2.js";
|
||||||
import { markLegacyAutoMergeStampsOnceImpl } from "../task-store/workflow-integrity.js";
|
import { markLegacyAutoMergeStampsOnceImpl } from "../task-store/workflow-integrity.js";
|
||||||
import { getTraitRegistry } from "../workflows/trait-registry.js";
|
import { getTraitRegistry } from "../workflows/trait-registry.js";
|
||||||
|
import { WorkflowSwitchRehomeFailedError } from "../workflows/workflow-reconciliation.js";
|
||||||
|
import { selectTaskWorkflowAndReconcileImpl } from "../task-store/workflow-definitions.js";
|
||||||
|
|
||||||
const input = { taskId: "FN-9177", stage: "executor" as const, reason: "test", timestamp: "2026-08-20T00:00:00.000Z" };
|
const input = { taskId: "FN-9177", stage: "executor" as const, reason: "test", timestamp: "2026-08-20T00:00:00.000Z" };
|
||||||
const façades = [
|
const façades = [
|
||||||
@@ -258,6 +269,42 @@ describe("core audit emitters tolerate hostile sinks", () => {
|
|||||||
|
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("workflow-switch torn reconciliation retains its forensic rejection", () => {
|
||||||
|
afterEach(() => {
|
||||||
|
workflowReader.current = undefined;
|
||||||
|
resolvedWorkflowIr.current = undefined;
|
||||||
|
});
|
||||||
|
|
||||||
|
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
|
||||||
|
"throws the committed torn error for %s audit sinks", async (mode) => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const hostile = hostileSink(mode);
|
||||||
|
let readCount = 0;
|
||||||
|
workflowReader.current = () => (readCount++ === 0 ? undefined : { id: "FN-9182", column: "todo" });
|
||||||
|
resolvedWorkflowIr.current = { version: 1, nodes: [], edges: [], columns: [{ id: "inbox", name: "Inbox", traits: [] }] };
|
||||||
|
const owner = {
|
||||||
|
...hostile.store,
|
||||||
|
asyncLayer: {},
|
||||||
|
selectTaskWorkflow: vi.fn(async () => []),
|
||||||
|
rehomeOccupant: vi.fn(async () => ({ moved: false })),
|
||||||
|
};
|
||||||
|
const operation = selectTaskWorkflowAndReconcileImpl(owner as never, "FN-9182", "WF-9182");
|
||||||
|
const rejection = expect(operation).rejects.toMatchObject({
|
||||||
|
name: "WorkflowSwitchRehomeFailedError",
|
||||||
|
taskId: "FN-9182",
|
||||||
|
workflowId: "WF-9182",
|
||||||
|
fromColumn: "todo",
|
||||||
|
intendedColumn: "inbox",
|
||||||
|
committed: true,
|
||||||
|
} satisfies Partial<WorkflowSwitchRehomeFailedError>);
|
||||||
|
if (mode === "never" || mode.startsWith("late-")) await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
|
||||||
|
hostile.settle();
|
||||||
|
await rejection;
|
||||||
|
if (mode !== "absent") expect(hostile.sink).toHaveBeenCalledTimes(1);
|
||||||
|
},
|
||||||
|
);
|
||||||
|
});
|
||||||
|
|
||||||
describe("remaining production owners retain hostile-sink isolation", () => {
|
describe("remaining production owners retain hostile-sink isolation", () => {
|
||||||
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
|
it.each(["absent", "throw", "reject", "never", "late-resolve", "late-reject"] as const)(
|
||||||
"keeps plugin trait degradation and merge-request projection non-blocking for %s sinks", async (mode) => {
|
"keeps plugin trait degradation and merge-request projection non-blocking for %s sinks", async (mode) => {
|
||||||
|
|||||||
@@ -1,5 +1,9 @@
|
|||||||
import { afterEach, describe, expect, it, vi } from "vitest";
|
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||||
import { CORE_RUN_AUDIT_EMIT_TIMEOUT_MS, emitBoundedRunAudit } from "../run-audit/emit-bounded-run-audit.js";
|
import {
|
||||||
|
CORE_RUN_AUDIT_EMIT_TIMEOUT_MS,
|
||||||
|
emitBoundedRunAudit,
|
||||||
|
emitBoundedRunAuditWithOutcome,
|
||||||
|
} from "../run-audit/emit-bounded-run-audit.js";
|
||||||
|
|
||||||
const event = { taskId: "FN-1", agentId: "system", runId: "run", domain: "database" as const, mutationType: "test:audit", target: "FN-1", metadata: {} };
|
const event = { taskId: "FN-1", agentId: "system", runId: "run", domain: "database" as const, mutationType: "test:audit", target: "FN-1", metadata: {} };
|
||||||
|
|
||||||
@@ -33,3 +37,55 @@ describe("emitBoundedRunAudit", () => {
|
|||||||
await expect(late).resolves.toBeUndefined();
|
await expect(late).resolves.toBeUndefined();
|
||||||
});
|
});
|
||||||
});
|
});
|
||||||
|
|
||||||
|
describe("emitBoundedRunAuditWithOutcome", () => {
|
||||||
|
it("classifies absent, successful, and failed sinks with original errors", async () => {
|
||||||
|
const log = { warn: vi.fn() };
|
||||||
|
const thrown = new Error("thrown");
|
||||||
|
const rejected = new Error("rejected");
|
||||||
|
|
||||||
|
await expect(emitBoundedRunAuditWithOutcome(undefined, event, { log })).resolves.toEqual({ outcome: "absent" });
|
||||||
|
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: 1 as never }, event, { log })).resolves.toEqual({ outcome: "absent" });
|
||||||
|
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => undefined }, event, { log })).resolves.toEqual({ outcome: "recorded" });
|
||||||
|
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => { throw thrown; } }, event, { log })).resolves.toEqual({ outcome: "failed", error: thrown });
|
||||||
|
await expect(emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => Promise.reject(rejected) }, event, { log })).resolves.toEqual({ outcome: "failed", error: rejected });
|
||||||
|
});
|
||||||
|
|
||||||
|
it("bounds never and late-settling sinks with explicit outcomes", async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const log = { warn: vi.fn() };
|
||||||
|
const never = emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => new Promise(() => undefined) }, event, { log });
|
||||||
|
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
|
||||||
|
await expect(never).resolves.toEqual({ outcome: "timed-out" });
|
||||||
|
|
||||||
|
let resolveLate!: () => void;
|
||||||
|
const lateResolve = emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => new Promise<void>((resolve) => { resolveLate = resolve; }) }, event, { timeoutMs: 1, log });
|
||||||
|
await vi.advanceTimersByTimeAsync(1);
|
||||||
|
resolveLate();
|
||||||
|
await expect(lateResolve).resolves.toEqual({ outcome: "timed-out" });
|
||||||
|
|
||||||
|
let rejectLate!: (error: Error) => void;
|
||||||
|
const lateReject = emitBoundedRunAuditWithOutcome({ recordRunAuditEvent: () => new Promise<void>((_resolve, reject) => { rejectLate = reject; }) }, event, { timeoutMs: 1, log });
|
||||||
|
await vi.advanceTimersByTimeAsync(1);
|
||||||
|
rejectLate(new Error("late rejection"));
|
||||||
|
await expect(lateReject).resolves.toEqual({ outcome: "timed-out" });
|
||||||
|
});
|
||||||
|
|
||||||
|
it("still invokes the void seam synchronously and resolves undefined for every sink state", async () => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const modes = [
|
||||||
|
undefined,
|
||||||
|
() => undefined,
|
||||||
|
() => { throw new Error("throw"); },
|
||||||
|
() => Promise.reject(new Error("reject")),
|
||||||
|
() => new Promise(() => undefined),
|
||||||
|
];
|
||||||
|
for (const recordRunAuditEvent of modes) {
|
||||||
|
const sink = vi.fn(recordRunAuditEvent);
|
||||||
|
const promise = emitBoundedRunAudit(recordRunAuditEvent === undefined ? undefined : { recordRunAuditEvent: sink }, event, { log: { warn: vi.fn() } });
|
||||||
|
if (recordRunAuditEvent !== undefined) expect(sink).toHaveBeenCalledTimes(1);
|
||||||
|
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
|
||||||
|
await expect(promise).resolves.toBeUndefined();
|
||||||
|
}
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|||||||
@@ -0,0 +1,76 @@
|
|||||||
|
import { afterEach, describe, expect, it, vi } from "vitest";
|
||||||
|
import { CORE_RUN_AUDIT_EMIT_TIMEOUT_MS } from "../run-audit/emit-bounded-run-audit.js";
|
||||||
|
import { reconcilePhantomCommittedReservationsAsync } from "../task-store/async/async-phantom-reservations.js";
|
||||||
|
|
||||||
|
type SinkMode = "absent" | "recorded" | "throw" | "reject" | "never" | "late-resolve" | "late-reject";
|
||||||
|
|
||||||
|
function createFixture(mode: SinkMode) {
|
||||||
|
let settle: (() => void) | undefined;
|
||||||
|
const sink = vi.fn(() => {
|
||||||
|
if (mode === "throw") throw new Error("forced audit failure");
|
||||||
|
if (mode === "reject") return Promise.reject(new Error("forced audit failure"));
|
||||||
|
if (mode === "never") return new Promise(() => undefined);
|
||||||
|
if (mode === "late-resolve") return new Promise<void>((resolve) => { settle = resolve; });
|
||||||
|
if (mode === "late-reject") return new Promise<void>((_resolve, reject) => { settle = () => reject(new Error("late audit failure")); });
|
||||||
|
return undefined;
|
||||||
|
});
|
||||||
|
const ids = ["FN-9182-a", "FN-9182-b"];
|
||||||
|
const reservations = ids.map((taskId) => ({ taskId, liveId: null, projectArchiveId: null, coldArchiveId: null }));
|
||||||
|
const query = {
|
||||||
|
from: vi.fn(), leftJoin: vi.fn(), where: vi.fn(), orderBy: vi.fn(async () => reservations),
|
||||||
|
};
|
||||||
|
query.from.mockReturnValue(query); query.leftJoin.mockReturnValue(query); query.where.mockReturnValue(query);
|
||||||
|
const txQuery = { from: vi.fn(), leftJoin: vi.fn(), where: vi.fn(async () => ids.map((taskId) => ({ taskId }))) };
|
||||||
|
txQuery.from.mockReturnValue(txQuery); txQuery.leftJoin.mockReturnValue(txQuery);
|
||||||
|
let deleteCount = 0;
|
||||||
|
const tx = {
|
||||||
|
select: vi.fn(() => txQuery),
|
||||||
|
delete: vi.fn(() => ({ where: vi.fn(() => ({ returning: vi.fn(async () => {
|
||||||
|
deleteCount += 1;
|
||||||
|
return deleteCount === 1 ? ids.map((taskId) => ({ taskId })) : [];
|
||||||
|
}) })) })),
|
||||||
|
};
|
||||||
|
const layer = {
|
||||||
|
projectId: "project-9182",
|
||||||
|
db: { select: vi.fn(() => query) },
|
||||||
|
transactionImmediate: async (work: (value: typeof tx) => unknown) => work(tx),
|
||||||
|
};
|
||||||
|
return {
|
||||||
|
store: {
|
||||||
|
getAsyncLayer: () => layer,
|
||||||
|
taskDir: (id: string) => `/definitely-missing-fn-9182/${id}`,
|
||||||
|
...(mode === "absent" ? {} : { recordRunAuditEvent: sink }),
|
||||||
|
},
|
||||||
|
sink,
|
||||||
|
settle: () => settle?.(),
|
||||||
|
ids,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
afterEach(() => vi.useRealTimers());
|
||||||
|
|
||||||
|
describe("phantom reservation audit outcomes", () => {
|
||||||
|
it.each(["absent", "recorded", "throw", "reject"] as const)("preserves reconciliation payloads for %s sinks", async (mode) => {
|
||||||
|
const fixture = createFixture(mode);
|
||||||
|
const result = await reconcilePhantomCommittedReservationsAsync(fixture.store as never);
|
||||||
|
if (mode === "throw" || mode === "reject") {
|
||||||
|
expect(result).toEqual({ reconciled: [], skipped: fixture.ids.map((id) => ({ id, reason: "audit-failed: forced audit failure" })) });
|
||||||
|
} else {
|
||||||
|
expect(result).toEqual({ reconciled: fixture.ids, skipped: [] });
|
||||||
|
}
|
||||||
|
if (mode !== "absent") expect(fixture.sink).toHaveBeenCalledTimes(2);
|
||||||
|
});
|
||||||
|
|
||||||
|
it.each(["never", "late-resolve", "late-reject"] as const)("bounds %s audit sinks as per-reservation skipped results", async (mode) => {
|
||||||
|
vi.useFakeTimers();
|
||||||
|
const fixture = createFixture(mode);
|
||||||
|
const operation = reconcilePhantomCommittedReservationsAsync(fixture.store as never);
|
||||||
|
await vi.advanceTimersByTimeAsync(CORE_RUN_AUDIT_EMIT_TIMEOUT_MS * fixture.ids.length);
|
||||||
|
fixture.settle();
|
||||||
|
await expect(operation).resolves.toEqual({
|
||||||
|
reconciled: [],
|
||||||
|
skipped: fixture.ids.map((id) => ({ id, reason: "audit-failed: timed-out" })),
|
||||||
|
});
|
||||||
|
expect(fixture.sink).toHaveBeenCalledTimes(2);
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -9,6 +9,12 @@ export type RunAuditSinkHost = {
|
|||||||
|
|
||||||
export type RunAuditLogger = { warn: (message: string) => void };
|
export type RunAuditLogger = { warn: (message: string) => void };
|
||||||
|
|
||||||
|
export type BoundedRunAuditResult =
|
||||||
|
| { outcome: "recorded" }
|
||||||
|
| { outcome: "absent" }
|
||||||
|
| { outcome: "failed"; error: unknown }
|
||||||
|
| { outcome: "timed-out" };
|
||||||
|
|
||||||
type RunAuditEvent = RunAuditEventInput | { mutationType: string; [key: string]: unknown };
|
type RunAuditEvent = RunAuditEventInput | { mutationType: string; [key: string]: unknown };
|
||||||
|
|
||||||
const defaultLog = createLogger("run-audit");
|
const defaultLog = createLogger("run-audit");
|
||||||
@@ -19,33 +25,52 @@ const defaultLog = createLogger("run-audit");
|
|||||||
* without creating its documented dependency cycle. Best-effort telemetry must not block, reject,
|
* without creating its documented dependency cycle. Best-effort telemetry must not block, reject,
|
||||||
* or otherwise alter the lifecycle operation which emitted it; no retry, queue, or backoff belongs here.
|
* or otherwise alter the lifecycle operation which emitted it; no retry, queue, or backoff belongs here.
|
||||||
*/
|
*/
|
||||||
export async function emitBoundedRunAudit(
|
export function emitBoundedRunAudit(
|
||||||
host: RunAuditSinkHost,
|
host: RunAuditSinkHost,
|
||||||
event: RunAuditEvent,
|
event: RunAuditEvent,
|
||||||
options: { timeoutMs?: number; log?: RunAuditLogger } = {},
|
options: { timeoutMs?: number; log?: RunAuditLogger } = {},
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
|
return emitBoundedRunAuditWithOutcome(host, event, options).then(() => undefined);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* FNXC:RunAudit 2026-08-20-07:14:
|
||||||
|
* FN-9182 supplies an explicit result for callers whose forensic throw ordering or caller-visible
|
||||||
|
* skipped payload derives from whether audit recording landed. This remains best-effort telemetry,
|
||||||
|
* not a durability guarantee: it bounds hostile sinks while preserving their original failure value.
|
||||||
|
*/
|
||||||
|
export async function emitBoundedRunAuditWithOutcome(
|
||||||
|
host: RunAuditSinkHost,
|
||||||
|
event: RunAuditEvent,
|
||||||
|
options: { timeoutMs?: number; log?: RunAuditLogger } = {},
|
||||||
|
): Promise<BoundedRunAuditResult> {
|
||||||
const log = options.log ?? defaultLog;
|
const log = options.log ?? defaultLog;
|
||||||
const sink = host?.recordRunAuditEvent;
|
const sink = host?.recordRunAuditEvent;
|
||||||
if (typeof sink !== "function") return;
|
if (typeof sink !== "function") return { outcome: "absent" };
|
||||||
|
|
||||||
let sinkPromise: Promise<unknown>;
|
let sinkPromise: Promise<unknown>;
|
||||||
try {
|
try {
|
||||||
|
// Invoke before the first await: synchronous timeline consumers observe every audit attempt.
|
||||||
sinkPromise = Promise.resolve(sink.call(host, event as RunAuditEventInput));
|
sinkPromise = Promise.resolve(sink.call(host, event as RunAuditEventInput));
|
||||||
} catch {
|
} catch (error) {
|
||||||
log.warn(`[run-audit] failed to record ${event.mutationType}`);
|
log.warn(`[run-audit] failed to record ${event.mutationType}`);
|
||||||
return;
|
return { outcome: "failed", error };
|
||||||
}
|
}
|
||||||
|
|
||||||
void sinkPromise.catch(() => undefined);
|
void sinkPromise.catch(() => undefined);
|
||||||
await new Promise<void>((resolve) => {
|
return new Promise<BoundedRunAuditResult>((resolve) => {
|
||||||
const timer = setTimeout(() => {
|
const timer = setTimeout(() => {
|
||||||
log.warn(`[run-audit] timed out recording ${event.mutationType}`);
|
log.warn(`[run-audit] timed out recording ${event.mutationType}`);
|
||||||
resolve();
|
resolve({ outcome: "timed-out" });
|
||||||
}, options.timeoutMs ?? CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
|
}, options.timeoutMs ?? CORE_RUN_AUDIT_EMIT_TIMEOUT_MS);
|
||||||
timer.unref?.();
|
timer.unref?.();
|
||||||
void sinkPromise.then(
|
void sinkPromise.then(
|
||||||
() => { clearTimeout(timer); resolve(); },
|
() => { clearTimeout(timer); resolve({ outcome: "recorded" }); },
|
||||||
() => { clearTimeout(timer); log.warn(`[run-audit] failed to record ${event.mutationType}`); resolve(); },
|
(error) => {
|
||||||
|
clearTimeout(timer);
|
||||||
|
log.warn(`[run-audit] failed to record ${event.mutationType}`);
|
||||||
|
resolve({ outcome: "failed", error });
|
||||||
|
},
|
||||||
);
|
);
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -4,6 +4,7 @@ import { and, asc, eq, inArray, isNull } from "drizzle-orm";
|
|||||||
import { alias } from "drizzle-orm/pg-core";
|
import { alias } from "drizzle-orm/pg-core";
|
||||||
import * as schema from "../../postgres/schema/index.js";
|
import * as schema from "../../postgres/schema/index.js";
|
||||||
import type { TaskStore } from "../../store.js";
|
import type { TaskStore } from "../../store.js";
|
||||||
|
import { emitBoundedRunAuditWithOutcome } from "../../run-audit/emit-bounded-run-audit.js";
|
||||||
|
|
||||||
export interface PhantomReservationReconcileResult {
|
export interface PhantomReservationReconcileResult {
|
||||||
reconciled: string[];
|
reconciled: string[];
|
||||||
@@ -144,24 +145,30 @@ export async function reconcilePhantomCommittedReservationsAsync(
|
|||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
if (pruned.prunedActivityLog > 0 || pruned.prunedAgents > 0) {
|
if (pruned.prunedActivityLog > 0 || pruned.prunedAgents > 0) {
|
||||||
try {
|
const auditResult = await emitBoundedRunAuditWithOutcome(store, {
|
||||||
await store.recordRunAuditEvent({
|
agentId: "self-healing",
|
||||||
agentId: "self-healing",
|
runId: `phantom-reservation:${taskId}`,
|
||||||
runId: `phantom-reservation:${taskId}`,
|
taskId,
|
||||||
taskId,
|
domain: "database",
|
||||||
domain: "database",
|
mutationType: "task:reconcile-phantom-committed-reservation",
|
||||||
mutationType: "task:reconcile-phantom-committed-reservation",
|
target: taskId,
|
||||||
target: taskId,
|
metadata: { reservationStatus: "committed", ...pruned },
|
||||||
metadata: { reservationStatus: "committed", ...pruned },
|
});
|
||||||
});
|
if (auditResult.outcome === "failed" || auditResult.outcome === "timed-out") {
|
||||||
} catch (error) {
|
|
||||||
/*
|
/*
|
||||||
FNXC:PostgresReservationRecovery 2026-07-14-21:55:
|
FNXC:PostgresReservationRecovery 2026-07-14-21:55:
|
||||||
Audit emission is isolated per reconciled reservation. One failed audit must not relabel earlier successful reconciliations as skipped or prevent later IDs from completing their own bookkeeping.
|
Audit emission is isolated per reconciled reservation. One failed audit must not relabel earlier successful reconciliations as skipped or prevent later IDs from completing their own bookkeeping.
|
||||||
|
|
||||||
|
FNXC:RunAudit 2026-08-20-07:14:
|
||||||
|
FN-9182 bounds this awaited audit and classifies its outcome. Preserve the
|
||||||
|
audit-failed prefix because callers already expose it per reservation; timeout
|
||||||
|
is a new bounded failure state rather than an indefinitely stalled sweep.
|
||||||
*/
|
*/
|
||||||
result.skipped.push({
|
result.skipped.push({
|
||||||
id: taskId,
|
id: taskId,
|
||||||
reason: `audit-failed: ${error instanceof Error ? error.message : String(error)}`,
|
reason: auditResult.outcome === "timed-out"
|
||||||
|
? "audit-failed: timed-out"
|
||||||
|
: `audit-failed: ${auditResult.error instanceof Error ? auditResult.error.message : String(auditResult.error)}`,
|
||||||
});
|
});
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -46,6 +46,10 @@ import { resolveWorkflowIrForTask } from "../workflows/workflow-ir-resolver.js";
|
|||||||
import { acquireTaskAdvisoryXactLock } from "./task-advisory-lock.js";
|
import { acquireTaskAdvisoryXactLock } from "./task-advisory-lock.js";
|
||||||
import { resolveProjectColumnsForRoles, REVIEW_ROLES } from "../project-lane-vocabulary.js";
|
import { resolveProjectColumnsForRoles, REVIEW_ROLES } from "../project-lane-vocabulary.js";
|
||||||
import type { InReviewDurationLanes } from "./async/async-audit.js";
|
import type { InReviewDurationLanes } from "./async/async-audit.js";
|
||||||
|
import { createLogger } from "../process/logger.js";
|
||||||
|
import { emitBoundedRunAuditWithOutcome } from "../run-audit/emit-bounded-run-audit.js";
|
||||||
|
|
||||||
|
const workflowDefinitionLog = createLogger("workflow-definitions");
|
||||||
|
|
||||||
export async function getAgentLogsByTimeRangeImpl(store: TaskStore,
|
export async function getAgentLogsByTimeRangeImpl(store: TaskStore,
|
||||||
taskId: string,
|
taskId: string,
|
||||||
@@ -849,33 +853,37 @@ export async function selectTaskWorkflowAndReconcileImpl(store: TaskStore,
|
|||||||
AWAITED, not fire-and-forget (PR #2512 review — CodeRabbit). The whole claim of
|
AWAITED, not fire-and-forget (PR #2512 review — CodeRabbit). The whole claim of
|
||||||
this branch is that the divergence is a fact on disk; `void`-ing the write and
|
this branch is that the divergence is a fact on disk; `void`-ing the write and
|
||||||
throwing on the next line meant the one artifact self-healing is meant to find
|
throwing on the next line meant the one artifact self-healing is meant to find
|
||||||
could silently be absent. The write is awaited and its own failure is swallowed
|
could silently be absent. The bounded outcome is still awaited before throwing,
|
||||||
so it can never mask the rejection the caller actually needs to see.
|
but a hostile sink cannot stall the switch forever.
|
||||||
|
|
||||||
NO ERROR PROSE (PR #2512 review — CodeRabbit). `outcome.error` is a propagated
|
NO ERROR PROSE (PR #2512 review — CodeRabbit). `outcome.error` is a propagated
|
||||||
`err.message`; persisting it would contradict this file's own "ids/columns/
|
`err.message`; persisting it would contradict this file's own "ids/columns/
|
||||||
outcomes-only" claim and the project rule that run-audit never stores error
|
outcomes-only" claim and the project rule that run-audit never stores error
|
||||||
prose. The bounded outcome code goes here; the human-readable reason travels on
|
prose. The bounded outcome code goes here; the human-readable reason travels on
|
||||||
the thrown error, which is not persisted.
|
the thrown error, which is not persisted.
|
||||||
|
|
||||||
|
FNXC:RunAudit 2026-08-20-07:14:
|
||||||
|
FN-9182 bounds the awaited forensic emission and reports its outcome without
|
||||||
|
changing this branch's typed rejection. A timeout is observable in logs rather
|
||||||
|
than indefinitely hiding a torn workflow switch from its caller.
|
||||||
*/
|
*/
|
||||||
try {
|
const auditResult = await emitBoundedRunAuditWithOutcome(store, {
|
||||||
await store.recordRunAuditEvent({
|
taskId,
|
||||||
taskId,
|
agentId: "system",
|
||||||
agentId: "system",
|
runId: `workflow-switch-torn-${taskId}`,
|
||||||
runId: `workflow-switch-torn-${taskId}`,
|
domain: "database",
|
||||||
domain: "database",
|
mutationType: "task:workflow-switch-torn",
|
||||||
mutationType: "task:workflow-switch-torn",
|
target: taskId,
|
||||||
target: taskId,
|
metadata: {
|
||||||
metadata: {
|
workflowId,
|
||||||
workflowId,
|
fromColumn,
|
||||||
fromColumn,
|
intendedColumn: decision.targetColumn,
|
||||||
intendedColumn: decision.targetColumn,
|
selectionCommitted: true,
|
||||||
selectionCommitted: true,
|
outcome: "rehome-rejected",
|
||||||
outcome: "rehome-rejected",
|
},
|
||||||
},
|
});
|
||||||
});
|
if (auditResult.outcome !== "recorded") {
|
||||||
} catch {
|
workflowDefinitionLog.warn(`[workflow-switch] audit ${auditResult.outcome} task=${taskId} workflow=${workflowId}`);
|
||||||
// An audit-write failure must not replace the rejection being reported.
|
|
||||||
}
|
}
|
||||||
throw new WorkflowSwitchRehomeFailedError({
|
throw new WorkflowSwitchRehomeFailedError({
|
||||||
taskId,
|
taskId,
|
||||||
|
|||||||
Reference in New Issue
Block a user