FN-9172: isolate executor run-audit emissions

Keep executor lifecycle transitions moving when run-audit sinks fail or stall.

- Add a shared bounded, best-effort executor audit emission seam.
- Route executor lifecycle audit writes through sink isolation.
- Cover absent, throwing, rejecting, hanging, and late-settling sinks.
- Document the emitter policy and publish a patch changeset.

Files changed:
 .changeset/fn-9172-executor-run-audit-isolation.md |   7 +
 AGENTS.md                                          |   7 +
 docs/run-audit.md                                  |   6 +-
 .../src/__tests__/emit-bounded-run-audit.test.ts   |  58 +++++
 .../executor-run-audit-emitter-isolation.test.ts   | 252 +++++++++++++++++++++
 .../src/executor/acquire-session-registry-path.ts  |   5 +-
 .../engine/src/executor/completion-finalization.ts |   3 +-
 .../engine/src/executor/create-task-done-tool.ts   |   3 +-
 .../engine/src/executor/emit-bounded-run-audit.ts  |  48 ++++
 .../executor/emit-merge-boundary-unproven-audit.ts |  68 ++----
 .../engine/src/executor/execute-workflow-graph.ts  |   5 +-
 .../engine/src/executor/handle-graph-failure.ts    |  18 +-
 ...dle-stale-in-review-parse-pause-abort-replay.ts |   7 +-
 ...ndle-stale-in-review-plan-pause-abort-replay.ts |   7 +-
 .../maybe-dispatch-workflow-work-engine.ts         |   7 +-
 .../src/executor/no-merge-complete-column.ts       |   3 +-
 .../reenter-paused-aborted-workflow-node.ts        |   7 +-
 .../src/executor/required-artifact-recovery.ts     |   3 +-
 packages/engine/src/executor/run-implementation.ts |   5 +-
 19 files changed, 433 insertions(+), 86 deletions(-)

Fusion-Task-Id: FN-9172

Fusion-Task-Lineage: 80c67c36-3bbb-4274-ac4c-9808c4cda3ce

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-19 20:35:19 -07:00
parent 84d9a59e7e
commit 3dea1bb851
19 changed files with 433 additions and 86 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Audit telemetry failures can no longer stall or abort task execution.
category: fix
dev: Routes executor telemetry through emitBoundedRunAudit with bounded sink isolation.

View File

@@ -275,6 +275,13 @@ Scoped exception (FN-5819/FN-8823): while project auto-merge is On, shared-branc
`testMode?: boolean` is now available in both project and global settings. If project `testMode === true` (or the resolved default provider is `"mock"` at any tier), every AI lane is forced to `mock/scripted`, overriding per-task and per-lane model selections. The dashboard exposes this via the Settings Modal "Enable test mode" toggle and a persistent "Test mode — no real AI calls" banner.
### Run Audit
<!--
FNXC:RunAudit 2026-08-20-03:12:
FN-9172 makes executor telemetry optional even when PostgreSQL or an extension sink stalls. Direct
`store.recordRunAuditEvent` calls under `packages/engine/src/executor/` are an anti-pattern: use
`emitBoundedRunAudit` so audit visibility never becomes a lifecycle dependency.
-->
- FN-9172: New executor run-audit emitters must use `emitBoundedRunAudit`; it absorbs absent, throwing, rejecting, hanging, and late-settling sinks without changing the owning branch.
- 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.

View File

@@ -72,4 +72,8 @@ Events that make durable-agent error states and their recovery inspectable.
## Maintenance contract
Adding a new catalogued run-audit event requires updating **both** the typed catalogue module (`packages/engine/src/run-audit/run-audit-catalogue.ts`) **and** this doc together — the parity test (`packages/engine/src/__tests__/run-audit-catalogue.test.ts`) fails if the documented event set and the catalogue module's set ever diverge, keeping the observability surface truthful as the real `DatabaseMutationType` union evolves. Removing an event likewise requires updating both in the same change.
Adding a new catalogued run-audit event requires updating **both** the typed catalogue module (`packages/engine/src/run-audit/run-audit-catalogue.ts`) **and** this doc together — the parity test (`packages/engine/src/__tests__/run-audit-catalogue.test.ts`) fails if the documented event set and the catalogue module's set ever diverge, keeping the observability surface truthful as the real `DatabaseMutationType` union evolves. Removing an event likewise requires updating both in the same change.
### Emit-seam policy
Executor telemetry must use `emitBoundedRunAudit`. 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.

View File

