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:
@@ -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.
|
||||
@@ -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" });
|
||||
});
|
||||
});
|
||||
@@ -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),
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user