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 { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||||
import { BUILTIN_WORKFLOWS, type WorkflowIr } from "@fusion/core";
|
import { BUILTIN_WORKFLOWS, type WorkflowIr } from "@fusion/core";
|
||||||
import "./executor-test-helpers.js";
|
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 { TaskExecutor } from "../executor.js";
|
||||||
import type { PluginRunner } from "../plugins/plugin-runner.js";
|
import type { PluginRunner } from "../plugins/plugin-runner.js";
|
||||||
import { WorkflowGraphExecutor } from "../workflows/workflow-graph-executor.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",
|
path: "/tmp/test/.worktrees/swift-falcon",
|
||||||
branch: "fusion/fn-ce-1",
|
branch: "fusion/fn-ce-1",
|
||||||
});
|
});
|
||||||
vi.spyOn(executor as any, "captureBaseCommitSha").mockResolvedValue(undefined);
|
|
||||||
|
|
||||||
const captured: { step?: any; worktreePath?: string } = {};
|
const captured: { step?: any; worktreePath?: string } = {};
|
||||||
vi.spyOn(executor as any, "executeWorkflowStep").mockImplementation(async (...args: any[]) => {
|
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",
|
path: "/tmp/test/.worktrees/fresh-ce-checkout",
|
||||||
branch: "fusion/fn-ce-1",
|
branch: "fusion/fn-ce-1",
|
||||||
});
|
});
|
||||||
vi.spyOn(executor as any, "captureBaseCommitSha").mockResolvedValue(undefined);
|
|
||||||
|
|
||||||
const captured: { worktreePath?: string } = {};
|
const captured: { worktreePath?: string } = {};
|
||||||
vi.spyOn(executor as any, "executeWorkflowStep").mockImplementation(async (...args: any[]) => {
|
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",
|
path: "/tmp/test/.worktrees/acquired-code-review",
|
||||||
branch: "fusion/fn-ce-1",
|
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 executeStep = vi.spyOn(executor as any, "executeWorkflowStep").mockResolvedValue({ success: true, output: "APPROVE" });
|
||||||
const requirements: any[] = [];
|
const requirements: any[] = [];
|
||||||
const codeReview = {
|
const codeReview = {
|
||||||
|
|||||||
@@ -21,7 +21,9 @@ accepted the old catch-all `hard-cancel`.
|
|||||||
*/
|
*/
|
||||||
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
import { afterEach, beforeEach, describe, expect, it, vi } from "vitest";
|
||||||
import "./executor-test-helpers.js";
|
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 { createMockStore, resetExecutorMocks } from "./executor-test-helpers.js";
|
||||||
import type { TaskDetail } from "@fusion/core";
|
import type { TaskDetail } from "@fusion/core";
|
||||||
|
|
||||||
@@ -250,7 +252,7 @@ describe("pause-abort provenance truthfulness (KB-PROV)", () => {
|
|||||||
column: "in-review",
|
column: "in-review",
|
||||||
steps: [{ name: "Implement", status: "done" }],
|
steps: [{ name: "Implement", status: "done" }],
|
||||||
});
|
});
|
||||||
const benign = (executor as any).isBenignInReviewPauseAbort(
|
const benign = isBenignInReviewPauseAbort(
|
||||||
live,
|
live,
|
||||||
{ disposition: "failed", outcome: "failure", visitedNodeIds: ["plan", "execute"], context: {} },
|
{ disposition: "failed", outcome: "failure", visitedNodeIds: ["plan", "execute"], context: {} },
|
||||||
provenance,
|
provenance,
|
||||||
@@ -288,7 +290,7 @@ describe("pause-abort provenance truthfulness (KB-PROV)", () => {
|
|||||||
const { executor } = makeExecutor();
|
const { executor } = makeExecutor();
|
||||||
const result = { disposition: "failed", outcome: "failure", visitedNodeIds: ["plan", "execute"], context: {} };
|
const result = { disposition: "failed", outcome: "failure", visitedNodeIds: ["plan", "execute"], context: {} };
|
||||||
const classify = (column: string, reviewLane: string) =>
|
const classify = (column: string, reviewLane: string) =>
|
||||||
(executor as any).isBenignInReviewPauseAbort(
|
isBenignInReviewPauseAbort(
|
||||||
makeTask({ column, steps: [{ name: "Implement", status: "done" }] }),
|
makeTask({ column, steps: [{ name: "Implement", status: "done" }] }),
|
||||||
result,
|
result,
|
||||||
"engine-abort",
|
"engine-abort",
|
||||||
|
|||||||
@@ -5,7 +5,9 @@ import os from "node:os";
|
|||||||
import path from "node:path";
|
import path from "node:path";
|
||||||
import { EventEmitter } from "node:events";
|
import { EventEmitter } from "node:events";
|
||||||
import type { Task, TaskStore } from "@fusion/core";
|
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 hasGit = spawnSync("git", ["--version"], { stdio: "pipe" }).status === 0;
|
||||||
const describeIfGit = hasGit ? describe : describe.skip;
|
const describeIfGit = hasGit ? describe : describe.skip;
|
||||||
@@ -62,10 +64,9 @@ describeIfGit("captureBaseCommitSha (real git)", () => {
|
|||||||
}
|
}
|
||||||
|
|
||||||
const store = createStore();
|
const store = createStore();
|
||||||
const executor = new TaskExecutor(store, repo);
|
|
||||||
const audit = { git: vi.fn().mockResolvedValue(undefined) };
|
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;
|
const firstBase = (store.updateTask as any).mock.calls[0][1].baseCommitSha as string;
|
||||||
expect(firstBase).toBeTruthy();
|
expect(firstBase).toBeTruthy();
|
||||||
|
|
||||||
@@ -74,7 +75,7 @@ describeIfGit("captureBaseCommitSha (real git)", () => {
|
|||||||
|
|
||||||
// Resume of the same task: baseCommitSha must be preserved so diff math
|
// Resume of the same task: baseCommitSha must be preserved so diff math
|
||||||
// stays stable across sessions (FN-4309/FN-4383).
|
// 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);
|
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 () => {
|
it("suspends an unrouted graph run before opening an implementation session", async () => {
|
||||||
const store = createMockStore();
|
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({
|
await executor.execute({
|
||||||
id: "FN-routing", title: "Routing fixture", description: "", column: "in-progress",
|
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);
|
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(
|
expect(() => selectImplementationSessionCall(
|
||||||
mockedCreateFnAgent.mock.calls.map(([options]) => options as { customTools?: Array<{ name?: string }> }),
|
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 { installTaskWorktreeIdentityGuard } from "../worktree/worktree-hooks.js";
|
||||||
import type * as ReviewerModule from "../execution/reviewer.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
|
// Mock external dependencies
|
||||||
vi.mock("../pi.js", () => ({
|
vi.mock("../pi.js", () => ({
|
||||||
createFnAgent: vi.fn(),
|
createFnAgent: vi.fn(),
|
||||||
|
|||||||
@@ -29,6 +29,9 @@ function createStore(overrides: Partial<Record<string, unknown>> = {}): TaskStor
|
|||||||
logEntry: vi.fn().mockResolvedValue(undefined),
|
logEntry: vi.fn().mockResolvedValue(undefined),
|
||||||
getSettings: vi.fn().mockResolvedValue({ autoMerge: false }),
|
getSettings: vi.fn().mockResolvedValue({ autoMerge: false }),
|
||||||
getRunContextFor: vi.fn(),
|
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),
|
on: emitter.on.bind(emitter),
|
||||||
...overrides,
|
...overrides,
|
||||||
}) as unknown as TaskStore & EventEmitter;
|
}) as unknown as TaskStore & EventEmitter;
|
||||||
@@ -209,11 +212,27 @@ describeIfGit("U1 KTD2 — verifyWorktreeInvariants iterates per worktree, prese
|
|||||||
expect(result.expected).toBe(BRANCH);
|
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();
|
fx = await createWorkspaceFixture();
|
||||||
const executor = workspaceExecutor(fx);
|
const executor = workspaceExecutor(fx);
|
||||||
const task = makeTask(TASK_ID, { branch: BRANCH, workspaceWorktrees: {} });
|
const task = makeTask(TASK_ID, { branch: BRANCH, workspaceWorktrees: {} });
|
||||||
const result = await (executor as any).verifyWorktreeInvariants(task);
|
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 });
|
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,
|
// 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.
|
// 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([
|
expect(callArgs.customTools!.map((tool) => tool.name)).toEqual([
|
||||||
"fn_task_create",
|
"fn_task_create",
|
||||||
"fn_task_log",
|
"fn_task_log",
|
||||||
@@ -3411,6 +3412,7 @@ describe("executeHeartbeat", () => {
|
|||||||
"fn_mission_update",
|
"fn_mission_update",
|
||||||
"fn_mission_set_status",
|
"fn_mission_set_status",
|
||||||
"fn_mission_delete",
|
"fn_mission_delete",
|
||||||
|
"fn_mission_reconcile",
|
||||||
"fn_milestone_add",
|
"fn_milestone_add",
|
||||||
"fn_milestone_update",
|
"fn_milestone_update",
|
||||||
"fn_milestone_delete",
|
"fn_milestone_delete",
|
||||||
@@ -3419,6 +3421,7 @@ describe("executeHeartbeat", () => {
|
|||||||
"fn_slice_delete",
|
"fn_slice_delete",
|
||||||
"fn_feature_add",
|
"fn_feature_add",
|
||||||
"fn_feature_update",
|
"fn_feature_update",
|
||||||
|
"fn_feature_repair_validation",
|
||||||
"fn_feature_set_status",
|
"fn_feature_set_status",
|
||||||
"fn_feature_delete",
|
"fn_feature_delete",
|
||||||
"fn_feature_link_task",
|
"fn_feature_link_task",
|
||||||
|
|||||||
@@ -158,12 +158,28 @@ export async function createPgLayer(): Promise<PgLayerFixture> {
|
|||||||
runtimeUrl: testUrl,
|
runtimeUrl: testUrl,
|
||||||
migrationUrl: testUrl,
|
migrationUrl: testUrl,
|
||||||
migrationUrlOverridden: false,
|
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 });
|
const schemaConn = await createConnectionSetFromUrl(backend, { poolMax: 1, connectTimeoutSeconds: 5 });
|
||||||
await applySchemaBaseline(schemaConn.migration);
|
await applySchemaBaseline(schemaConn.migration);
|
||||||
await schemaConn.close();
|
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 {
|
return {
|
||||||
layer,
|
layer,
|
||||||
dbName,
|
dbName,
|
||||||
|
|||||||
@@ -2,7 +2,8 @@ import { describe, expect, it, vi, beforeEach } from "vitest";
|
|||||||
import type { TaskDetail } from "@fusion/core";
|
import type { TaskDetail } from "@fusion/core";
|
||||||
import "../executor-test-helpers.js";
|
import "../executor-test-helpers.js";
|
||||||
import { PLAN_REVIEW_PROVIDER_FAILURE_HOLD_VALUE } from "../../workflows/workflow-graph-executor.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 { activeSessionRegistry } from "../../agents/active-session-registry.js";
|
||||||
import {
|
import {
|
||||||
createMockStore,
|
createMockStore,
|
||||||
@@ -80,8 +81,7 @@ describe("graphFailureValue optional-group materialized ids", () => {
|
|||||||
});
|
});
|
||||||
|
|
||||||
it("prefers the group's published value for a `group::template` failed node", () => {
|
it("prefers the group's published value for a `group::template` failed node", () => {
|
||||||
const executor = new TaskExecutor(createMockStore(), "/tmp/test");
|
const value = graphFailureValue({
|
||||||
const value = (executor as any).graphFailureValue({
|
|
||||||
visitedNodeIds: ["plan-review", "plan-review::plan-review-step"],
|
visitedNodeIds: ["plan-review", "plan-review::plan-review-step"],
|
||||||
context: {
|
context: {
|
||||||
"node:plan-review:value": PLAN_REVIEW_PROVIDER_FAILURE_HOLD_VALUE,
|
"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", () => {
|
it("falls back to the unqualified template value when the group has none", () => {
|
||||||
const executor = new TaskExecutor(createMockStore(), "/tmp/test");
|
const value = graphFailureValue({
|
||||||
const value = (executor as any).graphFailureValue({
|
|
||||||
visitedNodeIds: ["plan-review::plan-review-step"],
|
visitedNodeIds: ["plan-review::plan-review-step"],
|
||||||
context: { "node:plan-review-step:value": "exception" },
|
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", () => {
|
it("keeps resolving foreach `#` instance ids through the container key", () => {
|
||||||
const executor = new TaskExecutor(createMockStore(), "/tmp/test");
|
const value = graphFailureValue({
|
||||||
const value = (executor as any).graphFailureValue({
|
|
||||||
visitedNodeIds: ["steps#0:step-execute"],
|
visitedNodeIds: ["steps#0:step-execute"],
|
||||||
context: { "node:steps:value": "awaiting-user-input" },
|
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 () => {
|
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({
|
const { store, task, executor, mergeRequester } = makeHarness({
|
||||||
autoMerge: undefined,
|
autoMerge: true,
|
||||||
branchContext: { groupId: "BG-6735", source: "mission", assignmentMode: "shared" },
|
branchContext: { groupId: "BG-6735", source: "mission", assignmentMode: "shared" },
|
||||||
}, { autoMerge: false });
|
}, { autoMerge: false });
|
||||||
|
|
||||||
|
|||||||
@@ -15,6 +15,8 @@
|
|||||||
import type { Settings, Task, TaskStore } from "@fusion/core";
|
import type { Settings, Task, TaskStore } from "@fusion/core";
|
||||||
import { resolveExecutorFallbackModel, resolvePersistAgentThinkingLog } from "@fusion/core";
|
import { resolveExecutorFallbackModel, resolvePersistAgentThinkingLog } from "@fusion/core";
|
||||||
import { AgentLogger } from "../agents/agent-logger.js";
|
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 {
|
import {
|
||||||
createResolvedAgentSession,
|
createResolvedAgentSession,
|
||||||
resolveExecutorSessionModel,
|
resolveExecutorSessionModel,
|
||||||
@@ -82,6 +84,8 @@ export async function attemptExecutorVerificationFix(
|
|||||||
onAgentText: deps.onAgentText,
|
onAgentText: deps.onAgentText,
|
||||||
onAgentTool: deps.onAgentTool,
|
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
|
// Build skill selection context
|
||||||
let skillContext: Awaited<ReturnType<typeof buildSessionSkillContext>> | undefined;
|
let skillContext: Awaited<ReturnType<typeof buildSessionSkillContext>> | undefined;
|
||||||
@@ -109,6 +113,7 @@ export async function attemptExecutorVerificationFix(
|
|||||||
task.credentialInstanceId,
|
task.credentialInstanceId,
|
||||||
);
|
);
|
||||||
const { provider: executorProvider, modelId: executorModelId } = executorSessionModel;
|
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);
|
const executorFallback = resolveExecutorFallbackModel(settings);
|
||||||
|
|
||||||
@@ -157,6 +162,8 @@ Do not refactor, rename broadly, or make opportunistic improvements.
|
|||||||
...(skillContext?.skillSelectionContext ? { skillSelection: skillContext.skillSelectionContext } : {}),
|
...(skillContext?.skillSelectionContext ? { skillSelection: skillContext.skillSelectionContext } : {}),
|
||||||
...(skillContext && skillContext.additionalSkillPaths.length > 0 ? { additionalSkillPaths: skillContext.additionalSkillPaths } : {}),
|
...(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(
|
await deps.store.logEntry(
|
||||||
task.id,
|
task.id,
|
||||||
|
|||||||
@@ -144,6 +144,8 @@ export function buildExecuteWorkflowGraphDeps(host: any): any {
|
|||||||
"graphRethinkNarrations", "graphRouting", "graphSeamGoverningNodeId", "graphSeamSkillName",
|
"graphRethinkNarrations", "graphRouting", "graphSeamGoverningNodeId", "graphSeamSkillName",
|
||||||
"graphSeamThinkingLevel", "graphStepActiveContext", "graphStepRunOnce", "graphStepSessionPinned",
|
"graphSeamThinkingLevel", "graphStepActiveContext", "graphStepRunOnce", "graphStepSessionPinned",
|
||||||
"graphToolFailureRunCursors", "graphUnattendedRuns", "outerConcurrencyClaims",
|
"graphToolFailureRunCursors", "graphUnattendedRuns", "outerConcurrencyClaims",
|
||||||
|
// FNXC:AgentActivityStream 2026-08-15-22:15: FN-8864 gate-attribution retention map (restored post-wave-18).
|
||||||
|
"workflowGateActivityPrincipals",
|
||||||
]),
|
]),
|
||||||
...facadeMethods(host, [
|
...facadeMethods(host, [
|
||||||
"getRunContextFor", "advanceNoMergeWorkflowToCompleteColumn", "applyGraphRethinkReset",
|
"getRunContextFor", "advanceNoMergeWorkflowToCompleteColumn", "applyGraphRethinkReset",
|
||||||
@@ -163,9 +165,14 @@ export function buildHandleGraphFailureDeps(host: any): any {
|
|||||||
store: host.store,
|
store: host.store,
|
||||||
rootDir: host.rootDir,
|
rootDir: host.rootDir,
|
||||||
options: host.options as { stuckTaskDetector?: { untrackTask?: (taskId: string) => void }; [k: string]: unknown },
|
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, [
|
...facadeFields(host, [
|
||||||
"activeWorktrees", "completionFinalizedTaskIds", "graphExecuteSelfRequeued",
|
"activeWorktrees", "completionFinalizedTaskIds", "graphExecuteSelfRequeued",
|
||||||
"graphToolFailureRunCursors", "pausedAborted", "pausedAbortProvenance", "userCanceledTaskIds",
|
"graphToolFailureRunCursors", "pausedAborted", "pausedAbortProvenance", "userCanceledTaskIds",
|
||||||
|
"executing", "resumingUnpaused", "activeSessions", "activeStepExecutors",
|
||||||
|
"activeWorkflowStepSessions", "activeCliTaskSessions", "activeWorkflowGraphAbortControllers",
|
||||||
]),
|
]),
|
||||||
...facadeMethods(host, [
|
...facadeMethods(host, [
|
||||||
"getRunContextFor", "clearCompletedTaskWatchdog", "clearPausedAborted", "execute",
|
"getRunContextFor", "clearCompletedTaskWatchdog", "clearPausedAborted", "execute",
|
||||||
@@ -1343,7 +1350,12 @@ export function buildRouteResetParsePinMismatchToRetryDeps(host: any): any {
|
|||||||
export function buildCreateWorktreeFacadeDeps(host: any, tryCreateWorktree: any): any {
|
export function buildCreateWorktreeFacadeDeps(host: any, tryCreateWorktree: any): any {
|
||||||
return buildCreateWorktreeDeps(
|
return buildCreateWorktreeDeps(
|
||||||
host,
|
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,
|
tryCreateWorktree,
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -28,7 +28,9 @@ import {
|
|||||||
resolveWorkflowIrForTask,
|
resolveWorkflowIrForTask,
|
||||||
upsertWorkflowStepResult,
|
upsertWorkflowStepResult,
|
||||||
applySupersededFindingIds,
|
applySupersededFindingIds,
|
||||||
|
isTerminalStepResult,
|
||||||
} from "@fusion/core";
|
} from "@fusion/core";
|
||||||
|
import { resolveWorkflowGateActivityClaim } from "./workflow-gate-activity.js";
|
||||||
import type { ImplementationExit } from "./implementation-exit.js";
|
import type { ImplementationExit } from "./implementation-exit.js";
|
||||||
import type { WorkflowGraphTaskRunResult } from "../workflows/workflow-graph-task-runner.js";
|
import type { WorkflowGraphTaskRunResult } from "../workflows/workflow-graph-task-runner.js";
|
||||||
import { WorkflowGraphTaskRunner } 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>;
|
graphStepSessionPinned: Set<string>;
|
||||||
graphToolFailureRunCursors: Map<string, number>;
|
graphToolFailureRunCursors: Map<string, number>;
|
||||||
graphUnattendedRuns: Set<string>;
|
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>;
|
outerConcurrencyClaims: Set<string>;
|
||||||
processWideGraphRouting: Set<string>;
|
processWideGraphRouting: Set<string>;
|
||||||
getRunContextFor: (taskId: string) => EngineRunContext | undefined;
|
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.
|
* the same write. Exported for production-shaped graph-writer tests.
|
||||||
*/
|
*/
|
||||||
export async function persistWorkflowStepResult(
|
export async function persistWorkflowStepResult(
|
||||||
deps: Pick<ExecuteWorkflowGraphDeps, "store" | "getRunContextFor" | "readTaskArtifact">,
|
deps: Pick<ExecuteWorkflowGraphDeps, "store" | "getRunContextFor" | "readTaskArtifact">
|
||||||
|
& Partial<Pick<ExecuteWorkflowGraphDeps, "workflowGateActivityPrincipals" | "activeWorkflowPrincipals">>,
|
||||||
taskId: string,
|
taskId: string,
|
||||||
result: CoreWorkflowStepResult,
|
result: CoreWorkflowStepResult,
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
@@ -191,8 +196,76 @@ export async function persistWorkflowStepResult(
|
|||||||
} else {
|
} else {
|
||||||
await deps.store.updateTask(taskId, { workflowStepResults: existing }, deps.getRunContextFor(taskId));
|
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,
|
workflowAgentCapacity: deps.workflowAgentCapacity,
|
||||||
activeWorkflowAuthorities: deps.activeWorkflowAuthorities,
|
activeWorkflowAuthorities: deps.activeWorkflowAuthorities,
|
||||||
activeWorkflowPrincipals: deps.activeWorkflowPrincipals,
|
activeWorkflowPrincipals: deps.activeWorkflowPrincipals,
|
||||||
|
workflowGateActivityPrincipals: deps.workflowGateActivityPrincipals,
|
||||||
workflowCapacityAttemptIds,
|
workflowCapacityAttemptIds,
|
||||||
directWorkflowPrincipalWorkItemIds,
|
directWorkflowPrincipalWorkItemIds,
|
||||||
directWorkflowPrincipalHeldWorkItemIds,
|
directWorkflowPrincipalHeldWorkItemIds,
|
||||||
@@ -797,6 +871,10 @@ export async function executeWorkflowGraph(
|
|||||||
}
|
}
|
||||||
deps.activeWorkflowAuthorities.delete(task.id);
|
deps.activeWorkflowAuthorities.delete(task.id);
|
||||||
deps.activeWorkflowPrincipals.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) {
|
if (graphAbortController && deps.activeWorkflowGraphAbortControllers.get(task.id) === graphAbortController) {
|
||||||
deps.activeWorkflowGraphAbortControllers.delete(task.id);
|
deps.activeWorkflowGraphAbortControllers.delete(task.id);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ import {
|
|||||||
applyReviewSeverityGate,
|
applyReviewSeverityGate,
|
||||||
isOpenWorkflowReviewFinding,
|
isOpenWorkflowReviewFinding,
|
||||||
MAX_WORKFLOW_REVIEW_FINDINGS,
|
MAX_WORKFLOW_REVIEW_FINDINGS,
|
||||||
|
PLAN_REVIEW_GROUP_ID,
|
||||||
finalizePlanningSegment,
|
finalizePlanningSegment,
|
||||||
resolveExecutorFallbackModel,
|
resolveExecutorFallbackModel,
|
||||||
resolvePersistAgentThinkingLog,
|
resolvePersistAgentThinkingLog,
|
||||||
@@ -81,6 +82,10 @@ import {
|
|||||||
type WorkflowStepOutcome,
|
type WorkflowStepOutcome,
|
||||||
} from "./workflow-step-verdict.js";
|
} from "./workflow-step-verdict.js";
|
||||||
import { resolveDiffBaseRef } from "./worktree-git-refs.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);
|
const execAsync = promisify(exec);
|
||||||
|
|
||||||
@@ -127,7 +132,6 @@ export async function executeWorkflowStep(
|
|||||||
// assumptions and proceed instead of parking on a question. Explicit opt-in
|
// assumptions and proceed instead of parking on a question. Explicit opt-in
|
||||||
// only (default false = board run); see runGraphCustomNode / KTD-3.
|
// only (default false = board run); see runGraphCustomNode / KTD-3.
|
||||||
const unattended = stepOptions?.unattended === true;
|
const unattended = stepOptions?.unattended === true;
|
||||||
const isPlanReviewStep = workflowStep.id === "graph:plan-review-step" || workflowStep.name === "Plan Review";
|
|
||||||
/*
|
/*
|
||||||
FNXC:WorkflowReviewFindings 2026-08-05-06:29:
|
FNXC:WorkflowReviewFindings 2026-08-05-06:29:
|
||||||
reviewKind is carried from graph synthesis (cfg.reviewKind / optional-group context) so prompt
|
reviewKind is carried from graph synthesis (cfg.reviewKind / optional-group context) so prompt
|
||||||
@@ -141,6 +145,15 @@ export async function executeWorkflowStep(
|
|||||||
requireExternalIntegrationEvidence?: boolean;
|
requireExternalIntegrationEvidence?: boolean;
|
||||||
};
|
};
|
||||||
const optionalGroupId = workflowStepMetadata.optionalGroupId;
|
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 =
|
const isReviewTypeWorkflowStep =
|
||||||
isPlanReviewStep
|
isPlanReviewStep
|
||||||
|| workflowStepMetadata.reviewCanFixInline === true
|
|| workflowStepMetadata.reviewCanFixInline === true
|
||||||
@@ -182,6 +195,11 @@ export async function executeWorkflowStep(
|
|||||||
}
|
}
|
||||||
const workflowReviewSpecText = typeof workflowReviewSpecArtifact === "string" ? workflowReviewSpecArtifact : "";
|
const workflowReviewSpecText = typeof workflowReviewSpecArtifact === "string" ? workflowReviewSpecArtifact : "";
|
||||||
const planReviewSpecText = isPlanReviewStep ? workflowReviewSpecText : "";
|
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:
|
FNXC:PlanReview 2026-07-21-16:30:
|
||||||
@@ -265,33 +283,41 @@ export async function executeWorkflowStep(
|
|||||||
const approvedContractBlock = isReviewTypeWorkflowStep && !isPlanReviewStep
|
const approvedContractBlock = isReviewTypeWorkflowStep && !isPlanReviewStep
|
||||||
? `
|
? `
|
||||||
|
|
||||||
Approved Task Contract:
|
Approved Task Contract:
|
||||||
- PROMPT.md is the authoritative current contract for this review. It includes any approved planning revisions and scope decisions.
|
- 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.
|
- 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.
|
- 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.
|
- Scope exclusions do not waive security, correctness, or data-integrity defects in the approved implementation.
|
||||||
|
|
||||||
--- BEGIN APPROVED PROMPT.md ---
|
--- BEGIN APPROVED PROMPT.md ---
|
||||||
${workflowReviewSpecText}
|
${workflowReviewSpecText}
|
||||||
--- END APPROVED PROMPT.md ---`
|
--- 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
|
const scopeBlock = isPlanReviewStep
|
||||||
? `Plan Review Scope:
|
? `Plan Review Scope:
|
||||||
- Review the task plan artifact (PROMPT.md), reproduced verbatim below, and task metadata only.
|
- 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.
|
- 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.
|
- 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.
|
- If the plan is internally consistent, complete, scoped, and verifiable, approve even when the worktree contains unrelated changes from another task.
|
||||||
|
|
||||||
--- BEGIN PROMPT.md ---
|
--- BEGIN PROMPT.md ---
|
||||||
${planReviewSpecText}
|
${planReviewSpecText}
|
||||||
--- END PROMPT.md ---`
|
--- END PROMPT.md ---${planReviewConvergenceContext ? `\n\n${planReviewConvergenceContext}` : ""}`
|
||||||
: `Diff Scope (files changed by THIS task vs base):
|
: `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:
|
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.
|
- The modified-file list is the starting point and primary reporting scope, not a prohibition on reading code required to validate the change.
|
||||||
- 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.
|
- 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.
|
||||||
- Your wall-clock budget is short. Spending it browsing unmodified files will cause this step to time out and block merge.${approvedContractBlock}`;
|
- 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 latestTaskForUserComments = await deps.store.getTask(task.id).catch(() => task);
|
||||||
const workflowStepUserComments = selectUserCommentsForAgentContext(latestTaskForUserComments, { limit: null });
|
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.
|
- 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.
|
- 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.
|
- 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}
|
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);
|
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
|
// Determine primary model and an explicit fallback. Review-type workflow
|
||||||
// steps use the validator lane; ordinary workflow prompts use the executor
|
// 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;
|
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.
|
// 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;
|
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
|
const workflowFallback = isReviewTypeWorkflowStep
|
||||||
? resolveValidatorFallbackModel(settings)
|
? resolveValidatorFallbackModel(settings)
|
||||||
@@ -658,6 +688,8 @@ export async function executeWorkflowStep(
|
|||||||
...(additionalSkillPaths ? { additionalSkillPaths } : {}),
|
...(additionalSkillPaths ? { additionalSkillPaths } : {}),
|
||||||
...(readonlyCustomTools.allowed.length > 0 ? { customTools: readonlyCustomTools.allowed } : {}),
|
...(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(
|
const workflowModelDetails = formatModelMarkerDetails(
|
||||||
describeModel(session),
|
describeModel(session),
|
||||||
|
|||||||
@@ -75,6 +75,19 @@ export type HandleGraphFailureDeps = {
|
|||||||
pausedAborted: Set<string>;
|
pausedAborted: Set<string>;
|
||||||
pausedAbortProvenance: Map<string, PausedAbortProvenance>;
|
pausedAbortProvenance: Map<string, PausedAbortProvenance>;
|
||||||
userCanceledTaskIds: Set<string>;
|
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;
|
getRunContextFor: (taskId: string) => EngineRunContext | undefined;
|
||||||
clearCompletedTaskWatchdog: (taskId: string) => void;
|
clearCompletedTaskWatchdog: (taskId: string) => void;
|
||||||
clearPausedAborted: (taskId: string) => void;
|
clearPausedAborted: (taskId: string) => void;
|
||||||
@@ -972,10 +985,46 @@ export async function handleGraphFailure(
|
|||||||
status: null,
|
status: null,
|
||||||
error: null,
|
error: null,
|
||||||
}, deps.getRunContextFor(task.id));
|
}, 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 = () => {
|
const scheduleRetry = () => {
|
||||||
deps.execute(live).catch((err: unknown) =>
|
void (async () => {
|
||||||
executorLog.error(`Failed transient graph resume retry for ${task.id}:`, err),
|
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) {
|
if (TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS > 0) {
|
||||||
const handle = setTimeout(scheduleRetry, TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS);
|
const handle = setTimeout(scheduleRetry, TRANSIENT_GRAPH_RESUME_RETRY_BACKOFF_MS);
|
||||||
|
|||||||
@@ -16,7 +16,7 @@
|
|||||||
* clean completion handoffs.
|
* clean completion handoffs.
|
||||||
*/
|
*/
|
||||||
import type { Task, TaskDetail, TaskStore } from "@fusion/core";
|
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 { ensureWorkflowCompletionSummary } from "../workflows/workflow-completion-summary.js";
|
||||||
import { executorLog } from "../logger.js";
|
import { executorLog } from "../logger.js";
|
||||||
import type { EngineRunContext } from "../util/run-audit.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;
|
return handedOff;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,6 +3,10 @@
|
|||||||
* Optional step revision attempt accounting peeled from executor.ts.
|
* Optional step revision attempt accounting peeled from executor.ts.
|
||||||
*/
|
*/
|
||||||
import type { Task } from "@fusion/core";
|
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:";
|
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 {
|
export function optionalStepRevisionLogOutcome(details: string, key: string): string {
|
||||||
return `${details}\n${OPTIONAL_STEP_REVISION_KEY_MARKER} ${key}`;
|
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)
|
// Re-export for backward compatibility (tests import from executor.ts)
|
||||||
export { summarizeToolArgs } from "../agents/agent-logger.js";
|
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 {
|
export {
|
||||||
createAgentCreateTool,
|
createAgentCreateTool,
|
||||||
createAgentDeleteTool,
|
createAgentDeleteTool,
|
||||||
|
|||||||
@@ -57,6 +57,7 @@ import {
|
|||||||
resolvePersistAgentThinkingLog,
|
resolvePersistAgentThinkingLog,
|
||||||
resolveTaskLifecycleColumns,
|
resolveTaskLifecycleColumns,
|
||||||
resolveWorkflowIrForTask,
|
resolveWorkflowIrForTask,
|
||||||
|
resolveAgentActivityAttribution,
|
||||||
serializeRetryStormError,
|
serializeRetryStormError,
|
||||||
} from "@fusion/core";
|
} from "@fusion/core";
|
||||||
import type { AgentSession } from "@earendil-works/pi-coding-agent";
|
import type { AgentSession } from "@earendil-works/pi-coding-agent";
|
||||||
@@ -170,6 +171,8 @@ import {
|
|||||||
} from "./session-worktree-paths.js";
|
} from "./session-worktree-paths.js";
|
||||||
import { isWorkflowStepSkillDiscoverable, mergeAdditionalSkillPaths } from "./skill-path-helpers.js";
|
import { isWorkflowStepSkillDiscoverable, mergeAdditionalSkillPaths } from "./skill-path-helpers.js";
|
||||||
import { getExecutorSystemPrompt } from "./system-prompt.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 { createConfiguredCommandAbortError, createSeenSteeringIds } from "./task-predicates.js";
|
||||||
import {
|
import {
|
||||||
accumulateTokenUsage as accumulateTokenUsageImpl,
|
accumulateTokenUsage as accumulateTokenUsageImpl,
|
||||||
@@ -416,6 +419,9 @@ export async function runImplementation(
|
|||||||
runId: syntheticRunId,
|
runId: syntheticRunId,
|
||||||
agentId: task.assignedAgentId ?? "executor",
|
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)
|
// Build engine run context for audit instrumentation (FN-1404)
|
||||||
const engineRunContext: EngineRunContext = {
|
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 agentRotationEvent: import("../credential-instance-rotation.js").RotationEvent | undefined;
|
||||||
let agentRotationDeclined = false;
|
let agentRotationDeclined = false;
|
||||||
@@ -2041,11 +2050,14 @@ export async function runImplementation(
|
|||||||
// give the agent logger the context it needs to emit usage_events tool
|
// give the agent logger the context it needs to emit usage_events tool
|
||||||
// rows (KTD3). nodeId is sourced from the routed/effective node, null
|
// rows (KTD3). nodeId is sourced from the routed/effective node, null
|
||||||
// when the task has no node context.
|
// when the task has no node context.
|
||||||
agentLogger.setUsageContext({
|
attachAgentUsageTelemetry(agentLogger, {
|
||||||
|
store: deps.store,
|
||||||
model: executorModelId ?? null,
|
model: executorModelId ?? null,
|
||||||
provider: executorProvider ?? null,
|
provider: executorProvider ?? null,
|
||||||
nodeId: detail.effectiveNodeId ?? detail.nodeId ?? null,
|
nodeId: detail.effectiveNodeId ?? detail.nodeId ?? null,
|
||||||
agentId: engineRunContext.agentId ?? null,
|
agentId: engineRunContext.agentId ?? null,
|
||||||
|
taskId: task.id,
|
||||||
|
lane: "executor",
|
||||||
});
|
});
|
||||||
|
|
||||||
// Determine whether we're resuming a previous session (pause/resume)
|
// Determine whether we're resuming a previous session (pause/resume)
|
||||||
@@ -2165,6 +2177,14 @@ export async function runImplementation(
|
|||||||
});
|
});
|
||||||
session = createdSession.session;
|
session = createdSession.session;
|
||||||
sessionFile = createdSession.sessionFile;
|
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) {
|
} catch (sessionStartError) {
|
||||||
if (await deps.recoverMissingWorktreeSessionStartFailure(task, worktreePath, sessionStartError, audit)) {
|
if (await deps.recoverMissingWorktreeSessionStartFailure(task, worktreePath, sessionStartError, audit)) {
|
||||||
return;
|
return;
|
||||||
@@ -2617,6 +2637,8 @@ export async function runImplementation(
|
|||||||
taskId: task.id,
|
taskId: task.id,
|
||||||
});
|
});
|
||||||
retrySession = createdRetrySession.session;
|
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);
|
await deps.captureExecutorTokenUsageBaseline(task.id, retrySession);
|
||||||
captureSessionTokenBaseline(retrySession);
|
captureSessionTokenBaseline(retrySession);
|
||||||
if (createdRetrySession.sessionFile) {
|
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.`;
|
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 "";
|
if (!taskCreateWithheld && !delegateWithheld) return "";
|
||||||
const withheld = [
|
const withheld = [
|
||||||
...(taskCreateWithheld ? ["`fn_task_create`"] : []),
|
...(taskCreateWithheld ? ["`fn_task_create`"] : []),
|
||||||
...(delegateWithheld ? ["`fn_delegate_task`"] : []),
|
...(delegateWithheld ? ["`fn_delegate_task`"] : []),
|
||||||
].join(" and ");
|
].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
|
return `## Follow-up task creation is disabled for this session
|
||||||
|
|
||||||
This project's "Ephemeral agent follow-up tasks" policy withholds ${withheld}. ${
|
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"
|
taskCreateWithheld && delegateWithheld ? "them" : "it"
|
||||||
}, and do not retry.
|
}, 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. */
|
/** Resolve the executor system prompt from settings, falling back to the hardcoded constant. */
|
||||||
@@ -320,6 +330,7 @@ export function getExecutorSystemPrompt(
|
|||||||
getWithheldTaskCreationGuidance(
|
getWithheldTaskCreationGuidance(
|
||||||
toolAvailability?.taskCreateWithheld === true,
|
toolAvailability?.taskCreateWithheld === true,
|
||||||
toolAvailability?.delegateWithheld === true,
|
toolAvailability?.delegateWithheld === true,
|
||||||
|
maximumRecommendations,
|
||||||
),
|
),
|
||||||
].filter((section) => section.trim());
|
].filter((section) => section.trim());
|
||||||
return sections.join("\n\n");
|
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.
|
* `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 }>();
|
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 executing = new Set<string>();
|
||||||
protected resumingUnpaused = new Set<string>();
|
protected resumingUnpaused = new Set<string>();
|
||||||
protected approvalSuspended = 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 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); }
|
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); }
|
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 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, 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 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 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); }
|
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;
|
workflowAgentCapacity: WorkflowAgentCapacity;
|
||||||
activeWorkflowAuthorities: Map<string, ActiveWorkflowAuthority>;
|
activeWorkflowAuthorities: Map<string, ActiveWorkflowAuthority>;
|
||||||
activeWorkflowPrincipals: Map<string, { agentId: string; nodeInstanceId: string; agent?: Agent }>;
|
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>;
|
workflowCapacityAttemptIds: Set<string>;
|
||||||
directWorkflowPrincipalWorkItemIds: Set<string>;
|
directWorkflowPrincipalWorkItemIds: Set<string>;
|
||||||
/** Holds written for principal unavailability so terminalization can skip re-closing them. */
|
/** 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,
|
agentId: routed.route.agent.id,
|
||||||
nodeInstanceId,
|
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;
|
if (durableWorkItemId) context["workflow:work-item-id"] = durableWorkItemId;
|
||||||
context["workflow:principal-agent-id"] = routed.route.agent.id;
|
context["workflow:principal-agent-id"] = routed.route.agent.id;
|
||||||
context["workflow:principal-role"] = routed.route.role;
|
context["workflow:principal-role"] = routed.route.role;
|
||||||
|
|||||||
Reference in New Issue
Block a user