@@ -0,0 +1,58 @@
import { describe, expect, it, vi } from "vitest";
import { EXECUTOR_RUN_AUDIT_EMIT_TIMEOUT_MS, emitBoundedRunAudit } from "../executor/emit-bounded-run-audit.js";
const event = { taskId: "FN-9172", agentId: "executor", runId: "run", domain: "database", mutationType: "task:execution-blocked-parked", target: "FN-9172", metadata: { taskId: "FN-9172", outcome: "parked" } } as any;
describe("emitBoundedRunAudit", () => {
it.each([undefined, null, {}])("ignores an absent or non-function sink", async (store) => {
await expect(emitBoundedRunAudit(store as any, event)).resolves.toBeUndefined();
});
it("forwards the event unchanged to a healthy sink", async () => {
const recordRunAuditEvent = vi.fn().mockResolvedValue(undefined);
await emitBoundedRunAudit({ recordRunAuditEvent } as any, event);
expect(recordRunAuditEvent).toHaveBeenCalledOnce();
expect(recordRunAuditEvent).toHaveBeenCalledWith(event);
});
it.each([
["throws synchronously", vi.fn(() => { throw new Error("boom"); })],
["rejects", vi.fn().mockRejectedValue(new Error("boom"))],
])("contains a sink that %s", async (_name, recordRunAuditEvent) => {
await expect(emitBoundedRunAudit({ recordRunAuditEvent } as any, event)).resolves.toBeUndefined();
});
it("honors an override when a sink never settles", async () => {
vi.useFakeTimers();
try {
const pending = emitBoundedRunAudit({ recordRunAuditEvent: vi.fn(() => new Promise<void>(() => {})) } as any, event, { timeoutMs: 7 });
await vi.advanceTimersByTimeAsync(6);
let settled = false;
void pending.then(() => { settled = true; });
await Promise.resolve();
expect(settled).toBe(false);
await vi.advanceTimersByTimeAsync(1);
await expect(pending).resolves.toBeUndefined();
} finally { vi.useRealTimers(); }
});
it.each(["resolve", "reject"])("pre-observes a late %s after the default bound", async (outcome) => {
vi.useFakeTimers();
const unhandled = vi.fn();
process.on("unhandledRejection", unhandled);
try {
let settle!: () => void;
let reject!: (error: Error) => void;
const sinkPromise = new Promise<void>((resolve, rejectPromise) => { settle = resolve; reject = rejectPromise; });
const pending = emitBoundedRunAudit({ recordRunAuditEvent: vi.fn(() => sinkPromise) } as any, event);
await vi.advanceTimersByTimeAsync(EXECUTOR_RUN_AUDIT_EMIT_TIMEOUT_MS);
await pending;
if (outcome === "resolve") settle(); else reject(new Error("late"));
await Promise.resolve();
expect(unhandled).not.toHaveBeenCalled();
} finally {
process.off("unhandledRejection", unhandled);
vi.useRealTimers();
}
});
});

View File

