FN-9108: Restore routed workflow-principal session identity

Restore routed principal attribution throughout workflow prompt and review sessions.

- Thread the graph-selected principal into workflow-step execution.
- Resolve principal runtime configuration, skills, telemetry, and session attribution consistently.
- Fail closed when the routed principal is unavailable and add regression coverage.
- Add a patch changeset for the restored behavior.

Files changed:
 .../fn-9108-workflow-principal-session-identity.md |   7 +
 .../executor-workflow-step-principal.test.ts       | 154 +++++++++++++++++++++
 .../engine/src/executor/execute-workflow-step.ts   |  34 ++++-
 .../engine/src/executor/run-graph-custom-node.ts   |  11 +-
 4 files changed, 198 insertions(+), 8 deletions(-)

Fusion-Task-Id: FN-9108

Fusion-Task-Lineage: ec668551-8ff4-4ced-8c7f-a4aa04eb044f

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-15 17:45:32 -07:00
parent 872d260ab8
commit 66bbeaa28d
4 changed files with 198 additions and 8 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": patch
---
summary: Restore routed workflow-principal identity for prompt and review sessions.
category: fix
dev: Restores executeWorkflowStep and runGraphCustomNode principal threading.

View File

@@ -0,0 +1,154 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import "./executor-test-helpers.js";
import { TaskExecutor } from "../executor.js";
import {
createMockStore,
mockedCreateFnAgent,
mockedExecSync,
resetExecutorMocks,
} from "./executor-test-helpers.js";
type CapturedSession = {
defaultProvider?: string;
defaultModelId?: string;
runtimeHint?: string;
};
function captureSession(): { last?: CapturedSession } {
const holder: { last?: CapturedSession } = {};
mockedCreateFnAgent.mockImplementation(async (options: any) => {
holder.last = {
defaultProvider: options.defaultProvider,
defaultModelId: options.defaultModelId,
runtimeHint: options.runtimeHint,
};
const listeners: Array<(event: any) => void> = [];
return {
session: {
state: {},
subscribe: (listener: (event: any) => void) => {
listeners.push(listener);
return () => {};
},
prompt: vi.fn(async () => {
for (const listener of listeners) {
listener({ type: "message_update", assistantMessageEvent: { type: "text_delta", contentIndex: 0, delta: '{"verdict":"APPROVE","notes":""}' } });
}
}),
dispose: vi.fn(),
},
};
});
return holder;
}
function task(overrides: Record<string, unknown> = {}) {
return {
id: "FN-PRINCIPAL-1",
title: "Principal session",
description: "verify routed identity",
column: "in-progress" as const,
worktree: "/tmp/wt",
branch: "fusion/fn-principal-1",
baseCommitSha: "abc123",
dependencies: [],
steps: [{ name: "s", status: "in-progress" as const }],
currentStep: 0,
log: [],
createdAt: new Date().toISOString(),
updatedAt: new Date().toISOString(),
...overrides,
};
}
function step(overrides: Record<string, unknown> = {}) {
const now = new Date().toISOString();
return {
id: "step:principal",
name: "Code Review",
description: "",
mode: "prompt" as const,
phase: "pre-merge" as const,
gateMode: "advisory" as const,
prompt: "Review this.",
toolMode: "readonly" as const,
enabled: true,
createdAt: now,
updatedAt: now,
...overrides,
};
}
function makeExecutor() {
const store = createMockStore();
store.getSettings.mockResolvedValue({});
const agents = new Map([
["agent-principal", { id: "agent-principal", runtimeConfig: { model: "anthropic/claude-principal", modelProvider: "anthropic", modelId: "claude-principal", runtimeHint: "principal-hint" } }],
["agent-assigned", { id: "agent-assigned", runtimeConfig: { model: "openai/gpt-assigned", modelProvider: "openai", modelId: "gpt-assigned", runtimeHint: "assigned-hint" } }],
]);
const agentStore = { getAgent: vi.fn(async (id: string) => agents.get(id) ?? null), createAgent: vi.fn() };
return { store, executor: new TaskExecutor(store as any, "/tmp/test", { agentStore } as any) };
}
describe("executor workflow-step routed principal", () => {
beforeEach(() => {
resetExecutorMocks();
mockedExecSync.mockImplementation(() => Buffer.from(""));
});
it.each([
["review", step()],
["ordinary prompt", step({ name: "Implementation Prompt" })],
])("uses the routed principal runtime identity for a %s step", async (_lane, workflowStep) => {
const { executor } = makeExecutor();
const captured = captureSession();
const assignedConfig = vi.spyOn(executor as any, "getAssignedAgentRuntimeConfig");
await (executor as any).executeWorkflowStep(
task({ assignedAgentId: "agent-assigned" }), workflowStep, "/tmp/wt", {}, undefined,
{ principalAgentId: "agent-principal" },
);
expect(captured.last).toMatchObject({
defaultProvider: "anthropic",
defaultModelId: "claude-principal",
runtimeHint: "principal-hint",
});
expect(assignedConfig).not.toHaveBeenCalled();
});
it("threads only string graph principal context to prompt steps", async () => {
const { store, executor } = makeExecutor();
store.getTask.mockImplementation(async (id: string) => task({ id }));
const execute = vi.spyOn(executor as any, "executeWorkflowStep").mockResolvedValue({ success: true, output: "ok" });
const executeScript = vi.spyOn(executor as any, "executeScriptWorkflowStep").mockResolvedValue({ success: true, output: "ok" });
const promptNode = { id: "principal-prompt", kind: "prompt", config: { prompt: "Review" } };
const scriptNode = { id: "principal-script", kind: "script", config: { command: "true" } };
await (executor as any).runGraphCustomNode(promptNode, task(), {}, undefined, { "workflow:principal-agent-id": "agent-principal" });
await (executor as any).runGraphCustomNode(promptNode, task(), {}, undefined, { "workflow:principal-agent-id": 42 });
await (executor as any).runGraphCustomNode(scriptNode, task(), {}, undefined, { "workflow:principal-agent-id": "agent-principal" });
expect(execute.mock.calls[0][5]).toMatchObject({ principalAgentId: "agent-principal" });
expect(execute.mock.calls[1][5]).toMatchObject({ principalAgentId: undefined });
expect(execute).toHaveBeenCalledTimes(2);
expect(executeScript.mock.calls[0]).toHaveLength(5);
});
it("fails closed instead of falling back when the routed principal is unavailable", async () => {
const { executor } = makeExecutor();
await expect((executor as any).executeWorkflowStep(
task({ assignedAgentId: "agent-assigned" }), step(), "/tmp/wt", {}, undefined,
{ principalAgentId: "agent-missing" },
)).rejects.toThrow(/^workflow-principal-unavailable:agent-missing$/);
});
it("keeps unrouted steps on their assigned agent", async () => {
const { executor } = makeExecutor();
const captured = captureSession();
await (executor as any).executeWorkflowStep(
task({ assignedAgentId: "agent-assigned" }), step({ name: "Implementation Prompt" }), "/tmp/wt", {}, undefined,
);
expect(captured.last).toMatchObject({ defaultProvider: "openai", defaultModelId: "gpt-assigned", runtimeHint: "assigned-hint" });
});
});

