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:
7
.changeset/fix-wave18-peel-restorations.md
Normal file
7
.changeset/fix-wave18-peel-restorations.md
Normal 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.
|
||||
@@ -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 = {
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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);
|
||||
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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 }> }),
|
||||
|
||||
@@ -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(),
|
||||
|
||||
@@ -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 });
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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",
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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" },
|
||||
});
|
||||
|
||||
@@ -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 });
|
||||
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -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);
|
||||
}
|
||||
|
||||
@@ -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),
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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");
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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>();
|
||||
|
||||
@@ -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); }
|
||||
|
||||
22
packages/engine/src/executor/workflow-gate-activity.ts
Normal file
22
packages/engine/src/executor/workflow-gate-activity.ts
Normal 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");
|
||||
}
|
||||
@@ -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;
|
||||
|
||||
Reference in New Issue
Block a user