@@ -0,0 +1,252 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import type { Task, TaskDetail, WorkflowIr } from "@fusion/core";
import { activeSessionRegistry } from "../agents/active-session-registry.js";
import { MAX_EXECUTE_REQUEUE_LOOP_CYCLES, TaskExecutor } from "../executor.js";
import { createTaskDoneTool } from "../executor/create-task-done-tool.js";
import { parkCompletedBlockedTask } from "../executor/completion-finalization.js";
import { acquireSessionRegistryPath } from "../executor/acquire-session-registry-path.js";
import { advanceNoMergeWorkflowToCompleteColumn } from "../executor/no-merge-complete-column.js";
import { recoverMissingRequiredArtifacts } from "../executor/required-artifact-recovery.js";
import { handleStaleInReviewPlanPauseAbortReplay } from "../executor/handle-stale-in-review-plan-pause-abort-replay.js";
import { WorkflowGraphTaskRunner } from "../workflows/workflow-graph-task-runner.js";
import { createMockStore } from "./executor-test-helpers.js";
const NO_MERGE_IR = {
version: "v2", id: "wf-audit", name: "audit", nodes: [], edges: [],
columns: [
{ id: "working", name: "Working", traits: [] },
{ id: "complete", name: "Complete", traits: [{ trait: "complete" }] },
],
} as unknown as WorkflowIr;
const sinkStates = {
rejected: () => vi.fn().mockRejectedValue(new Error("audit sink down")),
synchronousThrow: () => vi.fn(() => { throw new Error("sync boom"); }),
hanging: () => vi.fn(() => new Promise<void>(() => {})),
};
async function settleBounded<T>(invoke: () => Promise<T>, hanging: boolean): Promise<T> {
if (!hanging) return await invoke();
vi.useFakeTimers();
try {
const result = invoke();
await vi.advanceTimersByTimeAsync(2_000);
const settled = await result;
// FNXC:RunAudit 2026-08-20-03:24: Run the zero-delay lifecycle continuation scheduled after the bounded emit.
await vi.runOnlyPendingTimersAsync();
return settled;
} finally {
vi.useRealTimers();
}
}
/**
* FNXC:RunAudit 2026-08-20-03:18:
* FN-9172 regression coverage drives executor entry points instead of the emit seam alone. The
* lifecycle write/log/return following each historical shape must remain observable when the
* store's real audit method rejects, hangs, or throws synchronously.
*/
describe("FN-9172 executor run-audit emitter isolation", () => {
afterEach(() => activeSessionRegistry.clear());
it.each(Object.entries(sinkStates))("keeps the terminal dispatch-loop park and token persistence moving after a %s sink", async (state, makeSink) => {
const store = createMockStore() as any;
const task = {
id: "FN-LOOP", title: "loop", description: "", column: "todo", dependencies: [],
steps: [{ name: "Implement", status: "in-progress" }], currentStep: 0, log: [],
status: null, error: null, paused: false, userPaused: false, autoMerge: true,
executeRequeueLoopCount: MAX_EXECUTE_REQUEUE_LOOP_CYCLES - 1,
executeRequeueLoopSignature: "execute|implementation-incomplete",
} as any;
store.getTask.mockResolvedValue(task);
store.getSettings.mockResolvedValue({ executorToolFailureRetryCount: 0 });
store.updateTask.mockImplementation(async (_id: string, patch: object) => Object.assign(task, patch));
store.recordRunAuditEvent = makeSink();
const executor = new TaskExecutor(store, "/repo") as any;
const persistTokenUsage = vi.spyOn(executor, "persistTokenUsage").mockResolvedValue(undefined);
await settleBounded(() => executor.handleGraphFailure(task, {
disposition: "failed", outcome: "failure", visitedNodeIds: ["execute"],
context: { "node:execute:value": "implementation-incomplete" },
}), state === "hanging");
expect(task).toMatchObject({ status: "failed", error: expect.stringMatching(/^EXECUTION_DISPATCH_LOOP_EXHAUSTED:/) });
expect(persistTokenUsage).toHaveBeenCalledWith("FN-LOOP");
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ mutationType: "task:execution-dispatch-loop-terminalized" }));
});
it.each(Object.entries(sinkStates))("keeps the tool-failure retry lane schedulable after a %s sink", async (state, makeSink) => {
const store = createMockStore() as any;
const task = {
id: "FN-RETRY", title: "retry", description: "", column: "in-progress", dependencies: [],
steps: [{ name: "Implement", status: "in-progress" }], currentStep: 0, log: [],
status: null, error: null, paused: false, userPaused: false, autoMerge: true,
toolFailureDetectorLogCursor: 0,
} as any;
store.getTask.mockResolvedValue(task);
store.getSettings.mockResolvedValue({ executorToolFailureRetryCount: 1, executorToolFailureRetryBackoffMs: 0 });
store.getAgentLogCount = vi.fn().mockResolvedValue(1);
store.getAgentLogs = vi.fn().mockResolvedValue([{ type: "tool_error" }]);
store.claimNextToolFailureRetry = vi.fn().mockResolvedValue({ outcome: "claimed", attempt: 1 });
store.recordRunAuditEvent = makeSink();
const executor = new TaskExecutor(store, "/repo") as any;
executor.graphToolFailureRunCursors.set(task.id, 0);
const execute = vi.spyOn(executor, "execute").mockResolvedValue(undefined);
await settleBounded(() => executor.handleGraphFailure(task, {
disposition: "failed", outcome: "failure", visitedNodeIds: ["steps#0:step-execute"],
context: { "node:steps#0:step-execute:value": "failure" },
}), state === "hanging");
await new Promise((resolve) => setTimeout(resolve, 0));
expect(store.updateTask).toHaveBeenCalledWith("FN-RETRY", { status: null, error: null }, undefined);
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ mutationType: "task:execution-tool-failure-retry" }));
expect(execute).toHaveBeenCalledWith(task);
});
it.each(Object.entries(sinkStates))("keeps Shape-A no-merge completion moving after a %s sink", async (state, makeSink) => {
const recordRunAuditEvent = makeSink();
const moveTask = vi.fn().mockResolvedValue(undefined);
const store = {
getTaskWorkflowSelection: () => ({ workflowId: "wf-audit", stepIds: [] }),
getTaskWorkflowSelectionAsync: async () => ({ workflowId: "wf-audit", stepIds: [] }),
getWorkflowDefinition: async () => ({ ir: NO_MERGE_IR }),
moveTask,
recordRunAuditEvent,
} as any;
await settleBounded(
() => advanceNoMergeWorkflowToCompleteColumn(store, { id: "FN-A", column: "working", steps: [] } as TaskDetail),
state === "hanging",
);
expect(moveTask).toHaveBeenCalledWith("FN-A", "complete", expect.any(Object));
expect(recordRunAuditEvent).toHaveBeenCalledOnce();
});
it.each(Object.entries(sinkStates))("keeps Shape-A artifact recovery applying its follow-up after a %s sink", async (state, makeSink) => {
const recordRunAuditEvent = makeSink();
const task = { id: "FN-ARTIFACT", column: "todo", steps: [], recoveryRetryCount: 0 } as Task;
const updateTask = vi.fn().mockResolvedValue(undefined);
const store = {
getTask: vi.fn().mockResolvedValue(task),
getTaskWorkflowSelection: () => ({ workflowId: "builtin:coding", stepIds: [] }),
getTaskWorkflowSelectionAsync: async () => ({ workflowId: "builtin:coding", stepIds: [] }),
getWorkflowDefinition: async () => ({ ir: { version: "v2", columns: [{ id: "todo", traits: [{ trait: "hold", config: { release: "manual" } }] }] } }),
logEntry: vi.fn().mockResolvedValue(undefined),
moveTask: vi.fn().mockResolvedValue(undefined),
updateTask,
recordRunAuditEvent,
} as any;
await settleBounded(() => recoverMissingRequiredArtifacts({
store,
getRunContextFor: () => undefined,
isRequiredArtifactRecoveryProtected: async () => false,
workflowLifecycleMovesInFlight: new Set(),
}, task, ["plan"], { source: "graph-entry" }), state === "hanging");
expect(recordRunAuditEvent).toHaveBeenCalledOnce();
expect(updateTask).toHaveBeenCalledWith("FN-ARTIFACT", expect.objectContaining({ status: "needs-replan" }), undefined);
});
it.each(Object.entries(sinkStates))("keeps completed-blocked parks returning true after a %s sink", async (state, makeSink) => {
const task = { id: "FN-PARK", column: "todo", status: null, paused: false, userPaused: false, steps: [{ status: "done" }] } as any;
const updateTask = vi.fn().mockResolvedValue(undefined);
const moveTask = vi.fn().mockResolvedValue(undefined);
const store = {
getTask: vi.fn().mockResolvedValue(task),
getTaskWorkflowSelection: () => undefined,
getTaskWorkflowSelectionAsync: async () => undefined,
moveTask,
updateTask,
logEntry: vi.fn().mockResolvedValue(undefined),
recordRunAuditEvent: makeSink(),
} as any;
const result = await settleBounded(() => parkCompletedBlockedTask({
store, getRunContextFor: () => undefined, getTaskCompletionBlocker: async () => "dependency",
}, task, "dependency", "test", true), state === "hanging");
expect(result).toBe(true);
expect(updateTask).toHaveBeenCalledWith("FN-PARK", expect.objectContaining({ paused: true, status: "queued" }), undefined);
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ mutationType: "task:completed-blocked-parked" }));
});
it.each(Object.entries(sinkStates))("keeps the fn_task_done blocked exit persisting tokens after a %s sink", async (state, makeSink) => {
const task = { id: "FN-DONE", column: "in-progress", dependencies: [], log: [], steps: [{ status: "in-progress" }] } as any;
const persistTokenUsage = vi.fn().mockResolvedValue(undefined);
const store = {
getTask: vi.fn().mockResolvedValue(task), updateTask: vi.fn().mockResolvedValue(undefined),
moveTask: vi.fn().mockResolvedValue(undefined), logEntry: vi.fn().mockResolvedValue(undefined), recordRunAuditEvent: makeSink(),
} as any;
const tool = createTaskDoneTool({
store, getRunContextFor: () => undefined, persistTokenUsage, workflowLifecycleMovesInFlight: new Set(),
} as any, task.id, "/repo", "", new Map());
const result = await settleBounded(() => tool.execute("call", {
outcome: "blocked", reason: "upstream dependency", blockedBy: ["FN-9999"],
}), state === "hanging");
expect(result.content[0]?.text).toContain("parked as blocked");
expect(persistTokenUsage).toHaveBeenCalledWith("FN-DONE");
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ mutationType: "task:execution-blocked-parked" }));
});
it.each(Object.entries(sinkStates))("keeps Shape-B stale-plan replay returning true after a %s sink", async (state, makeSink) => {
const recordRunAuditEvent = makeSink();
const persistTokenUsage = vi.fn().mockResolvedValue(undefined);
const live = { id: "FN-B", column: "review", status: null, error: null, autoMerge: true, steps: [] } as TaskDetail;
const result = await settleBounded(() => handleStaleInReviewPlanPauseAbortReplay({
store: { getSettings: vi.fn().mockResolvedValue({ globalPause: false, enginePaused: false }), logEntry: vi.fn(), recordRunAuditEvent } as any,
getRunContextFor: () => undefined,
resolveResumeLanes: async () => ({ review: "review" }) as any,
isLiveSharedBranchGroupMember: async () => false,
clearPausedAborted: vi.fn(),
activeWorktrees: new Map(),
persistTokenUsage,
}, live, {
interruptedNodeId: "plan", visitedNodeIds: ["plan"], context: { "node:plan:value": "aborted" },
} as any, "global-pause", true, false), state === "hanging");
expect(result).toBe(true);
expect(persistTokenUsage).toHaveBeenCalledWith("FN-B");
expect(recordRunAuditEvent).toHaveBeenCalledOnce();
});
it.each(Object.entries(sinkStates))("keeps Shape-C workflow suspension returning after a %s sink", async (state, makeSink) => {
const store = createMockStore() as any;
const task = { id: "FN-C", column: "in-progress", steps: [], dependencies: [] } as any;
store.getTask.mockResolvedValue(task);
store.getSettings.mockResolvedValue({ experimentalFeatures: { workflowGraphExecutor: true } });
store.getTaskWorkflowSelectionAsync = vi.fn().mockResolvedValue({ workflowId: "wf-c", stepIds: [] });
store.getWorkflowDefinition = vi.fn().mockResolvedValue({ id: "wf-c", ir: { version: "v2", columns: [], nodes: [], edges: [] } });
store.recordRunAuditEvent = makeSink();
const run = vi.spyOn(WorkflowGraphTaskRunner.prototype, "run").mockResolvedValue({
disposition: "suspended", outcome: "failure", visitedNodeIds: [], suspension: { nodeId: "wait", reason: "manual", fromColumn: "in-progress", toColumn: "todo" },
} as any);
try {
const executor = new TaskExecutor(store, "/repo") as any;
await settleBounded(() => executor.executeWorkflowGraph(task), state === "hanging");
expect(store.recordRunAuditEvent).toHaveBeenCalledWith(expect.objectContaining({ mutationType: "task:workflow-run-suspended" }));
} finally {
run.mockRestore();
}
});
it.each(Object.entries(sinkStates))("keeps synchronous Shape-D session acquisition alive after a %s sink", async (_state, makeSink) => {
const recordRunAuditEvent = makeSink();
const path = "/repo/.worktrees/reclaimed";
activeSessionRegistry.registerPath(path, { taskId: "FN-OLD", kind: "executor", ownerKey: "old" });
const now = Date.now;
vi.spyOn(Date, "now").mockReturnValue(now() + 60_000);
try {
expect(() => acquireSessionRegistryPath({ store: { recordRunAuditEvent } as any, hasLiveTaskSessionSurface: () => false }, "FN-D", path, "executor", "new")).not.toThrow();
expect(activeSessionRegistry.lookupByPath(path)).toMatchObject({ taskId: "FN-D" });
// Fire-and-forget emission begins synchronously; it is never permitted to throw into this API.
expect(recordRunAuditEvent).toHaveBeenCalledOnce();
} finally {
vi.restoreAllMocks();
}
});
});