View File

@@ -125,7 +125,7 @@ export async function executeWorkflowStep(
worktreePath: string,
settings: Settings,
taskEnv?: NodeJS.ProcessEnv,
stepOptions?: { unattended?: boolean },
stepOptions?: { unattended?: boolean; principalAgentId?: string },
): Promise<WorkflowStepOutcome> {
let toolMode: "coding" | "readonly" = workflowStep.toolMode || "readonly";
// (U3) Genuinely-unattended run — set FUSION_HEADLESS=1 below so skills record
@@ -442,6 +442,25 @@ CRITICAL SCOPING RULES — read before doing anything else:
You have access to the file system to review changes.${inlineFixBlock}${verdictBlock}`;
/*
* FNXC:WorkflowAgentRouting 2026-08-07-04:45:
* The graph admission fence chooses the permanent identity before this
* session exists. Resolve that exact agent for its model, skills, audit,
* and log attribution; never fall back to task ownership after routing.
*
* FNXC:WorkflowAgentRouting 2026-08-15-23:41:
* Wave-18 peel #3317 dropped this threading. FN-9108 restores the
* FN-8764/FN-8821 routed-principal contract and regression coverage.
*/
const workflowPrincipal = stepOptions?.principalAgentId
? await deps.getAuthoritativeAssignedAgent(stepOptions.principalAgentId)
: undefined;
if (stepOptions?.principalAgentId && !workflowPrincipal) {
throw new Error(`workflow-principal-unavailable:${stepOptions.principalAgentId}`);
}
const sessionTask = workflowPrincipal
? { ...task, assignedAgentId: workflowPrincipal.id }
: task;
const agentLogger = new AgentLogger({
store: deps.store,
taskId: task.id,
@@ -457,7 +476,7 @@ CRITICAL SCOPING RULES — read before doing anything else:
},
});
// 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" });
attachAgentUsageTelemetry(agentLogger, { store: deps.store, agentId: sessionTask.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, lane: "executor" });
// Determine primary model and an explicit fallback. Review-type workflow
// steps use the validator lane; ordinary workflow prompts use the executor
@@ -466,7 +485,8 @@ CRITICAL SCOPING RULES — read before doing anything else:
// steps to inherit project execution-lane model settings before defaults.
// Review gates are independent validation surfaces and must not silently use
// the same implementation model merely because they execute in this method.
const assignedRuntimeConfig = await deps.getAssignedAgentRuntimeConfig(task.assignedAgentId);
const assignedRuntimeConfig = workflowPrincipal?.runtimeConfig
?? await deps.getAssignedAgentRuntimeConfig(task.assignedAgentId);
const laneModel = isReviewTypeWorkflowStep
? resolveValidatorSessionModel(
task.validatorModelProvider,
@@ -487,7 +507,7 @@ CRITICAL SCOPING RULES — read before doing anything else:
const primaryModelId = useOverride ? workflowStep.modelId : laneModel.modelId;
// FNXC:ProviderAuth 2026-08-01-08:39: A workflow-step model override has no paired instance selection, so only the resolved primary task lane may carry its requested credential instance. Fallback attempts must retain their provider-default behavior rather than inheriting a primary-provider identity.
const primaryCredentialInstanceId = useOverride ? undefined : laneModel.credentialInstanceId;
attachAgentUsageTelemetry(agentLogger, { store: deps.store, agentId: task.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, model: primaryModelId ?? null, provider: primaryProvider ?? null, lane: "executor" });
attachAgentUsageTelemetry(agentLogger, { store: deps.store, agentId: sessionTask.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, model: primaryModelId ?? null, provider: primaryProvider ?? null, lane: "executor" });
const workflowFallback = isReviewTypeWorkflowStep
? resolveValidatorFallbackModel(settings)
@@ -514,13 +534,13 @@ CRITICAL SCOPING RULES — read before doing anything else:
// Build skill selection context for workflow step session
const skillContext = await buildSessionSkillContext({
agentStore: deps.options.agentStore!,
task,
task: sessionTask,
sessionPurpose: "executor",
projectRootDir: deps.rootDir,
pluginRunner: deps.options.pluginRunner,
});
const workflowAgent = await deps.getAuthoritativeAssignedAgent(task.assignedAgentId);
const workflowAgent = workflowPrincipal ?? await deps.getAuthoritativeAssignedAgent(task.assignedAgentId);
const workflowRuntimeHint = extractRuntimeHint(workflowAgent?.runtimeConfig);
// Signal to skills running in this step (e.g. compound-engineering ce-plan /
// ce-work) that they are inside a Fusion autonomous workflow step, NOT an
@@ -691,7 +711,7 @@ CRITICAL SCOPING RULES — read before doing anything else:
: {}),
});
// 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" });
emitAgentSessionStart({ store: deps.store, agentId: sessionTask.assignedAgentId ?? null, taskId: task.id, nodeId: task.effectiveNodeId ?? task.nodeId ?? null, model: primaryModelId ?? null, provider: primaryProvider ?? null, lane: "executor" });
const workflowModelDetails = formatModelMarkerDetails(
describeModel(session),

View File

@@ -451,9 +451,18 @@ export async function runGraphCustomNode(
// sets FUSION_HEADLESS=1 only when this is explicitly true.
const unattended = deps.graphUnattendedRuns.has(live.id);
/*
* FNXC:WorkflowAgentRouting 2026-08-15-23:41:
* FN-8764/FN-8821 routing selects a principal before graph execution. Wave-18
* peel #3317 dropped the handoff to prompt sessions; FN-9108 restores it.
* Non-string context is unrouted so untrusted graph values cannot become IDs.
*/
const principalAgentId = typeof graphContext?.["workflow:principal-agent-id"] === "string"
? graphContext["workflow:principal-agent-id"]
: undefined;
let outcome: WorkflowStepOutcome = mode === "script"
? await deps.executeScriptWorkflowStep(live, step, worktreePath, settings, nodeEnv)
: await deps.executeWorkflowStep(live, step, worktreePath, settings, nodeEnv, { unattended });
: await deps.executeWorkflowStep(live, step, worktreePath, settings, nodeEnv, { unattended, principalAgentId });
/*
* FNXC:WorkflowReviewFindings 2026-08-05-06:29:
* Script nodes retain their exit-code verdict semantics, but an explicitly classified review