diff --git a/.changeset/fix-wave18-peel-restorations.md b/.changeset/fix-wave18-peel-restorations.md new file mode 100644 index 0000000000..46234ab583 --- /dev/null +++ b/.changeset/fix-wave18-peel-restorations.md @@ -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. diff --git a/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts b/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts index 431a95c69d..8c68fc1b18 100644 --- a/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts +++ b/packages/engine/src/__tests__/ce-workflow-step-executor.test.ts @@ -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 = { diff --git a/packages/engine/src/__tests__/executor-abort-provenance.test.ts b/packages/engine/src/__tests__/executor-abort-provenance.test.ts index c06a15743f..048dc365ae 100644 --- a/packages/engine/src/__tests__/executor-abort-provenance.test.ts +++ b/packages/engine/src/__tests__/executor-abort-provenance.test.ts @@ -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", diff --git a/packages/engine/src/__tests__/executor-base-commit-capture.real-git.test.ts b/packages/engine/src/__tests__/executor-base-commit-capture.real-git.test.ts index e20cd93049..bcb0d0397b 100644 --- a/packages/engine/src/__tests__/executor-base-commit-capture.real-git.test.ts +++ b/packages/engine/src/__tests__/executor-base-commit-capture.real-git.test.ts @@ -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); diff --git a/packages/engine/src/__tests__/executor-review-verdicts.test.ts b/packages/engine/src/__tests__/executor-review-verdicts.test.ts index e368683de4..c7dcb51440 100644 --- a/packages/engine/src/__tests__/executor-review-verdicts.test.ts +++ b/packages/engine/src/__tests__/executor-review-verdicts.test.ts @@ -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", diff --git a/packages/engine/src/__tests__/executor-step-session.test.ts b/packages/engine/src/__tests__/executor-step-session.test.ts index 106b19e184..a593a5f8b2 100644 --- a/packages/engine/src/__tests__/executor-step-session.test.ts +++ b/packages/engine/src/__tests__/executor-step-session.test.ts @@ -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 }> }), diff --git a/packages/engine/src/__tests__/executor-test-helpers.ts b/packages/engine/src/__tests__/executor-test-helpers.ts index 147d5097dc..79e0f0db75 100644 --- a/packages/engine/src/__tests__/executor-test-helpers.ts +++ b/packages/engine/src/__tests__/executor-test-helpers.ts @@ -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:` +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 & { + TaskExecutor: new (store: unknown, rootDir: string, options?: Record) => unknown; + }; + class DefaultRoutedTaskExecutor extends actual.TaskExecutor { + constructor(store: any, rootDir: string, options: Record = {}) { + 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(), diff --git a/packages/engine/src/__tests__/executor-workspace-capture.test.ts b/packages/engine/src/__tests__/executor-workspace-capture.test.ts index fa19b2f915..0cd34aec99 100644 --- a/packages/engine/src/__tests__/executor-workspace-capture.test.ts +++ b/packages/engine/src/__tests__/executor-workspace-capture.test.ts @@ -29,6 +29,9 @@ function createStore(overrides: Partial> = {}): 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 }); }); }); diff --git a/packages/engine/src/__tests__/heartbeat-executor.test.ts b/packages/engine/src/__tests__/heartbeat-executor.test.ts index 337c2409f7..d75321e511 100644 --- a/packages/engine/src/__tests__/heartbeat-executor.test.ts +++ b/packages/engine/src/__tests__/heartbeat-executor.test.ts @@ -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", diff --git a/packages/engine/src/__tests__/reliability-interactions/_helpers.ts b/packages/engine/src/__tests__/reliability-interactions/_helpers.ts index 34c1073418..223355b605 100644 --- a/packages/engine/src/__tests__/reliability-interactions/_helpers.ts +++ b/packages/engine/src/__tests__/reliability-interactions/_helpers.ts @@ -158,12 +158,28 @@ export async function createPgLayer(): Promise { 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, diff --git a/packages/engine/src/__tests__/reliability-interactions/graph-node-missing-worktree-recovery.test.ts b/packages/engine/src/__tests__/reliability-interactions/graph-node-missing-worktree-recovery.test.ts index d49ae383e8..94d4b7a809 100644 --- a/packages/engine/src/__tests__/reliability-interactions/graph-node-missing-worktree-recovery.test.ts +++ b/packages/engine/src/__tests__/reliability-interactions/graph-node-missing-worktree-recovery.test.ts @@ -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" }, }); diff --git a/packages/engine/src/__tests__/reliability-interactions/merge-node-paused-abort-retryable.test.ts b/packages/engine/src/__tests__/reliability-interactions/merge-node-paused-abort-retryable.test.ts index a4d1d6fe36..20218d7170 100644 --- a/packages/engine/src/__tests__/reliability-interactions/merge-node-paused-abort-retryable.test.ts +++ b/packages/engine/src/__tests__/reliability-interactions/merge-node-paused-abort-retryable.test.ts @@ -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 }); diff --git a/packages/engine/src/executor/attempt-executor-verification-fix.ts b/packages/engine/src/executor/attempt-executor-verification-fix.ts index b67b126fac..a7440e7449 100644 --- a/packages/engine/src/executor/attempt-executor-verification-fix.ts +++ b/packages/engine/src/executor/attempt-executor-verification-fix.ts @@ -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> | 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, diff --git a/packages/engine/src/executor/deps-bags.ts b/packages/engine/src/executor/deps-bags.ts index 7dc240e2e2..1ff9ab2a48 100644 --- a/packages/engine/src/executor/deps-bags.ts +++ b/packages/engine/src/executor/deps-bags.ts @@ -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, ...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, ); } diff --git a/packages/engine/src/executor/execute-workflow-graph.ts b/packages/engine/src/executor/execute-workflow-graph.ts index 1b487d6b71..3aa0d21af5 100644 --- a/packages/engine/src/executor/execute-workflow-graph.ts +++ b/packages/engine/src/executor/execute-workflow-graph.ts @@ -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; graphToolFailureRunCursors: Map; graphUnattendedRuns: Set; + /** FNXC:AgentActivityStream 2026-08-15-22:15: FN-8864 node-scoped routed-principal retention for gate attribution (restored post-wave-18). */ + workflowGateActivityPrincipals: Map; outerConcurrencyClaims: Set; processWideGraphRouting: Set; 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, + deps: Pick + & Partial>, taskId: string, result: CoreWorkflowStepResult, ): Promise { @@ -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); } diff --git a/packages/engine/src/executor/execute-workflow-step.ts b/packages/engine/src/executor/execute-workflow-step.ts index 21028ff4d5..d84c65a99e 100644 --- a/packages/engine/src/executor/execute-workflow-step.ts +++ b/packages/engine/src/executor/execute-workflow-step.ts @@ -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), diff --git a/packages/engine/src/executor/handle-graph-failure.ts b/packages/engine/src/executor/handle-graph-failure.ts index ed6bde5c05..c1aec50591 100644 --- a/packages/engine/src/executor/handle-graph-failure.ts +++ b/packages/engine/src/executor/handle-graph-failure.ts @@ -75,6 +75,19 @@ export type HandleGraphFailureDeps = { pausedAborted: Set; pausedAbortProvenance: Map; userCanceledTaskIds: Set; + /* + 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; + resumingUnpaused: Set; + activeSessions: Map; + activeStepExecutors: Map; + activeWorkflowStepSessions: Map; + activeCliTaskSessions: Map; + activeWorkflowGraphAbortControllers: Map; + processWideGraphRouting: Set; 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); diff --git a/packages/engine/src/executor/handoff-task-to-review.ts b/packages/engine/src/executor/handoff-task-to-review.ts index b715f42d69..d2732d95f8 100644 --- a/packages/engine/src/executor/handoff-task-to-review.ts +++ b/packages/engine/src/executor/handoff-task-to-review.ts @@ -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; } diff --git a/packages/engine/src/executor/optional-step-revision.ts b/packages/engine/src/executor/optional-step-revision.ts index be69a87015..a3d55456e0 100644 --- a/packages/engine/src/executor/optional-step-revision.ts +++ b/packages/engine/src/executor/optional-step-revision.ts @@ -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, 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, + 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"); +} diff --git a/packages/engine/src/executor/public-reexports.ts b/packages/engine/src/executor/public-reexports.ts index 13d83dc70f..fc94cb7705 100644 --- a/packages/engine/src/executor/public-reexports.ts +++ b/packages/engine/src/executor/public-reexports.ts @@ -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, diff --git a/packages/engine/src/executor/run-implementation.ts b/packages/engine/src/executor/run-implementation.ts index 7100c9f551..10e4326129 100644 --- a/packages/engine/src/executor/run-implementation.ts +++ b/packages/engine/src/executor/run-implementation.ts @@ -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) { diff --git a/packages/engine/src/executor/system-prompt.ts b/packages/engine/src/executor/system-prompt.ts index 7da5eb1048..eba4aa4e7b 100644 --- a/packages/engine/src/executor/system-prompt.ts +++ b/packages/engine/src/executor/system-prompt.ts @@ -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"); diff --git a/packages/engine/src/executor/task-executor-state.ts b/packages/engine/src/executor/task-executor-state.ts index 716affc1e5..3a1a68f4b6 100644 --- a/packages/engine/src/executor/task-executor-state.ts +++ b/packages/engine/src/executor/task-executor-state.ts @@ -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(); + /** + * 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(); protected executing = new Set(); protected resumingUnpaused = new Set(); protected approvalSuspended = new Set(); diff --git a/packages/engine/src/executor/task-executor-worktree-pure-facades.ts b/packages/engine/src/executor/task-executor-worktree-pure-facades.ts index b3be579bf1..cebdbeec8f 100644 --- a/packages/engine/src/executor/task-executor-worktree-pure-facades.ts +++ b/packages/engine/src/executor/task-executor-worktree-pure-facades.ts @@ -50,8 +50,8 @@ export abstract class TaskExecutorWorktreePureFacades extends TaskExecutorState protected async recoverMissingWorktreeSessionStartFailure(...args: FacadeRestArgs): ReturnType { return impl.recoverMissingWorktreeSessionStartFailureImpl(bags.buildRecoverMissingWorktreeSessionStartFailureDeps(this), ...args); } protected async emitWorktreeReanchoredAudit(...args: FacadeRestArgs): ReturnType { return impl.emitWorktreeReanchoredAuditImpl(bags.buildStoreRunContextDeps(this), ...args); } listWorktreeHolders(): Array<{ taskId: string; worktreePath: string }> { return impl.listWorktreeHoldersImpl(this.activeWorktrees); } - protected async tryCreateWorktree(...args: FacadeRestArgs): 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): 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): 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): 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): ReturnType { return impl.cleanupConflictingWorktreeImpl(bags.buildCleanupConflictingWorktreeDeps(this), ...args); } protected async resolveWorktreeStartPoint(startPoint: string, taskId: string): ReturnType { return impl.resolveWorktreeStartPointImpl(this.rootDir, this.store, startPoint, taskId); } protected async squashImportDepIntoWorktree(...args: FacadeAfterFirst): ReturnType { return impl.squashImportDepIntoWorktreeImpl(this.store, ...args); } diff --git a/packages/engine/src/executor/workflow-gate-activity.ts b/packages/engine/src/executor/workflow-gate-activity.ts new file mode 100644 index 0000000000..359ea0df84 --- /dev/null +++ b/packages/engine/src/executor/workflow-gate-activity.ts @@ -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"); +} diff --git a/packages/engine/src/executor/workflow-principal-before-node.ts b/packages/engine/src/executor/workflow-principal-before-node.ts index c36b79292f..c890c8caea 100644 --- a/packages/engine/src/executor/workflow-principal-before-node.ts +++ b/packages/engine/src/executor/workflow-principal-before-node.ts @@ -43,6 +43,8 @@ export type WorkflowPrincipalBeforeNodeDeps = { workflowAgentCapacity: WorkflowAgentCapacity; activeWorkflowAuthorities: Map; activeWorkflowPrincipals: Map; + /** FNXC:AgentActivityStream 2026-08-15-22:15: FN-8864 node-scoped routed-principal retention for gate attribution (restored post-wave-18). */ + workflowGateActivityPrincipals: Map; workflowCapacityAttemptIds: Set; directWorkflowPrincipalWorkItemIds: Set; /** 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;