fix(engine): restore shipped behaviors dropped by the wave-18 executor peel

Full-suite repair, engine harness cluster (~85 red suites). Two causes:
(1) graph dispatch now fails closed without options.agentStore
(FN-8764/FN-8821, intended) — the shared executor test harness now
provisions the workflow-routing agent-store fixture for bare
TaskExecutor constructions, with explicit opt-out for the two tests
asserting the fail-closed park; (2) the wave-18 'pure peel' (#3317)
rebuilt executor.ts from a stale base and silently deleted shipped
behaviors, restored here: FN-8864 agent-activity writers (task
started/handed-off, workflow gate pass/fail, gate principal attribution
via executor/workflow-gate-activity.ts), FN-8768 Plan Review group
recognition, convergence primer, and modified-file review scoping,
FN-6782's fire-time guard on transient resume-after-restart retries,
FN-8868 session usage telemetry boundaries, recommendation-route
withheld-tool guidance, and the per-instance worktree retry cap.
Stale expectations updated for intended changes (FN-8823 shared-member
hold, FN-9060 zero-acquire fail-closed, heartbeat tool inventory,
peeled-module seams, PG harness provisioning). Verified: 23 files /
711 tests green, engine typecheck clean, merge gate green.

Known follow-ups (not addressed here): step-session error routing may
still bypass FN-5866 non-continuable classification
(post-done-continuation-no-wedge red), scheduler mission-loop trigger
gap (mission-validation-trigger-gap red), and executeWorkflowStep lost
routed workflow-principal session identity threading (untested drop
from the same peel).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
gsxdsm
2026-08-15 15:52:52 -07:00
parent 987878bd17
commit 56a087dac9
26 changed files with 438 additions and 59 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Restore agent-activity telemetry, Plan Review convergence, and restart-retry safety guards lost in an executor refactor.
category: fix
dev: The wave-18 executor peel (#3317) was built from a stale base and silently dropped shipped behaviors; restored — FN-8864 agent-activity writers (task started/handed-off, workflow gate pass/fail, gate principal attribution via new `executor/workflow-gate-activity.ts`), FN-8768 Plan Review group recognition + convergence primer + modified-file review scoping, FN-6782's fire-time guard on transient resume-after-restart retries, FN-8868 session usage telemetry boundaries, recommendation-route withheld-tool guidance, and the per-instance worktree retry cap. Graph dispatch requiring `options.agentStore` (FN-8764/FN-8821) is intended behavior; the shared test harness now provisions it.

View File

@@ -29,6 +29,13 @@ import { dirname, join } from "node:path";
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import { BUILTIN_WORKFLOWS, type WorkflowIr } from "@fusion/core";
import "./executor-test-helpers.js";
// captureBaseCommitSha was peeled off TaskExecutor into executor/worktree-git-refs.ts (wave 18),
// so the old per-test `vi.spyOn(executor, "captureBaseCommitSha")` seam no longer exists.
// Stub it at its module home; everything else in that module stays real.
vi.mock("../executor/worktree-git-refs.js", async (importOriginal) => ({
...(await importOriginal() as object),
captureBaseCommitSha: vi.fn().mockResolvedValue(undefined),
}));
import { TaskExecutor } from "../executor.js";
import type { PluginRunner } from "../plugins/plugin-runner.js";
import { WorkflowGraphExecutor } from "../workflows/workflow-graph-executor.js";
@@ -334,7 +341,6 @@ describe("CE workflow-step executor integration", () => {
path: "/tmp/test/.worktrees/swift-falcon",
branch: "fusion/fn-ce-1",
});
vi.spyOn(executor as any, "captureBaseCommitSha").mockResolvedValue(undefined);
const captured: { step?: any; worktreePath?: string } = {};
vi.spyOn(executor as any, "executeWorkflowStep").mockImplementation(async (...args: any[]) => {
@@ -404,7 +410,6 @@ describe("CE workflow-step executor integration", () => {
path: "/tmp/test/.worktrees/fresh-ce-checkout",
branch: "fusion/fn-ce-1",
});
vi.spyOn(executor as any, "captureBaseCommitSha").mockResolvedValue(undefined);
const captured: { worktreePath?: string } = {};
vi.spyOn(executor as any, "executeWorkflowStep").mockImplementation(async (...args: any[]) => {
@@ -481,7 +486,6 @@ describe("CE workflow-step executor integration", () => {
path: "/tmp/test/.worktrees/acquired-code-review",
branch: "fusion/fn-ce-1",
});
vi.spyOn(executor as any, "captureBaseCommitSha").mockResolvedValue(undefined);
const executeStep = vi.spyOn(executor as any, "executeWorkflowStep").mockResolvedValue({ success: true, output: "APPROVE" });
const requirements: any[] = [];
const codeReview = {

View File

@@ -21,7 +21,9 @@ accepted the old catch-all `hard-cancel`.
*/
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
import "./executor-test-helpers.js";
import { TaskExecutor } from "../executor.js";
// isBenignInReviewPauseAbort was peeled off TaskExecutor into executor/graph-resume-predicates.ts
// (wave 18) and stays re-exported from executor.js; call the module function, not an instance method.
import { TaskExecutor, isBenignInReviewPauseAbort } from "../executor.js";
import { createMockStore, resetExecutorMocks } from "./executor-test-helpers.js";
import type { TaskDetail } from "@fusion/core";
@@ -250,7 +252,7 @@ describe("pause-abort provenance truthfulness (KB-PROV)", () => {
column: "in-review",
steps: [{ name: "Implement", status: "done" }],
});
const benign = (executor as any).isBenignInReviewPauseAbort(
const benign = isBenignInReviewPauseAbort(
live,
{ disposition: "failed", outcome: "failure", visitedNodeIds: ["plan", "execute"], context: {} },
provenance,
@@ -288,7 +290,7 @@ describe("pause-abort provenance truthfulness (KB-PROV)", () => {
const { executor } = makeExecutor();
const result = { disposition: "failed", outcome: "failure", visitedNodeIds: ["plan", "execute"], context: {} };
const classify = (column: string, reviewLane: string) =>
(executor as any).isBenignInReviewPauseAbort(
isBenignInReviewPauseAbort(
makeTask({ column, steps: [{ name: "Implement", status: "done" }] }),
result,
"engine-abort",

View File

@@ -5,7 +5,9 @@ import os from "node:os";
import path from "node:path";
import { EventEmitter } from "node:events";
import type { Task, TaskStore } from "@fusion/core";
import { TaskExecutor } from "../executor.js";
// captureBaseCommitSha was peeled off TaskExecutor into executor/worktree-git-refs.ts (wave 18);
// it now takes the store explicitly instead of reading it from the executor instance.
import { captureBaseCommitSha } from "../executor.js";
const hasGit = spawnSync("git", ["--version"], { stdio: "pipe" }).status === 0;
const describeIfGit = hasGit ? describe : describe.skip;
@@ -62,10 +64,9 @@ describeIfGit("captureBaseCommitSha (real git)", () => {
}
const store = createStore();
const executor = new TaskExecutor(store, repo);
const audit = { git: vi.fn().mockResolvedValue(undefined) };
await (executor as any).captureBaseCommitSha(makeTask(), repo, audit, { isResume: false });
await captureBaseCommitSha(store, makeTask(), repo, audit, { isResume: false });
const firstBase = (store.updateTask as any).mock.calls[0][1].baseCommitSha as string;
expect(firstBase).toBeTruthy();
@@ -74,7 +75,7 @@ describeIfGit("captureBaseCommitSha (real git)", () => {
// Resume of the same task: baseCommitSha must be preserved so diff math
// stays stable across sessions (FN-4309/FN-4383).
await (executor as any).captureBaseCommitSha(makeTask(firstBase), repo, audit, { isResume: true });
await captureBaseCommitSha(store, makeTask(firstBase), repo, audit, { isResume: true });
expect((store.updateTask as any).mock.calls).toHaveLength(1);

View File

@@ -349,7 +349,9 @@ describe("workflow routing fixture", () => {
it("suspends an unrouted graph run before opening an implementation session", async () => {
const store = createMockStore();
const executor = new TaskExecutor(store, "/tmp/test", {});
// Explicit `agentStore: undefined` opts out of the harness's default routing agent store
// (executor-test-helpers fills it for bare constructions) so the unrouted suspend stays testable.
const executor = new TaskExecutor(store, "/tmp/test", { agentStore: undefined });
await executor.execute({
id: "FN-routing", title: "Routing fixture", description: "", column: "in-progress",

View File

@@ -126,7 +126,9 @@ describe("workflow routing harness guards", () => {
};
store.getTask.mockResolvedValue(task as any);
await new TaskExecutor(store, "/tmp/test").execute(task as any);
// Explicit `agentStore: undefined` opts out of the harness's default routing agent store
// (executor-test-helpers fills it for bare constructions) so the fail-closed park stays testable.
await new TaskExecutor(store, "/tmp/test", { agentStore: undefined }).execute(task as any);
expect(() => selectImplementationSessionCall(
mockedCreateFnAgent.mock.calls.map(([options]) => options as { customTools?: Array<{ name?: string }> }),

View File

@@ -3,6 +3,35 @@ import type { Mock } from "vitest";
import { installTaskWorktreeIdentityGuard } from "../worktree/worktree-hooks.js";
import type * as ReviewerModule from "../execution/reviewer.js";
/*
FNXC:EngineTests 2026-08-15-01:20:
Graph dispatch made the agent store MANDATORY: admitWorkflowPrincipalBeforeNode (FN-8764/FN-8821,
2026-08-07..09) fails closed with `workflow-principal-routing-unavailable:no-agent-store:<role>`
before any model session is created, so a bare `new TaskExecutor(store, root)` never reaches
createFnAgent and every fn_task_done/customTools capture stays undefined. Hundreds of legacy
executor suites construct the executor bare; rather than edit each construction site, fill the
missing agentStore at the ONE seam every such test goes through — the TaskExecutor constructor —
with the same createWorkflowRoutingAgentStore fixture the migrated suites pass explicitly.
An options bag that mentions `agentStore` at all (including an explicit `agentStore: undefined`)
always wins, so suites asserting the fail-closed no-agent-store park keep that behavior by
opting out explicitly. Mirrors the withSessionDefaults precedent below: supply only what a
test omitted, never override what it controls.
*/
vi.mock("../executor.js", async (importOriginal) => {
const actual = (await importOriginal()) as Record<string, unknown> & {
TaskExecutor: new (store: unknown, rootDir: string, options?: Record<string, unknown>) => unknown;
};
class DefaultRoutedTaskExecutor extends actual.TaskExecutor {
constructor(store: any, rootDir: string, options: Record<string, unknown> = {}) {
const filled = "agentStore" in options
? options
: { ...options, agentStore: createWorkflowRoutingAgentStore(store).agentStore };
super(store, rootDir, filled);
}
}
return { ...actual, TaskExecutor: DefaultRoutedTaskExecutor };
});
// Mock external dependencies
vi.mock("../pi.js", () => ({
createFnAgent: vi.fn(),

View File

@@ -29,6 +29,9 @@ function createStore(overrides: Partial<Record<string, unknown>> = {}): TaskStor
logEntry: vi.fn().mockResolvedValue(undefined),
getSettings: vi.fn().mockResolvedValue({ autoMerge: false }),
getRunContextFor: vi.fn(),
// verifyWorktreeInvariants resolves the live external-execution route through store.getTask
// (FNXC:ExternalExecutionCheckout 2026-08-09-23:53); undefined falls back to the passed task.
getTask: vi.fn().mockResolvedValue(undefined),
on: emitter.on.bind(emitter),
...overrides,
}) as unknown as TaskStore & EventEmitter;
@@ -209,11 +212,27 @@ describeIfGit("U1 KTD2 — verifyWorktreeInvariants iterates per worktree, prese
expect(result.expected).toBe(BRANCH);
});
it("regression: a zero-acquire workspace task (empty map) verifies vacuously → {ok:true}", async () => {
/*
FN-9060 (ebd345da0f) intentionally replaced the old vacuous {ok:true} for an UNPROVEN zero-acquire
workspace map: completion verification and per-repo review now classify the empty map through the
shared classifyWorkspaceZeroAcquire predicate, so an unproven empty map fails closed as no_commits
and only proven commit-free work (noCommitsExpected / no-op sentinel / prompt-derived eligibility)
verifies clean. Assert both halves of that contract.
*/
it("regression (FN-9060): an UNPROVEN zero-acquire workspace task (empty map) fails closed as no_commits", async () => {
fx = await createWorkspaceFixture();
const executor = workspaceExecutor(fx);
const task = makeTask(TASK_ID, { branch: BRANCH, workspaceWorktrees: {} });
const result = await (executor as any).verifyWorktreeInvariants(task);
expect(result.ok).toBe(false);
expect(result.reason).toBe("no_commits");
});
it("regression (FN-9060): a PROVEN commit-free zero-acquire workspace task (noCommitsExpected) verifies clean", async () => {
fx = await createWorkspaceFixture();
const executor = workspaceExecutor(fx);
const task = makeTask(TASK_ID, { branch: BRANCH, workspaceWorktrees: {}, noCommitsExpected: true });
const result = await (executor as any).verifyWorktreeInvariants(task);
expect(result).toEqual({ ok: true });
});
});

View File

@@ -3388,7 +3388,8 @@ describe("executeHeartbeat", () => {
*/
// fn_artifact_register/list/view, agent config/provisioning, mission hierarchy, ideation, goals/evaluations/identity,
// task read discovery (incl. logs_read), workflow discovery/authoring, task promotion, bounded research, clarification, web fetch, memory, and fn_heartbeat_done.
expect(callArgs.customTools).toHaveLength(66);
// FN-8948 added fn_mission_reconcile and the feature-validation repair surface added fn_feature_repair_validation. Count rose 66→68; keep exact so new tools fail loudly.
expect(callArgs.customTools).toHaveLength(68);
expect(callArgs.customTools!.map((tool) => tool.name)).toEqual([
"fn_task_create",
"fn_task_log",
@@ -3411,6 +3412,7 @@ describe("executeHeartbeat", () => {
"fn_mission_update",
"fn_mission_set_status",
"fn_mission_delete",
"fn_mission_reconcile",
"fn_milestone_add",
"fn_milestone_update",
"fn_milestone_delete",
@@ -3419,6 +3421,7 @@ describe("executeHeartbeat", () => {
"fn_slice_delete",
"fn_feature_add",
"fn_feature_update",
"fn_feature_repair_validation",
"fn_feature_set_status",
"fn_feature_delete",
"fn_feature_link_task",

View File

@@ -158,12 +158,28 @@ export async function createPgLayer(): Promise<PgLayerFixture> {
runtimeUrl: testUrl,
migrationUrl: testUrl,
migrationUrlOverridden: false,
/*
This helper creates its own local postmaster endpoint, making it the test equivalent of a
lifecycle-proven direct session transport (mirrors core pg-test-harness, FNXC:PlanningDependencyReseed
2026-08-04-00:54). Without it, every hold-column task mutation that routes through the planning
lifecycle advisory lock throws `Planning lifecycle lock transport unavailable: direct-session-unavailable`
and the self-healing sweeps under test silently recover 0 rows.
*/
directSessionUrl: testUrl,
directSessionProvenance: "migration-override",
};
const schemaConn = await createConnectionSetFromUrl(backend, { poolMax: 1, connectTimeoutSeconds: 5 });
await applySchemaBaseline(schemaConn.migration);
await schemaConn.close();
const connections = await createConnectionSetFromUrl(backend, { poolMax: 5, connectTimeoutSeconds: 5 });
const layer = createAsyncDataLayer(connections);
/*
FN-8764 built-in workflow-owner provisioning during AgentStore.init() requires a bound
asyncLayer.projectId (the project-scoped advisory lock hashes it), so the reliability layer
binds one explicitly — mirroring the CLI extension harness rather than the project-agnostic
core default.
*/
const projectId = dbName;
const connections = await createConnectionSetFromUrl(backend, { poolMax: 5, connectTimeoutSeconds: 5, projectId });
const layer = createAsyncDataLayer(connections, { projectId });
return {
layer,
dbName,

View File

@@ -2,7 +2,8 @@ import { describe, expect, it, vi, beforeEach } from "vitest";
import type { TaskDetail } from "@fusion/core";
import "../executor-test-helpers.js";
import { PLAN_REVIEW_PROVIDER_FAILURE_HOLD_VALUE } from "../../workflows/workflow-graph-executor.js";
import { TaskExecutor } from "../../executor.js";
// graphFailureValue was peeled off TaskExecutor into executor/graph-failure-pure.ts (wave 18); use the re-exported free function.
import { TaskExecutor, graphFailureValue } from "../../executor.js";
import { activeSessionRegistry } from "../../agents/active-session-registry.js";
import {
createMockStore,
@@ -80,8 +81,7 @@ describe("graphFailureValue optional-group materialized ids", () => {
});
it("prefers the group's published value for a `group::template` failed node", () => {
const executor = new TaskExecutor(createMockStore(), "/tmp/test");
const value = (executor as any).graphFailureValue({
const value = graphFailureValue({
visitedNodeIds: ["plan-review", "plan-review::plan-review-step"],
context: {
"node:plan-review:value": PLAN_REVIEW_PROVIDER_FAILURE_HOLD_VALUE,
@@ -92,8 +92,7 @@ describe("graphFailureValue optional-group materialized ids", () => {
});
it("falls back to the unqualified template value when the group has none", () => {
const executor = new TaskExecutor(createMockStore(), "/tmp/test");
const value = (executor as any).graphFailureValue({
const value = graphFailureValue({
visitedNodeIds: ["plan-review::plan-review-step"],
context: { "node:plan-review-step:value": "exception" },
});
@@ -101,8 +100,7 @@ describe("graphFailureValue optional-group materialized ids", () => {
});
it("keeps resolving foreach `#` instance ids through the container key", () => {
const executor = new TaskExecutor(createMockStore(), "/tmp/test");
const value = (executor as any).graphFailureValue({
const value = graphFailureValue({
visitedNodeIds: ["steps#0:step-execute"],
context: { "node:steps:value": "awaiting-user-input" },
});

View File

@@ -116,8 +116,12 @@ describe("merge-node paused-abort retry classification (FN-6735)", () => {
});
it("allows shared-branch-group local integration to retry even when global autoMerge is off", async () => {
// FN-8823 (3dd824d04e): under project auto-merge Off, a shared-branch member is HELD unless its
// task explicitly opts in with autoMerge: true (see AGENTS.md "autoMerge: false callout"). The
// invariant this test guards — shared local integration retries under global Off — now requires
// that explicit member opt-in; an undefined member autoMerge is fenced by design.
const { store, task, executor, mergeRequester } = makeHarness({
autoMerge: undefined,
autoMerge: true,
branchContext: { groupId: "BG-6735", source: "mission", assignmentMode: "shared" },
}, { autoMerge: false });

View File

@@ -15,6 +15,8 @@
import type { Settings, Task, TaskStore } from "@fusion/core";
import { resolveExecutorFallbackModel, resolvePersistAgentThinkingLog } from "@fusion/core";
import { AgentLogger } from "../agents/agent-logger.js";
// FNXC:CommandCenterActivity 2026-08-15-22:15: FN-8868 usage telemetry + session boundaries (restored post-wave-18).
import { attachAgentUsageTelemetry, emitAgentSessionStart } from "../agents/agent-usage-telemetry.js";
import {
createResolvedAgentSession,
resolveExecutorSessionModel,
@@ -82,6 +84,8 @@ export async function attemptExecutorVerificationFix(
onAgentText: deps.onAgentText,
onAgentTool: deps.onAgentTool,
});
// FNXC:CommandCenterActivity 2026-08-15-22:15: FN-8868 usage telemetry (restored post-wave-18).
attachAgentUsageTelemetry(logger, { store: deps.store, agentId: task.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, lane: "executor" });
// Build skill selection context
let skillContext: Awaited<ReturnType<typeof buildSessionSkillContext>> | undefined;
@@ -109,6 +113,7 @@ export async function attemptExecutorVerificationFix(
task.credentialInstanceId,
);
const { provider: executorProvider, modelId: executorModelId } = executorSessionModel;
attachAgentUsageTelemetry(logger, { store: deps.store, agentId: task.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, model: executorModelId ?? null, provider: executorProvider ?? null, lane: "executor" });
const executorFallback = resolveExecutorFallbackModel(settings);
@@ -157,6 +162,8 @@ Do not refactor, rename broadly, or make opportunistic improvements.
...(skillContext?.skillSelectionContext ? { skillSelection: skillContext.skillSelectionContext } : {}),
...(skillContext && skillContext.additionalSkillPaths.length > 0 ? { additionalSkillPaths: skillContext.additionalSkillPaths } : {}),
});
// FNXC:CommandCenterActivity 2026-08-15-22:15: session boundary for the verification-fix runtime session (restored post-wave-18).
emitAgentSessionStart({ store: deps.store, agentId: task.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, model: executorModelId ?? null, provider: executorProvider ?? null, lane: "executor" });
await deps.store.logEntry(
task.id,

View File

@@ -144,6 +144,8 @@ export function buildExecuteWorkflowGraphDeps(host: any): any {
"graphRethinkNarrations", "graphRouting", "graphSeamGoverningNodeId", "graphSeamSkillName",
"graphSeamThinkingLevel", "graphStepActiveContext", "graphStepRunOnce", "graphStepSessionPinned",
"graphToolFailureRunCursors", "graphUnattendedRuns", "outerConcurrencyClaims",
// FNXC:AgentActivityStream 2026-08-15-22:15: FN-8864 gate-attribution retention map (restored post-wave-18).
"workflowGateActivityPrincipals",
]),
...facadeMethods(host, [
"getRunContextFor", "advanceNoMergeWorkflowToCompleteColumn", "applyGraphRethinkReset",
@@ -163,9 +165,14 @@ export function buildHandleGraphFailureDeps(host: any): any {
store: host.store,
rootDir: host.rootDir,
options: host.options as { stuckTaskDetector?: { untrackTask?: (taskId: string) => void }; [k: string]: unknown },
// FNXC:WorkflowLifecycle 2026-08-15-22:15: liveness surfaces for the transient-resume fire-time guard
// (restored post-wave-18 — the peel replaced the guarded scheduled retry with an unguarded execute()).
processWideGraphRouting: host.constructor.processWideGraphRouting as Set<string>,
...facadeFields(host, [
"activeWorktrees", "completionFinalizedTaskIds", "graphExecuteSelfRequeued",
"graphToolFailureRunCursors", "pausedAborted", "pausedAbortProvenance", "userCanceledTaskIds",
"executing", "resumingUnpaused", "activeSessions", "activeStepExecutors",
"activeWorkflowStepSessions", "activeCliTaskSessions", "activeWorkflowGraphAbortControllers",
]),
...facadeMethods(host, [
"getRunContextFor", "clearCompletedTaskWatchdog", "clearPausedAborted", "execute",
@@ -1343,7 +1350,12 @@ export function buildRouteResetParsePinMismatchToRetryDeps(host: any): any {
export function buildCreateWorktreeFacadeDeps(host: any, tryCreateWorktree: any): any {
return buildCreateWorktreeDeps(
host,
{ maxWorktreeRetries: MAX_WORKTREE_RETRIES, worktreeRetryDelaysMs: [...WORKTREE_RETRY_DELAYS] },
/*
FNXC:CodeOrganization 2026-08-15-22:15: pre-peel executor.ts read `this.MAX_WORKTREE_RETRIES`
(a private instance field), so an instance-level override was part of the contract (tests pin it
to 1 to keep retry-exhaustion paths fast). Honor that override before the module constant.
*/
{ maxWorktreeRetries: (host as { MAX_WORKTREE_RETRIES?: number }).MAX_WORKTREE_RETRIES ?? MAX_WORKTREE_RETRIES, worktreeRetryDelaysMs: [...WORKTREE_RETRY_DELAYS] },
tryCreateWorktree,
);
}

View File

@@ -28,7 +28,9 @@ import {
resolveWorkflowIrForTask,
upsertWorkflowStepResult,
applySupersededFindingIds,
isTerminalStepResult,
} from "@fusion/core";
import { resolveWorkflowGateActivityClaim } from "./workflow-gate-activity.js";
import type { ImplementationExit } from "./implementation-exit.js";
import type { WorkflowGraphTaskRunResult } from "../workflows/workflow-graph-task-runner.js";
import { WorkflowGraphTaskRunner } from "../workflows/workflow-graph-task-runner.js";
@@ -78,6 +80,8 @@ export type ExecuteWorkflowGraphDeps = {
graphStepSessionPinned: Set<string>;
graphToolFailureRunCursors: Map<string, number>;
graphUnattendedRuns: Set<string>;
/** FNXC:AgentActivityStream 2026-08-15-22:15: FN-8864 node-scoped routed-principal retention for gate attribution (restored post-wave-18). */
workflowGateActivityPrincipals: Map<string, string>;
outerConcurrencyClaims: Set<string>;
processWideGraphRouting: Set<string>;
getRunContextFor: (taskId: string) => EngineRunContext | undefined;
@@ -134,7 +138,8 @@ export function clearPrincipalHoldBackoff(taskId: string): void {
* the same write. Exported for production-shaped graph-writer tests.
*/
export async function persistWorkflowStepResult(
deps: Pick<ExecuteWorkflowGraphDeps, "store" | "getRunContextFor" | "readTaskArtifact">,
deps: Pick<ExecuteWorkflowGraphDeps, "store" | "getRunContextFor" | "readTaskArtifact">
& Partial<Pick<ExecuteWorkflowGraphDeps, "workflowGateActivityPrincipals" | "activeWorkflowPrincipals">>,
taskId: string,
result: CoreWorkflowStepResult,
): Promise<void> {
@@ -191,8 +196,76 @@ export async function persistWorkflowStepResult(
} else {
await deps.store.updateTask(taskId, { workflowStepResults: existing }, deps.getRunContextFor(taskId));
}
} catch {
// Result recording is additive visibility — never affect the graph run.
/*
FNXC:AgentActivityStream 2026-08-09-09:38 (restored 2026-08-15-22:15 after wave-18 shell-ification dropped it):
Terminal graph gate results are emitted at this shared persistence sink, not at individual node
implementations (FN-8864). Node ids are operator-authored, so metadata sanitation records unknown
ids as the closed `custom` enum rather than retaining prose.
FNXC:AgentActivityStream 2026-08-09-11:50:
Workflow `skipped` is terminal and non-blocking, so it is a passed gate for activity consumers;
advisory failures and failures remain failed. Preserve the exact closed status in metadata rather
than deriving a replacement that loses the gate outcome.
*/
if (isTerminalStepResult(result)) {
const persistedResult = existing.find((entry) => entry.workflowStepId === result.workflowStepId) ?? resultToPersist;
const passed = result.status === "passed"
|| result.status === "skipped"
|| result.verdict === "APPROVE"
|| result.verdict === "APPROVE_WITH_NOTES"
|| result.verdict === "CLOSE_NO_OP";
try {
await deps.store.recordAgentActivity({
type: passed ? "workflow:gate-passed" : "workflow:gate-failed",
/*
FNXC:AgentActivityStream 2026-08-09-13:30:
A workflow gate belongs to the principal that actually ran its node. The active
routing fence preserves reviewer overrides and column bindings; falling back to
the task assignee is only for unclassified nodes that have no routed principal.
*/
attributionClaim: resolveWorkflowGateActivityClaim(
deps.workflowGateActivityPrincipals?.get(`${taskId}\0${result.workflowStepId}`)
?? deps.activeWorkflowPrincipals?.get(taskId)?.agentId,
live?.assignedAgentId,
),
taskId,
occurredAt: result.completedAt ?? result.startedAt ?? new Date().toISOString(),
/*
FNXC:AgentActivityStream 2026-08-09-19:03:
`priorAttempts` is intentionally bounded, so its length cannot identify retries:
after the retention cap it would make later gate attempts collide and disappear.
A graph attempt's persisted startedAt is its natural, replay-stable identity;
pending→terminal updates keep that value while a new dispatch gets a new one.
*/
discriminator: `${result.workflowStepId}:${result.startedAt ?? result.completedAt ?? result.status}`,
metadata: {
stepId: result.workflowStepId,
status: result.status,
attempt: persistedResult.priorAttempts?.length ?? 0,
},
});
/*
FNXC:AgentActivityStream 2026-08-09-13:59:
Once the terminal event is durable, discard the retained node identity so a later
run cannot inherit an earlier gate's routed principal.
*/
deps.workflowGateActivityPrincipals?.delete(`${taskId}\0${result.workflowStepId}`);
} catch (error) {
/*
FNXC:AgentActivityStream 2026-08-09-13:43:
Activity is observability only: warn so a failed append is diagnosable, but never
let it change the workflow gate result or interrupt graph execution.
*/
executorLog.warn(`[agent-activity] ${taskId}: failed to record workflow gate activity: ${error instanceof Error ? error.message : String(error)}`);
}
}
} catch (error) {
/*
FNXC:AgentActivityStream 2026-08-09-13:43:
Persisting the underlying step and its activity row is additive visibility. Log a
failed persistence attempt without converting an otherwise valid graph run into a failure.
*/
executorLog.warn(`[agent-activity] ${taskId}: failed to persist workflow step result: ${error instanceof Error ? error.message : String(error)}`);
}
}
@@ -444,6 +517,7 @@ export async function executeWorkflowGraph(
workflowAgentCapacity: deps.workflowAgentCapacity,
activeWorkflowAuthorities: deps.activeWorkflowAuthorities,
activeWorkflowPrincipals: deps.activeWorkflowPrincipals,
workflowGateActivityPrincipals: deps.workflowGateActivityPrincipals,
workflowCapacityAttemptIds,
directWorkflowPrincipalWorkItemIds,
directWorkflowPrincipalHeldWorkItemIds,
@@ -797,6 +871,10 @@ export async function executeWorkflowGraph(
}
deps.activeWorkflowAuthorities.delete(task.id);
deps.activeWorkflowPrincipals.delete(task.id);
// FNXC:AgentActivityStream 2026-08-15-22:15: drop FN-8864 gate-attribution retention for this run (restored post-wave-18).
for (const key of deps.workflowGateActivityPrincipals.keys()) {
if (key.startsWith(`${task.id}\0`)) deps.workflowGateActivityPrincipals.delete(key);
}
if (graphAbortController && deps.activeWorkflowGraphAbortControllers.get(task.id) === graphAbortController) {
deps.activeWorkflowGraphAbortControllers.delete(task.id);
}

View File

@@ -20,6 +20,7 @@ import {
applyReviewSeverityGate,
isOpenWorkflowReviewFinding,
MAX_WORKFLOW_REVIEW_FINDINGS,
PLAN_REVIEW_GROUP_ID,
finalizePlanningSegment,
resolveExecutorFallbackModel,
resolvePersistAgentThinkingLog,
@@ -81,6 +82,10 @@ import {
type WorkflowStepOutcome,
} from "./workflow-step-verdict.js";
import { resolveDiffBaseRef } from "./worktree-git-refs.js";
// FNXC:PlanReviewConvergence 2026-08-15-22:15: FN-8768 convergence primer + revision-key classifier (restored post-wave-18).
import { buildGraphPlanReviewConvergenceContext, optionalStepRevisionKey } from "./optional-step-revision.js";
// FNXC:CommandCenterActivity 2026-08-15-22:15: FN-8868 usage telemetry + session boundaries (restored post-wave-18).
import { attachAgentUsageTelemetry, emitAgentSessionStart } from "../agents/agent-usage-telemetry.js";
const execAsync = promisify(exec);
@@ -127,7 +132,6 @@ export async function executeWorkflowStep(
// assumptions and proceed instead of parking on a question. Explicit opt-in
// only (default false = board run); see runGraphCustomNode / KTD-3.
const unattended = stepOptions?.unattended === true;
const isPlanReviewStep = workflowStep.id === "graph:plan-review-step" || workflowStep.name === "Plan Review";
/*
FNXC:WorkflowReviewFindings 2026-08-05-06:29:
reviewKind is carried from graph synthesis (cfg.reviewKind / optional-group context) so prompt
@@ -141,6 +145,15 @@ export async function executeWorkflowStep(
requireExternalIntegrationEvidence?: boolean;
};
const optionalGroupId = workflowStepMetadata.optionalGroupId;
/*
FNXC:PlanReviewConvergence 2026-08-04-06:35 (FN-8768; restored 2026-08-15-22:15 after the wave-18
executor.ts shell-ification dropped it): a RENAMED inner step of the canonical Plan Review optional
group is still Plan Review — classify by group id, not only by the default id/name.
*/
const isPlanReviewStep = workflowStep.id === "graph:plan-review-step"
|| workflowStep.name === "Plan Review"
|| optionalGroupId === PLAN_REVIEW_GROUP_ID;
const planReviewRevisionKey = optionalStepRevisionKey(optionalGroupId, workflowStep.name);
const isReviewTypeWorkflowStep =
isPlanReviewStep
|| workflowStepMetadata.reviewCanFixInline === true
@@ -182,6 +195,11 @@ export async function executeWorkflowStep(
}
const workflowReviewSpecText = typeof workflowReviewSpecArtifact === "string" ? workflowReviewSpecArtifact : "";
const planReviewSpecText = isPlanReviewStep ? workflowReviewSpecText : "";
// FNXC:PlanReviewConvergence 2026-08-04-06:35 (FN-8768; restored 2026-08-15-22:15): cumulative
// prior-feedback primer + attempt-three severity ratchet for repeat Plan Review attempts.
const planReviewConvergenceContext = isPlanReviewStep
? buildGraphPlanReviewConvergenceContext(task, planReviewRevisionKey)
: "";
/*
FNXC:PlanReview 2026-07-21-16:30:
@@ -265,33 +283,41 @@ export async function executeWorkflowStep(
const approvedContractBlock = isReviewTypeWorkflowStep && !isPlanReviewStep
? `
Approved Task Contract:
- PROMPT.md is the authoritative current contract for this review. It includes any approved planning revisions and scope decisions.
- The Task Description is historical input only. Do not enforce superseded requirements from the original Task Description when they conflict with PROMPT.md.
- Do not request behavior that PROMPT.md explicitly defers, excludes, or forbids. Review the implementation against the approved contract reproduced below.
- Scope exclusions do not waive security, correctness, or data-integrity defects in the approved implementation.
Approved Task Contract:
- PROMPT.md is the authoritative current contract for this review. It includes any approved planning revisions and scope decisions.
- The Task Description is historical input only. Do not enforce superseded requirements from the original Task Description when they conflict with PROMPT.md.
- Do not request behavior that PROMPT.md explicitly defers, excludes, or forbids. Review the implementation against the approved contract reproduced below.
- Scope exclusions do not waive security, correctness, or data-integrity defects in the approved implementation.
--- BEGIN APPROVED PROMPT.md ---
${workflowReviewSpecText}
--- END APPROVED PROMPT.md ---`
--- BEGIN APPROVED PROMPT.md ---
${workflowReviewSpecText}
--- END APPROVED PROMPT.md ---`
: "";
/*
FNXC:CodeReviewCompleteness 2026-08-04-00:20 (FN-8768 / #3327; restored 2026-08-15-22:15 after the
wave-18 executor.ts shell-ification regressed this block to its pre-FN-8768 wording):
The modified-file list is a starting scope, not a read prohibition — reviewers may read callers,
helpers, and tests needed to validate the change, while unrelated pre-existing issues stay out of
scope. Plan Review appends the convergence primer so repeat attempts stop re-raising settled blockers.
*/
const scopeBlock = isPlanReviewStep
? `Plan Review Scope:
- Review the task plan artifact (PROMPT.md), reproduced verbatim below, and task metadata only.
- The plan is embedded in this prompt — do NOT go looking for a PROMPT.md file in the worktree; it lives at the project root (\`.fusion/tasks/${task.id}/PROMPT.md\`), outside this worktree, so review the embedded copy.
- Do NOT judge current implementation diffs, uncommitted worktree changes, or unrelated repository changes.
- If the plan is internally consistent, complete, scoped, and verifiable, approve even when the worktree contains unrelated changes from another task.
- Review the task plan artifact (PROMPT.md), reproduced verbatim below, and task metadata only.
- The plan is embedded in this prompt — do NOT go looking for a PROMPT.md file in the worktree; it lives at the project root (\`.fusion/tasks/${task.id}/PROMPT.md\`), outside this worktree, so review the embedded copy.
- Do NOT judge current implementation diffs, uncommitted worktree changes, or unrelated repository changes.
- If the plan is internally consistent, complete, scoped, and verifiable, approve even when the worktree contains unrelated changes from another task.
--- BEGIN PROMPT.md ---
${planReviewSpecText}
--- END PROMPT.md ---`
--- BEGIN PROMPT.md ---
${planReviewSpecText}
--- END PROMPT.md ---${planReviewConvergenceContext ? `\n\n${planReviewConvergenceContext}` : ""}`
: `Diff Scope (files changed by THIS task vs base):
${scopeFileBlock}${diffShortstat ? `\nDiff stat: ${diffShortstat}` : ""}
${scopeFileBlock}${diffShortstat ? `\nDiff stat: ${diffShortstat}` : ""}
CRITICAL SCOPING RULES — read before doing anything else:
- Review ONLY the files listed above. Do NOT analyze unmodified files or unrelated parts of the codebase.
- If NONE of the files in the diff scope are relevant to your review category (e.g. a UX/design reviewer with no UI/CSS/component files in scope, a security reviewer with no auth/network code in scope, an a11y reviewer with no markup changes), respond IMMEDIATELY with a single short approval line such as "No relevant changes in scope — approved." and STOP. Do not start exploring the codebase.
- Your wall-clock budget is short. Spending it browsing unmodified files will cause this step to time out and block merge.${approvedContractBlock}`;
CRITICAL SCOPING RULES — read before doing anything else:
- The modified-file list is the starting point and primary reporting scope, not a prohibition on reading code required to validate the change.
- Read necessary callers, selectors, shared helpers, consumers, and tests outside that list when they establish production reachability, invariant coverage, or API/UI parity. Do not report unrelated pre-existing issues.
- If NONE of the modified files are relevant to your review category, confirm that from the list and fast-bail without broad repository exploration.
- Keep adjacent reads bounded to the changed behavior and its immediate production/test chain so the review finishes within its wall-clock budget.${approvedContractBlock}`;
const latestTaskForUserComments = await deps.store.getTask(task.id).catch(() => task);
const workflowStepUserComments = selectUserCommentsForAgentContext(latestTaskForUserComments, { limit: null });
@@ -391,7 +417,8 @@ export async function executeWorkflowStep(
- If you find an in-scope issue you can fix safely, edit the relevant files in this same session, run the smallest relevant verification, and then return APPROVE or APPROVE_WITH_NOTES.
- Return REVISE only when the issue is still present, cannot be safely fixed in this reviewer session, needs broader executor remediation, or needs user input.
- Plan Review may use fn_task_prompt_write to replace the task's PROMPT.md with the complete revised plan. Do not implement product code from Plan Review.
- Code Review and Browser Verification may fix implementation issues inside the assigned task worktree. Report each self-fixed issue as a finding with resolution resolved-in-review; list a fixed prior-lane finding in supersededFindingIds.`
- Code Review and Browser Verification may fix implementation issues inside the assigned task worktree. Report each self-fixed issue as a finding with resolution resolved-in-review; list a fixed prior-lane finding in supersededFindingIds.
- After any inline edit, treat your own change as untrusted: re-read the fresh diff, restart the mandatory review procedure from its requirements ledger and production-reachability checks, and rerun the smallest relevant verification. Never approve solely because the local fix compiles or its narrow test passes.`
: "";
const systemPrompt = `You are a workflow step agent executing: ${workflowStep.name}
@@ -429,6 +456,8 @@ export async function executeWorkflowStep(
deps.options.onAgentTool?.(taskId, toolName, detail);
},
});
// FNXC:CommandCenterActivity 2026-08-15-22:15: FN-8868 usage telemetry (restored post-wave-18).
attachAgentUsageTelemetry(agentLogger, { store: deps.store, agentId: task.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, lane: "executor" });
// Determine primary model and an explicit fallback. Review-type workflow
// steps use the validator lane; ordinary workflow prompts use the executor
@@ -458,6 +487,7 @@ export async function executeWorkflowStep(
const primaryModelId = useOverride ? workflowStep.modelId : laneModel.modelId;
// FNXC:ProviderAuth 2026-08-01-08:39: A workflow-step model override has no paired instance selection, so only the resolved primary task lane may carry its requested credential instance. Fallback attempts must retain their provider-default behavior rather than inheriting a primary-provider identity.
const primaryCredentialInstanceId = useOverride ? undefined : laneModel.credentialInstanceId;
attachAgentUsageTelemetry(agentLogger, { store: deps.store, agentId: task.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, model: primaryModelId ?? null, provider: primaryProvider ?? null, lane: "executor" });
const workflowFallback = isReviewTypeWorkflowStep
? resolveValidatorFallbackModel(settings)
@@ -658,6 +688,8 @@ export async function executeWorkflowStep(
...(additionalSkillPaths ? { additionalSkillPaths } : {}),
...(readonlyCustomTools.allowed.length > 0 ? { customTools: readonlyCustomTools.allowed } : {}),
});
// FNXC:CommandCenterActivity 2026-08-15-22:15: session boundary for the workflow-step runtime session (restored post-wave-18).
emitAgentSessionStart({ store: deps.store, agentId: task.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, model: primaryModelId ?? null, provider: primaryProvider ?? null, lane: "executor" });
const workflowModelDetails = formatModelMarkerDetails(
describeModel(session),

View File

@@ -75,6 +75,19 @@ export type HandleGraphFailureDeps = {
pausedAborted: Set<string>;
pausedAbortProvenance: Map<string, PausedAbortProvenance>;
userCanceledTaskIds: Set<string>;
/*
FNXC:WorkflowLifecycle 2026-08-15-22:15 (restored post-wave-18): liveness surfaces consumed by the
transient-resume scheduled retry's FIRE-TIME guard, so a retry armed before a pause/delete/move/park
cannot re-dispatch a task that is no longer in a safe WIP resume state.
*/
executing: Set<string>;
resumingUnpaused: Set<string>;
activeSessions: Map<string, unknown>;
activeStepExecutors: Map<string, unknown>;
activeWorkflowStepSessions: Map<string, unknown>;
activeCliTaskSessions: Map<string, unknown>;
activeWorkflowGraphAbortControllers: Map<string, AbortController>;
processWideGraphRouting: Set<string>;
getRunContextFor: (taskId: string) => EngineRunContext | undefined;
clearCompletedTaskWatchdog: (taskId: string) => void;
clearPausedAborted: (taskId: string) => void;
@@ -972,10 +985,46 @@ export async function handleGraphFailure(
status: null,
error: null,
}, deps.getRunContextFor(task.id));
/*
FNXC:WorkflowLifecycle 2026-08-15-22:15 (restored post-wave-18, FN-6782 family): the scheduled
retry must re-read the LIVE row at fire time and re-verify it is still in a safe WIP resume
state. The peel had replaced this with an unguarded `execute(live)` on the stale snapshot, so a
task deleted/paused/moved/parked/canceled — or one that already has an active run — between arming
and firing would be re-dispatched anyway.
*/
const scheduleRetry = () => {
deps.execute(live).catch((err: unknown) =>
executorLog.error(`Failed transient graph resume retry for ${task.id}:`, err),
);
void (async () => {
try {
const resumeTask = await deps.store.getTask(task.id);
const resumeFailureState = resumeTask as Task & { lastError?: unknown; failureReason?: unknown };
if (
resumeTask.deletedAt
|| resumeTask.paused
|| resumeTask.userPaused
|| deps.userCanceledTaskIds.has(task.id)
|| resumeTask.status != null
|| resumeTask.error != null
|| resumeFailureState.lastError != null
|| resumeFailureState.failureReason != null
|| resumeTask.column !== failureLanes.wip
|| (await resolveTerminalColumnsFor(deps.store, resumeTask.id)).includes(resumeTask.column)
|| deps.executing.has(task.id)
|| deps.activeSessions.has(task.id)
|| deps.activeStepExecutors.has(task.id)
|| deps.activeWorkflowStepSessions.has(task.id)
|| deps.activeCliTaskSessions.has(task.id)
|| deps.activeWorkflowGraphAbortControllers.has(task.id)
|| deps.resumingUnpaused.has(task.id)
|| deps.processWideGraphRouting.has(task.id)
) {
executorLog.debug(`${task.id}: skipping transient graph resume retry — task is no longer in a safe WIP resume state`);
return;
}
await deps.execute(resumeTask);
} catch (err) {
executorLog.error(`Failed transient graph resume retry for ${task.id}:`, err);
}
})();
};
if (TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS > 0) {
const handle = setTimeout(scheduleRetry, TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS);

View File

@@ -16,7 +16,7 @@
* clean completion handoffs.
*/
import type { Task, TaskDetail, TaskStore } from "@fusion/core";
import { isMergeRequestContractShadowEnabled } from "@fusion/core";
import { isMergeRequestContractShadowEnabled, resolveAgentActivityAttribution } from "@fusion/core";
import { ensureWorkflowCompletionSummary } from "../workflows/workflow-completion-summary.js";
import { executorLog } from "../logger.js";
import type { EngineRunContext } from "../util/run-audit.js";
@@ -65,5 +65,8 @@ export async function handoffTaskToReview(
});
}
// FNXC:AgentActivityStream 2026-08-09-09:09 (restored 2026-08-15-22:15 after wave-18 shell-ification dropped it):
// FN-8864 durable task:handed-off activity at the review-handoff choke point; monitoring never blocks handoff.
try { await deps.store.recordAgentActivity({ type: "task:handed-off", attributionClaim: resolveAgentActivityAttribution([{ id: agentId ?? task.assignedAgentId ?? "executor", provenance: agentId || task.assignedAgentId ? "roster" : "lane" }], "executor"), taskId: task.id, occurredAt: new Date().toISOString(), discriminator: `${runId ?? ""}:${reason}`, metadata: { runId, reason, source: "executor" } }); } catch { /* monitoring never blocks review handoff */ }
return handedOff;
}

View File

@@ -3,6 +3,10 @@
* Optional step revision attempt accounting peeled from executor.ts.
*/
import type { Task } from "@fusion/core";
import {
collectPlanReviewFeedbackHistory,
countPlanReviewRevisionAttempts,
} from "../plan-review-feedback-history.js";
export const OPTIONAL_STEP_REVISION_KEY_MARKER = "Workflow revision key:";
@@ -34,3 +38,37 @@ export function countOptionalStepRevisionAttempts(task: Pick<Task, "log">, key:
export function optionalStepRevisionLogOutcome(details: string, key: string): string {
return `${details}\n${OPTIONAL_STEP_REVISION_KEY_MARKER} ${key}`;
}
/*
FNXC:PlanReviewConvergence 2026-08-04-06:35 (FN-8768; restored 2026-08-15-22:15 after the wave-18
executor.ts shell-ification dropped it): Retry numbering uses the uncapped durable attempt ledger,
while prompt prose uses the separately bounded, deduplicated same-episode decision history.
*/
export function buildGraphPlanReviewConvergenceContext(
task: Pick<Task, "workflowStepResults">,
revisionKey: string,
): string {
const priorAttemptCount = countPlanReviewRevisionAttempts(task.workflowStepResults, { revisionKey });
const attempt = priorAttemptCount + 1;
if (attempt <= 1) return "";
const history = collectPlanReviewFeedbackHistory(task.workflowStepResults, { revisionKey });
const lines = [
`## Convergence — Plan Review attempt ${attempt}`,
"Treat the cumulative prior feedback below as a decision primer. Verify each prior blocker against the current PROMPT.md before looking for new findings.",
"- Do not re-raise a resolved or semantically duplicate blocker.",
"- A newly blocking finding must identify the revision that introduced it, the prior blocker that genuinely masked it, or why it is independently delivery-blocking for correctness, security, data safety, or executability. Record an earlier reviewer miss explicitly; never demote a critical defect merely because it was missed before.",
];
if (attempt >= 3) {
lines.push(
"- Severity ratchet (attempt 3+): only delivery-blocking critical defects may return REVISE; important/minor wording or implementation-detail findings are advisory.",
);
}
if (history.length > 0) {
lines.push("", "### Cumulative prior Plan Review ledger");
history.forEach((feedback, index) => {
lines.push(`#### PR${index + 1}`, feedback);
});
}
return lines.join("\n");
}

View File

@@ -21,6 +21,8 @@ export type {
// Re-export for backward compatibility (tests import from executor.ts)
export { summarizeToolArgs } from "../agents/agent-logger.js";
// FNXC:AgentActivityStream 2026-08-15-22:15: FN-8864 claim helper restored post-wave-18 (see workflow-gate-activity.ts).
export { resolveWorkflowGateActivityClaim } from "./workflow-gate-activity.js";
export {
createAgentCreateTool,
createAgentDeleteTool,

View File

@@ -57,6 +57,7 @@ import {
resolvePersistAgentThinkingLog,
resolveTaskLifecycleColumns,
resolveWorkflowIrForTask,
resolveAgentActivityAttribution,
serializeRetryStormError,
} from "@fusion/core";
import type { AgentSession } from "@earendil-works/pi-coding-agent";
@@ -170,6 +171,8 @@ import {
} from "./session-worktree-paths.js";
import { isWorkflowStepSkillDiscoverable, mergeAdditionalSkillPaths } from "./skill-path-helpers.js";
import { getExecutorSystemPrompt } from "./system-prompt.js";
// FNXC:CommandCenterActivity 2026-08-15-22:15: FN-8868 usage telemetry + session boundaries (restored post-wave-18).
import { attachAgentUsageTelemetry, emitAgentSessionStart } from "../agents/agent-usage-telemetry.js";
import { createConfiguredCommandAbortError, createSeenSteeringIds } from "./task-predicates.js";
import {
accumulateTokenUsage as accumulateTokenUsageImpl,
@@ -416,6 +419,9 @@ export async function runImplementation(
runId: syntheticRunId,
agentId: task.assignedAgentId ?? "executor",
});
// FNXC:AgentActivityStream 2026-08-09-09:09 (restored 2026-08-15-22:15 after wave-18 shell-ification dropped it):
// FN-8864 durable task:started activity at the implementation entry; monitoring never blocks execution.
try { await deps.store.recordAgentActivity({ type: "task:started", attributionClaim: resolveAgentActivityAttribution([{ id: task.assignedAgentId ?? "executor", provenance: task.assignedAgentId ? "roster" : "lane" }], "executor"), taskId: task.id, occurredAt: new Date().toISOString(), discriminator: syntheticRunId, metadata: { runId: syntheticRunId } }); } catch { /* monitoring never blocks execution */ }
// Build engine run context for audit instrumentation (FN-1404)
const engineRunContext: EngineRunContext = {
@@ -1989,6 +1995,9 @@ export async function runImplementation(
}
},
});
// FNXC:CommandCenterActivity 2026-08-09-15:06 (restored 2026-08-15-22:15 after the wave-18 peel dropped it):
// wire the usage-event store early so tool rows emitted before model resolution still land.
attachAgentUsageTelemetry(agentLogger, { store: deps.store, agentId: engineRunContext.agentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, lane: "executor" });
let agentRotationEvent: import("../credential-instance-rotation.js").RotationEvent | undefined;
let agentRotationDeclined = false;
@@ -2041,11 +2050,14 @@ export async function runImplementation(
// give the agent logger the context it needs to emit usage_events tool
// rows (KTD3). nodeId is sourced from the routed/effective node, null
// when the task has no node context.
agentLogger.setUsageContext({
attachAgentUsageTelemetry(agentLogger, {
store: deps.store,
model: executorModelId ?? null,
provider: executorProvider ?? null,
nodeId: detail.effectiveNodeId ?? detail.nodeId ?? null,
agentId: engineRunContext.agentId ?? null,
taskId: task.id,
lane: "executor",
});
// Determine whether we're resuming a previous session (pause/resume)
@@ -2165,6 +2177,14 @@ export async function runImplementation(
});
session = createdSession.session;
sessionFile = createdSession.sessionFile;
/*
FNXC:CommandCenterActivity 2026-08-09-15:06 (restored 2026-08-15-22:15 after the wave-18 peel dropped it):
Reopening a persisted executor session after pause continues one logical AgentSession.
Emit its session boundary only for a fresh manager so resumed work cannot inflate Sessions.
*/
if (!isResuming) {
emitAgentSessionStart({ store: deps.store, agentId: engineRunContext.agentId ?? null, taskId: task.id, nodeId: detail.effectiveNodeId ?? detail.nodeId ?? null, model: executorModelId ?? null, provider: executorProvider ?? null, lane: "executor" });
}
} catch (sessionStartError) {
if (await deps.recoverMissingWorktreeSessionStartFailure(task, worktreePath, sessionStartError, audit)) {
return;
@@ -2617,6 +2637,8 @@ export async function runImplementation(
taskId: task.id,
});
retrySession = createdRetrySession.session;
// FNXC:CommandCenterActivity 2026-08-09-15:18 (restored 2026-08-15-22:15): a retry builds a distinct runtime session, so it needs its own boundary only after construction succeeds.
emitAgentSessionStart({ store: deps.store, agentId: engineRunContext.agentId ?? null, taskId: task.id, nodeId: detail.effectiveNodeId ?? detail.nodeId ?? null, model: executorModelId ?? null, provider: executorProvider ?? null, lane: "executor" });
await deps.captureExecutorTokenUsageBaseline(task.id, retrySession);
captureSessionTokenBaseline(retrySession);
if (createdRetrySession.sessionFile) {

View File

@@ -282,12 +282,22 @@ Recommendation capture is disabled for this project (maxRecommendationsPerTask i
At the final accepted \`fn_task_done(outcome="completed")\` checkpoint, evaluate optional, non-blocking work discovered outside this task. Send at most ${maximum} task-ready recommendations, each with a stable unique \`id\`, \`title\`, \`description\`, and \`category\`, or explicitly send \`recommendations: []\` when none genuinely qualify. Example populated payload: \`recommendations: [{ id: "follow-up-export", title: "Add task export", description: "Provide a CSV export for completed tasks.", category: "feature" }]\`. Do not fabricate filler or include required current-task work, blockers, secrets, executable commands, reasoning, or duplicate ids. Recommendations are only for completed outcomes; never send them with \`outcome="blocked"\`. Use immediate task creation/delegation only for an explicit task requirement, necessary dependency coordination, or operator direction.`;
}
function getWithheldTaskCreationGuidance(taskCreateWithheld: boolean, delegateWithheld: boolean): string {
/*
FNXC:TaskRecommendations 2026-08-15-22:15:
Restored the recommendation-aware withheld-tool guidance (pre-peel executor.ts shape). The 2026-08-10
partial restore recovered this function from a pre-FN-8850 base, so a withheld session was still pointed
at fn_task_log instead of the completion recommendation route its validator accepts — and when capture is
disabled the prompt must say so rather than invite an unavailable write.
*/
function getWithheldTaskCreationGuidance(taskCreateWithheld: boolean, delegateWithheld: boolean, maximum: number): string {
if (!taskCreateWithheld && !delegateWithheld) return "";
const withheld = [
...(taskCreateWithheld ? ["`fn_task_create`"] : []),
...(delegateWithheld ? ["`fn_delegate_task`"] : []),
].join(" and ");
const recommendationRoute = maximum > 0
? `For optional, non-blocking discoveries, use the available completion recommendation route at accepted completion (or \`recommendations: []\` if none qualify).`
: "Recommendation capture is disabled, so retain non-blocking context in an honest task log or completion summary without inventing a follow-up.";
return `## Follow-up task creation is disabled for this session
This project's "Ephemeral agent follow-up tasks" policy withholds ${withheld}. ${
@@ -296,7 +306,7 @@ This project's "Ephemeral agent follow-up tasks" policy withholds ${withheld}. $
taskCreateWithheld && delegateWithheld ? "them" : "it"
}, and do not retry.
Ignore any instruction above that tells you to file follow-up work with ${withheld}. When you find out-of-scope work, record it instead with \`fn_task_log(message="follow-up: ...")\` and include it in your \`fn_task_done\` summary so the operator sees it. If the work genuinely blocks this task, use \`fn_task_done(outcome="blocked", reason="...")\` rather than trying to create a task for it.`;
Ignore any instruction above that tells you to file follow-up work with ${withheld}. ${recommendationRoute} If the work genuinely blocks this task, use \`fn_task_done(outcome="blocked", reason="...")\` rather than trying to create a task for it.`;
}
/** Resolve the executor system prompt from settings, falling back to the hardcoded constant. */
@@ -320,6 +330,7 @@ export function getExecutorSystemPrompt(
getWithheldTaskCreationGuidance(
toolAvailability?.taskCreateWithheld === true,
toolAvailability?.delegateWithheld === true,
maximumRecommendations,
),
].filter((section) => section.trim());
return sections.join("\n\n");

View File

@@ -63,6 +63,13 @@ export abstract class TaskExecutorState {
* `agent` is the admission-time row so session identity does not depend on a second getAgent round-trip.
*/
protected activeWorkflowPrincipals = new Map<string, { agentId: string; nodeInstanceId: string; agent?: Agent }>();
/**
* FNXC:AgentActivityStream 2026-08-09-13:59 (restored 2026-08-15-22:15 after wave-18 shell-ification dropped it):
* Node-scoped routed-principal retention (`taskId\0nodeId` -> agentId) that outlives release-principal,
* because the graph can emit a terminal gate result AFTER the per-attempt reservation is released while
* the activity outbox must still attribute that gate to the exact routed principal (FN-8864).
*/
protected workflowGateActivityPrincipals = new Map<string, string>();
protected executing = new Set<string>();
protected resumingUnpaused = new Set<string>();
protected approvalSuspended = new Set<string>();

View File

@@ -50,8 +50,8 @@ export abstract class TaskExecutorWorktreePureFacades extends TaskExecutorState
protected async recoverMissingWorktreeSessionStartFailure(...args: FacadeRestArgs<typeof impl.recoverMissingWorktreeSessionStartFailureImpl>): ReturnType<typeof impl.recoverMissingWorktreeSessionStartFailureImpl> { return impl.recoverMissingWorktreeSessionStartFailureImpl(bags.buildRecoverMissingWorktreeSessionStartFailureDeps(this), ...args); }
protected async emitWorktreeReanchoredAudit(...args: FacadeRestArgs<typeof impl.emitWorktreeReanchoredAuditImpl>): ReturnType<typeof impl.emitWorktreeReanchoredAuditImpl> { return impl.emitWorktreeReanchoredAuditImpl(bags.buildStoreRunContextDeps(this), ...args); }
listWorktreeHolders(): Array<{ taskId: string; worktreePath: string }> { return impl.listWorktreeHoldersImpl(this.activeWorktrees); }
protected async tryCreateWorktree(...args: FacadeRestArgs<typeof impl.tryCreateWorktreeImpl>): Promise<{ path: string; branch: string }> { return impl.tryCreateWorktreeImpl(bags.buildWorktreeCreateConflictFacadeDeps(this, constants.MAX_WORKTREE_RETRIES, bindHandleWorktreeConflict(this), bindTryCreateWorktree(this)), ...args); }
protected async handleWorktreeConflict(...args: FacadeRestArgs<typeof impl.handleWorktreeConflictImpl>): Promise<{ path: string; branch: string } | null> { return impl.handleWorktreeConflictImpl(bags.buildWorktreeCreateConflictFacadeDeps(this, constants.MAX_WORKTREE_RETRIES, bindHandleWorktreeConflict(this), bindTryCreateWorktree(this)), ...args); }
protected async tryCreateWorktree(...args: FacadeRestArgs<typeof impl.tryCreateWorktreeImpl>): Promise<{ path: string; branch: string }> { return impl.tryCreateWorktreeImpl(bags.buildWorktreeCreateConflictFacadeDeps(this, (this as { MAX_WORKTREE_RETRIES?: number }).MAX_WORKTREE_RETRIES ?? constants.MAX_WORKTREE_RETRIES, bindHandleWorktreeConflict(this), bindTryCreateWorktree(this)), ...args); }
protected async handleWorktreeConflict(...args: FacadeRestArgs<typeof impl.handleWorktreeConflictImpl>): Promise<{ path: string; branch: string } | null> { return impl.handleWorktreeConflictImpl(bags.buildWorktreeCreateConflictFacadeDeps(this, (this as { MAX_WORKTREE_RETRIES?: number }).MAX_WORKTREE_RETRIES ?? constants.MAX_WORKTREE_RETRIES, bindHandleWorktreeConflict(this), bindTryCreateWorktree(this)), ...args); }
protected async cleanupConflictingWorktree(...args: FacadeRestArgs<typeof impl.cleanupConflictingWorktreeImpl>): ReturnType<typeof impl.cleanupConflictingWorktreeImpl> { return impl.cleanupConflictingWorktreeImpl(bags.buildCleanupConflictingWorktreeDeps(this), ...args); }
protected async resolveWorktreeStartPoint(startPoint: string, taskId: string): ReturnType<typeof impl.resolveWorktreeStartPointImpl> { return impl.resolveWorktreeStartPointImpl(this.rootDir, this.store, startPoint, taskId); }
protected async squashImportDepIntoWorktree(...args: FacadeAfterFirst<typeof impl.squashImportDepIntoWorktreeImpl>): ReturnType<typeof impl.squashImportDepIntoWorktreeImpl> { return impl.squashImportDepIntoWorktreeImpl(this.store, ...args); }

View File

@@ -0,0 +1,22 @@
/**
* FNXC:AgentActivityStream 2026-08-15-22:15:
* FN-8864 gate-attribution claim helper, restored after the wave-18 executor.ts
* shell-ification (#3317) dropped it: FN-8864 landed hours before wave 18 the same
* day, and the squashed shell rewrite was built from a pre-FN-8864 base, silently
* deleting the executor's agent-activity writers. Keep this module the single home
* for the pure claim logic; the writers live at their lifecycle seams
* (run-implementation, handoff-task-to-review, execute-workflow-graph).
*
* FNXC:AgentActivityStream 2026-08-09-13:30:
* Workflow-gate activity must credit the routed node principal, because that route carries a
* reviewer override or column binding that task assignment alone cannot express. The outbox
* boundary still roster-proves this claim before it can become an org-map agent attribution.
*/
import { resolveAgentActivityAttribution } from "@fusion/core";
export function resolveWorkflowGateActivityClaim(routedPrincipalAgentId: string | undefined, assignedAgentId: string | undefined) {
const agentId = routedPrincipalAgentId ?? assignedAgentId ?? "executor";
return resolveAgentActivityAttribution([
{ id: agentId, provenance: routedPrincipalAgentId || assignedAgentId ? "roster" : "lane" },
], "executor");
}

View File

@@ -43,6 +43,8 @@ export type WorkflowPrincipalBeforeNodeDeps = {
workflowAgentCapacity: WorkflowAgentCapacity;
activeWorkflowAuthorities: Map<string, ActiveWorkflowAuthority>;
activeWorkflowPrincipals: Map<string, { agentId: string; nodeInstanceId: string; agent?: Agent }>;
/** FNXC:AgentActivityStream 2026-08-15-22:15: FN-8864 node-scoped routed-principal retention for gate attribution (restored post-wave-18). */
workflowGateActivityPrincipals: Map<string, string>;
workflowCapacityAttemptIds: Set<string>;
directWorkflowPrincipalWorkItemIds: Set<string>;
/** Holds written for principal unavailability so terminalization can skip re-closing them. */
@@ -353,6 +355,13 @@ deps.activeWorkflowPrincipals.set(nodeTask.id, {
agentId: routed.route.agent.id,
nodeInstanceId,
});
/*
* FNXC:AgentActivityStream 2026-08-09-13:59 (restored 2026-08-15-22:15 after wave-18 shell-ification dropped it):
* Keep this node-scoped record past release-principal. The graph can emit its terminal
* result after the per-attempt reservation is released, while the activity outbox must
* still attribute that gate to this exact routed principal (FN-8864).
*/
deps.workflowGateActivityPrincipals.set(`${nodeTask.id}\0${node.id}`, routed.route.agent.id);
if (durableWorkItemId) context["workflow:work-item-id"] = durableWorkItemId;
context["workflow:principal-agent-id"] = routed.route.agent.id;
context["workflow:principal-role"] = routed.route.role;