View File

@@ -22,6 +22,7 @@ import {
} from "../agents/active-session-registry.js";
import { executorLog } from "../logger.js";
import { generateSyntheticRunId } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
export type AcquireSessionRegistryPathDeps = {
store: TaskStore;
@@ -45,7 +46,7 @@ export function acquireSessionRegistryPath(
executorLog.warn(
`${taskId}: reclaimed a stale active-session entry on ${registryPath} from dead task ${outcome.holderTaskId} (idle ${outcome.ageMs}ms)`,
);
void deps.store.recordRunAuditEvent?.({
void emitBoundedRunAudit(deps.store, {
taskId,
agentId: "executor",
runId: generateSyntheticRunId("session-path-reclaim", taskId),
@@ -53,6 +54,6 @@ export function acquireSessionRegistryPath(
mutationType: "session:reclaim-stale-foreign-path",
target: taskId,
metadata: { taskId, holderTaskId: outcome.holderTaskId, kind, ageMs: outcome.ageMs },
})?.catch?.(() => undefined);
});
}
}

View File

@@ -8,6 +8,7 @@ import { evaluateSkipBypassTaint } from "@fusion/core";
import { COMPLETED_BLOCKED_PAUSE_REASON } from "../self-healing.js";
import { executorLog } from "../logger.js";
import { generateSyntheticRunId, type EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { isTaskWorkComplete } from "./task-predicates.js";
import {
resolveReboundColumnFor,
@@ -89,7 +90,7 @@ export async function parkCompletedBlockedTask(
}, deps.getRunContextFor(task.id));
executorLog.log(`${task.id}: ${message}`);
await deps.store.logEntry(task.id, message, undefined, deps.getRunContextFor(task.id));
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: task.id,
agentId: "executor",
runId: generateSyntheticRunId("completed-blocked-park", task.id),

View File

@@ -29,6 +29,7 @@ import {
import { moveTaskToReplanColumn, resolveReplanTargetColumn } from "../execution/replan-target.js";
import { mergeEffectiveSettings } from "../project/effective-settings.js";
import { generateSyntheticRunId, type EngineRunContext, type RunAuditor } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { executorLog } from "../logger.js";
import { resolveReboundColumnFor } from "./lifecycle-columns.js";
import { evaluateTaskDoneRefusal } from "./task-done-refusal.js";
@@ -228,7 +229,7 @@ export function createTaskDoneTool(
deps.getRunContextFor(taskId),
);
}
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId,
agentId: "executor",
runId: generateSyntheticRunId("execution-blocked", taskId),

View File

@@ -0,0 +1,48 @@
import type { RunAuditEventInput, TaskStore } from "@fusion/core";
import { executorLog } from "../logger.js";
export const EXECUTOR_RUN_AUDIT_EMIT_TIMEOUT_MS = 2_000;
/**
* FNXC:RunAudit 2026-08-20-03:02:
* FN-9172 requires executor audit writes to remain best-effort telemetry rather than lifecycle
* dependencies. Every sink state is swallow-log-and-bound: no retry, backoff, or queueing may
* abort, stall, delay, or alter the owning executor branch.
*/
export async function emitBoundedRunAudit(
store: TaskStore | null | undefined,
event: RunAuditEventInput,
options: { timeoutMs?: number } = {},
): Promise<void> {
const sink = store?.recordRunAuditEvent;
if (typeof sink !== "function") return;
let sinkPromise: Promise<unknown>;
try {
sinkPromise = Promise.resolve(sink.call(store, event));
} catch {
executorLog.warn(`[run-audit] failed to record ${event.mutationType}`);
return;
}
// Observe late rejection before the bounded wait returns so it cannot become unhandled.
void sinkPromise.catch(() => undefined);
await new Promise<void>((resolve) => {
const timer = setTimeout(() => {
executorLog.warn(`[run-audit] timed out recording ${event.mutationType}`);
resolve();
}, options.timeoutMs ?? EXECUTOR_RUN_AUDIT_EMIT_TIMEOUT_MS);
timer.unref?.();
void sinkPromise.then(
() => {
clearTimeout(timer);
resolve();
},
() => {
clearTimeout(timer);
executorLog.warn(`[run-audit] failed to record ${event.mutationType}`);
resolve();
},
);
});
}

View File

@@ -1,6 +1,6 @@
import type { TaskStore } from "@fusion/core";
import { executorLog } from "../logger.js";
import { generateSyntheticRunId } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import type { MergeBoundaryUnprovenReasonCode } from "./workflow-merge-boundary.js";
export const MERGE_BOUNDARY_UNPROVEN_AUDIT_EMIT_TIMEOUT_MS = 2_000;
@@ -30,53 +30,23 @@ export async function emitMergeBoundaryUnprovenParked(
store: TaskStore | null | undefined,
payload: MergeBoundaryUnprovenParkedAuditPayload,
): Promise<void> {
const sink = store?.recordRunAuditEvent;
if (typeof sink !== "function") return;
let sinkPromise: Promise<unknown>;
try {
sinkPromise = Promise.resolve(sink.call(store, {
await emitBoundedRunAudit(store, {
taskId: payload.taskId,
agentId: "executor",
runId: payload.runId ?? generateSyntheticRunId("merge-boundary-unproven-park", payload.taskId),
domain: "database",
mutationType: "task:merge-boundary-unproven-parked",
target: payload.taskId,
metadata: {
taskId: payload.taskId,
agentId: "executor",
runId: payload.runId ?? generateSyntheticRunId("merge-boundary-unproven-park", payload.taskId),
domain: "database",
mutationType: "task:merge-boundary-unproven-parked",
target: payload.taskId,
metadata: {
taskId: payload.taskId,
nodeId: payload.nodeId,
failureValue: payload.failureValue,
source: payload.source,
...(payload.reasonCode === undefined ? {} : { reasonCode: payload.reasonCode }),
...(payload.missingInstanceCount === undefined ? {} : { missingInstanceCount: payload.missingInstanceCount }),
priorColumn: payload.priorColumn,
priorStatus: payload.priorStatus ?? null,
outcome: payload.outcome,
},
}));
} catch {
executorLog.warn("[run-audit] failed to record task:merge-boundary-unproven-parked");
return;
}
// Observe late rejection before the bounded wait returns so it cannot become unhandled.
void sinkPromise.catch(() => undefined);
await new Promise<void>((resolve) => {
const timer = setTimeout(() => {
executorLog.warn("[run-audit] timed out recording task:merge-boundary-unproven-parked");
resolve();
}, MERGE_BOUNDARY_UNPROVEN_AUDIT_EMIT_TIMEOUT_MS);
timer.unref?.();
void sinkPromise.then(
() => {
clearTimeout(timer);
resolve();
},
() => {
clearTimeout(timer);
executorLog.warn("[run-audit] failed to record task:merge-boundary-unproven-parked");
resolve();
},
);
});
nodeId: payload.nodeId,
failureValue: payload.failureValue,
source: payload.source,
...(payload.reasonCode === undefined ? {} : { reasonCode: payload.reasonCode }),
...(payload.missingInstanceCount === undefined ? {} : { missingInstanceCount: payload.missingInstanceCount }),
priorColumn: payload.priorColumn,
priorStatus: payload.priorStatus ?? null,
outcome: payload.outcome,
},
}, { timeoutMs: MERGE_BOUNDARY_UNPROVEN_AUDIT_EMIT_TIMEOUT_MS });
}

View File

@@ -42,6 +42,7 @@ import {
import { getActiveNotificationService } from "../util/notifier.js";
import { executorLog } from "../logger.js";
import type { EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { takePreHeldExecutorSlot } from "../concurrency/concurrency.js";
import { resolveCompleteColumnFor } from "./lifecycle-columns.js";
import { nextPlanReviewAttemptCount, PLAN_REVIEW_FEEDBACK_HISTORY_LIMIT } from "../plan-review-feedback-history.js";
@@ -779,7 +780,7 @@ export async function executeWorkflowGraph(
* Record suspension so an invisible wait is greppable (ids/outcomes-only audit).
*/
const suspension = result.suspension;
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: task.id,
agentId: "executor",
runId: resolvedRunId ?? `workflow-run-suspended:${task.id}`,
@@ -796,7 +797,7 @@ export async function executeWorkflowGraph(
continuationNodeId: continuation?.nodeId ?? null,
continuationState: continuation?.state ?? null,
},
}).catch(() => undefined);
});
executorLog.log(
`[workflow-graph] ${task.id} suspended at node '${suspension?.nodeId ?? "unknown"}' (${suspension?.reason ?? "unknown"})`,
);

View File

@@ -30,6 +30,7 @@ import { getPromptPath } from "../execution/spec-staleness.js";
import { moveTaskToReplanColumn, resolveReplanTargetColumn } from "../execution/replan-target.js";
import { executorLog } from "../logger.js";
import { generateSyntheticRunId, type EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { MERGE_BOUNDARY_UNPROVEN_VALUE } from "../workflows/workflow-merge-nodes.js";
import { emitMergeBoundaryUnprovenParked } from "./emit-merge-boundary-unproven-audit.js";
import { PAUSE_ABORT_PARK_ERROR_MARKER, PAUSE_ABORT_PARK_OPERATOR_MARKER } from "../self-healing.js";
@@ -896,7 +897,12 @@ export async function handleGraphFailure(
executeRequeueLoopCount: nextCount,
executeRequeueLoopSignature: signature,
}, deps.getRunContextFor(task.id));
await deps.store.recordRunAuditEvent?.({
/*
* FNXC:RunAudit 2026-08-20-03:02:
* FN-9172 keeps terminal-park telemetry bounded because updateTask has already landed;
* an audit failure must not skip persistTokenUsage or this branch's return.
*/
await emitBoundedRunAudit(deps.store, {
taskId: task.id,
agentId: "executor",
runId: generateSyntheticRunId("execution-dispatch-loop", task.id),
@@ -1090,7 +1096,7 @@ export async function handleGraphFailure(
if (claim.outcome === "claimed") {
await deps.store.updateTask(task.id, { status: null, error: null }, deps.getRunContextFor(task.id));
await deps.store.logEntry(task.id, `Consecutive tool-call failures — auto-retrying same model (${claim.attempt}/${maxToolFailureRetries}) instead of parking`, undefined, deps.getRunContextFor(task.id));
await deps.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("tool-failure-retry", task.id), domain: "database", mutationType: "task:execution-tool-failure-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", attempt: claim.attempt, maxAttempts: maxToolFailureRetries, consecutiveToolFailures: threshold, mode: "same-model" } });
await emitBoundedRunAudit(deps.store, { taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("tool-failure-retry", task.id), domain: "database", mutationType: "task:execution-tool-failure-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", attempt: claim.attempt, maxAttempts: maxToolFailureRetries, consecutiveToolFailures: threshold, mode: "same-model" } });
const schedule = () => { void (async () => { const resume = await deps.store.getTask(task.id); if (resume && !resume.deletedAt && !resume.paused && !resume.userPaused && resume.column === wipColumn) await deps.execute(resume); })().catch((error) => executorLog.error(`${task.id}: tool-failure retry failed`, error)); };
const delay = resolveConsecutiveToolFailureRetryBackoffMs(settings);
setTimeout(schedule, delay).unref?.();
@@ -1147,7 +1153,7 @@ export async function handleGraphFailure(
}, deps.getRunContextFor(task.id));
if (claimedEscalation) {
await deps.store.logEntry(task.id, "Same-model retries exhausted — escalating to alternate model/node (one attempt) instead of parking", undefined, deps.getRunContextFor(task.id));
await deps.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-retry", task.id), domain: "database", mutationType: "task:execution-escalation-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hasModelTarget, hasNodeTarget, priorConsecutiveToolFailureRetryCount: priorEscalationRetryCount } });
await emitBoundedRunAudit(deps.store, { taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-retry", task.id), domain: "database", mutationType: "task:execution-escalation-retry", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hasModelTarget, hasNodeTarget, priorConsecutiveToolFailureRetryCount: priorEscalationRetryCount } });
if (!hasNodeTarget) {
const scheduleEscalation = () => { void (async () => { const resumeTask = await deps.store.getTask(task.id); if (resumeTask && !resumeTask.deletedAt && !resumeTask.paused && !resumeTask.userPaused && resumeTask.column === wipColumn) await deps.execute(resumeTask); })().catch((error) => executorLog.error(`${task.id}: escalation retry failed`, error)); };
const handle = setTimeout(scheduleEscalation, resolveConsecutiveToolFailureRetryBackoffMs(settings));
@@ -1187,10 +1193,10 @@ export async function handleGraphFailure(
}, deps.getRunContextFor(task.id));
if (!cursorOwnedTerminalPark) return;
if (await deps.store.markToolFailureRetryExhaustedAudit(task.id)) {
await deps.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("tool-failure-retry-exhausted", task.id), domain: "database", mutationType: "task:execution-tool-failure-retry-exhausted", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", attempts: maxToolFailureRetries, limit: maxToolFailureRetries, outcome: "terminal-park" } });
await emitBoundedRunAudit(deps.store, { taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("tool-failure-retry-exhausted", task.id), domain: "database", mutationType: "task:execution-tool-failure-retry-exhausted", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", attempts: maxToolFailureRetries, limit: maxToolFailureRetries, outcome: "terminal-park" } });
}
if (escalationAttemptFailed) {
await deps.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-exhausted", task.id), domain: "database", mutationType: "task:execution-escalation-exhausted", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hadModelTarget: escalationHadModelTarget, hadNodeTarget: escalationHadNodeTarget } });
await emitBoundedRunAudit(deps.store, { taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-exhausted", task.id), domain: "database", mutationType: "task:execution-escalation-exhausted", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hadModelTarget: escalationHadModelTarget, hadNodeTarget: escalationHadNodeTarget } });
}
executorLog.warn(`${task.id}: ${message}`);
await deps.store.logEntry(task.id, message, undefined, deps.getRunContextFor(task.id));
@@ -1225,7 +1231,7 @@ export async function handleGraphFailure(
return { error: message, status: "failed" };
}, deps.getRunContextFor(task.id));
if (!escalationTerminalParked) return;
await deps.store.recordRunAuditEvent?.({ taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-exhausted", task.id), domain: "database", mutationType: "task:execution-escalation-exhausted", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hadModelTarget: escalationHadModelTarget, hadNodeTarget: escalationHadNodeTarget } });
await emitBoundedRunAudit(deps.store, { taskId: task.id, agentId: "executor", runId: generateSyntheticRunId("escalation-exhausted", task.id), domain: "database", mutationType: "task:execution-escalation-exhausted", target: task.id, metadata: { taskId: task.id, nodeId: failedNode ?? "unknown", hadModelTarget: escalationHadModelTarget, hadNodeTarget: escalationHadNodeTarget } });
} else {
// status "failed" doubles as the self-healing exemption: review-task
// revival sweeps skip tasks carrying a non-null status, preventing the

View File

@@ -17,6 +17,7 @@ import { graphFailureValue, isStalePauseAbortParkFailure } from "./graph-failure
import { isTerminalMergeGraphFailureValue } from "./task-predicates.js";
import type { ResumeLanes } from "./resolve-resume-lanes.js";
import { generateSyntheticRunId, type EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { executorLog } from "../logger.js";
import { WORKFLOW_NODE_ENGINE_PAUSE_ABORT_KIND } from "../workflows/workflow-graph-executor.js";
@@ -96,8 +97,7 @@ export async function handleStaleInReviewParsePauseAbortReplay(
await deps.store.logEntry(live.id, message, undefined, deps.getRunContextFor(live.id));
await deps.store.logEntry(live.id, "Auto-recovered: retrying stale in-review parse pause/resume replay — failure notification suppressed", undefined, deps.getRunContextFor(live.id));
await deps.store.updateTask(live.id, { graphResumeRetryCount: nextRetries, status: null, error: null }, deps.getRunContextFor(live.id));
try {
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: live.id,
agentId: "executor",
runId: generateSyntheticRunId("workflow-stale-parse-retry", live.id),
@@ -114,9 +114,6 @@ export async function handleStaleInReviewParsePauseAbortReplay(
mode: "preserved-in-review-retry-graph",
},
});
} catch (error) {
executorLog.warn(`${live.id}: failed to record stale parse replay retry audit: ${error instanceof Error ? error.message : String(error)}`);
}
await deps.persistTokenUsage(live.id);
const scheduleRetry = () => {

View File

@@ -14,6 +14,7 @@ import { graphFailureValue, isMergeGraphFailure, isStalePauseAbortParkFailure }
import { isTerminalMergeGraphFailureValue } from "./task-predicates.js";
import type { ResumeLanes } from "./resolve-resume-lanes.js";
import { generateSyntheticRunId, type EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { executorLog } from "../logger.js";
import { WORKFLOW_NODE_ENGINE_PAUSE_ABORT_KIND } from "../workflows/workflow-graph-executor.js";
@@ -77,8 +78,7 @@ export async function handleStaleInReviewPlanPauseAbortReplay(
await deps.store.updateTask(live.id, { status: null, error: null }, deps.getRunContextFor(live.id));
await deps.store.logEntry(live.id, "Auto-recovered: cleared stale in-review plan pause/resume replay failure — failure notification suppressed", undefined, deps.getRunContextFor(live.id));
}
try {
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: live.id,
agentId: "executor",
runId: generateSyntheticRunId("workflow-stale-plan-replay", live.id),
@@ -94,9 +94,6 @@ export async function handleStaleInReviewPlanPauseAbortReplay(
mode: "preserved-in-review",
},
});
} catch (error) {
executorLog.warn(`${live.id}: failed to record stale plan replay audit: ${error instanceof Error ? error.message : String(error)}`);
}
await deps.persistTokenUsage(live.id);
return true;
}

View File

@@ -7,6 +7,7 @@ import type { Task, TaskDetail, TaskStore, WorkflowIr, WorkflowWorkEngineDispatc
import { getWorkflowExtensionRegistry, resolveWorkflowIrForTask } from "@fusion/core";
import { executorLog } from "../logger.js";
import { generateSyntheticRunId } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
export type MaybeDispatchWorkflowWorkEngineDeps = {
store: TaskStore;
@@ -73,8 +74,7 @@ export async function maybeDispatchWorkflowWorkEngine(
task.id,
result.message ?? `Workflow work engine ${extensionId} claimed execution`,
);
try {
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: task.id,
agentId: "workflow-work-engine",
runId: result.runId ?? generateSyntheticRunId("workflow-work-engine", task.id),
@@ -87,9 +87,6 @@ export async function maybeDispatchWorkflowWorkEngine(
pluginId: definition.pluginId,
},
});
} catch (error) {
executorLog.warn(`${task.id}: failed to record workflow work-engine claim audit: ${error instanceof Error ? error.message : String(error)}`);
}
return true;
}

View File

@@ -34,6 +34,7 @@ import {
} from "@fusion/core";
import { executorLog } from "../logger.js";
import { generateSyntheticRunId } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
export async function advanceNoMergeWorkflowToCompleteColumn(
store: TaskStore,
@@ -80,7 +81,7 @@ export async function advanceNoMergeWorkflowToCompleteColumn(
}
// ids/outcomes-only metadata — no prose, no node/run internals.
await store.recordRunAuditEvent?.({
await emitBoundedRunAudit(store, {
taskId: task.id,
agentId: "executor",
runId: generateSyntheticRunId("workflow-no-merge-completion", task.id),

View File

@@ -9,6 +9,7 @@ import type { TaskDetail, TaskStore } from "@fusion/core";
import type { WorkflowGraphTaskRunResult } from "../workflows/workflow-graph-task-runner.js";
import type { PausedAbortProvenance } from "./paused-abort-provenance.js";
import { generateSyntheticRunId, type EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { executorLog } from "../logger.js";
import type { ResumeLanes } from "./resolve-resume-lanes.js";
@@ -58,8 +59,7 @@ export async function reenterPausedAbortedWorkflowNode(
await deps.store.logEntry(live.id, message, undefined, deps.getRunContextFor(live.id));
await deps.store.logEntry(live.id, `Auto-recovered: re-entering paused-aborted workflow graph node '${nodeId}' — failure notification suppressed`, undefined, deps.getRunContextFor(live.id));
await deps.store.updateTask(live.id, { graphResumeRetryCount: nextRetries, status: null, error: null }, deps.getRunContextFor(live.id));
try {
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: live.id,
agentId: "executor",
runId: generateSyntheticRunId("workflow-node-reentry", live.id),
@@ -76,9 +76,6 @@ export async function reenterPausedAbortedWorkflowNode(
mode: preservedInReview ? "preserved-in-review" : live.column === reentryLanes.hold ? "reexecuted-from-todo" : "reentered-graph",
},
});
} catch (error) {
executorLog.warn(`${live.id}: failed to record paused-node graph re-entry audit: ${error instanceof Error ? error.message : String(error)}`);
}
await deps.persistTokenUsage(live.id);
const scheduleRetry = () => {

View File

@@ -7,6 +7,7 @@ import type { Task, TaskStore } from "@fusion/core";
import { computeRecoveryDecision, formatDelay, MAX_RECOVERY_RETRIES } from "../healing/recovery-policy.js";
import { moveTaskToReplanColumn, resolveReplanTargetColumn } from "../execution/replan-target.js";
import { generateSyntheticRunId, type EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { resolveTerminalColumnsFor } from "./lifecycle-columns.js";
export type RequiredArtifactRecoveryDeps = {
@@ -63,7 +64,7 @@ export async function recoverMissingRequiredArtifacts(
const context = deps.getRunContextFor(task.id);
const action = decision.shouldRetry ? "replan" : "park-failed";
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: task.id,
agentId: "executor",
runId: context?.runId ?? generateSyntheticRunId("required-artifact-missing", task.id),

View File

@@ -191,6 +191,7 @@ import { resolveDedicatedPlannerColumnsForTask } from "../planner-lane-resolutio
import { mergeEffectiveSettings } from "../project/effective-settings.js";
import { buildStepFailureMessage, emitProactiveStatus, sanitizeFailureReason } from "../project/proactive-status.js";
import { createRunAuditor, generateSyntheticRunId, type EngineRunContext } from "../util/run-audit.js";
import { emitBoundedRunAudit } from "./emit-bounded-run-audit.js";
import { acquireTaskWorktree, WorktreeBaseRefreshError } from "../worktree/worktree-acquisition.js";
import { resolveWorktreesDir } from "../worktree/worktree-paths.js";
import {
@@ -1811,7 +1812,7 @@ export async function runImplementation(
const taskCreateWithheld = !isAgentTaskCreateToolAvailable(settings, executionCallerIsEphemeral);
const delegateWithheld = !isAgentDelegateTaskToolAvailable(settings, executionCallerIsEphemeral);
if (taskCreateWithheld || delegateWithheld) {
await deps.store.recordRunAuditEvent?.({
await emitBoundedRunAudit(deps.store, {
taskId: task.id,
agentId: identityAgent?.id ?? "executor",
runId: deps.getRunContextFor(task.id)?.runId ?? generateSyntheticRunId("task-create-withheld", task.id),
@@ -1825,7 +1826,7 @@ export async function runImplementation(
withheldDelegateTask: delegateWithheld,
lane: "execution-session",
},
}).catch(() => undefined);
});
}
/*
FNXC:AgentProvisioningGate 2026-07-26-13:20: