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,
|
worktreePath: string,
|
||||||
settings: Settings,
|
settings: Settings,
|
||||||
taskEnv?: NodeJS.ProcessEnv,
|
taskEnv?: NodeJS.ProcessEnv,
|
||||||
stepOptions?: { unattended?: boolean },
|
stepOptions?: { unattended?: boolean; principalAgentId?: string },
|
||||||
): Promise<WorkflowStepOutcome> {
|
): Promise<WorkflowStepOutcome> {
|
||||||
let toolMode: "coding" | "readonly" = workflowStep.toolMode || "readonly";
|
let toolMode: "coding" | "readonly" = workflowStep.toolMode || "readonly";
|
||||||
// (U3) Genuinely-unattended run — set FUSION_HEADLESS=1 below so skills record
|
// (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}`;
|
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({
|
const agentLogger = new AgentLogger({
|
||||||
store: deps.store,
|
store: deps.store,
|
||||||
taskId: task.id,
|
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).
|
// 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
|
// 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
|
||||||
@@ -466,7 +485,8 @@ CRITICAL SCOPING RULES — read before doing anything else:
|
|||||||
// steps to inherit project execution-lane model settings before defaults.
|
// steps to inherit project execution-lane model settings before defaults.
|
||||||
// Review gates are independent validation surfaces and must not silently use
|
// Review gates are independent validation surfaces and must not silently use
|
||||||
// the same implementation model merely because they execute in this method.
|
// 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
|
const laneModel = isReviewTypeWorkflowStep
|
||||||
? resolveValidatorSessionModel(
|
? resolveValidatorSessionModel(
|
||||||
task.validatorModelProvider,
|
task.validatorModelProvider,
|
||||||
@@ -487,7 +507,7 @@ CRITICAL SCOPING RULES — read before doing anything else:
|
|||||||
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" });
|
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
|
const workflowFallback = isReviewTypeWorkflowStep
|
||||||
? resolveValidatorFallbackModel(settings)
|
? resolveValidatorFallbackModel(settings)
|
||||||
@@ -514,13 +534,13 @@ CRITICAL SCOPING RULES — read before doing anything else:
|
|||||||
// Build skill selection context for workflow step session
|
// Build skill selection context for workflow step session
|
||||||
const skillContext = await buildSessionSkillContext({
|
const skillContext = await buildSessionSkillContext({
|
||||||
agentStore: deps.options.agentStore!,
|
agentStore: deps.options.agentStore!,
|
||||||
task,
|
task: sessionTask,
|
||||||
sessionPurpose: "executor",
|
sessionPurpose: "executor",
|
||||||
projectRootDir: deps.rootDir,
|
projectRootDir: deps.rootDir,
|
||||||
pluginRunner: deps.options.pluginRunner,
|
pluginRunner: deps.options.pluginRunner,
|
||||||
});
|
});
|
||||||
|
|
||||||
const workflowAgent = await deps.getAuthoritativeAssignedAgent(task.assignedAgentId);
|
const workflowAgent = workflowPrincipal ?? await deps.getAuthoritativeAssignedAgent(task.assignedAgentId);
|
||||||
const workflowRuntimeHint = extractRuntimeHint(workflowAgent?.runtimeConfig);
|
const workflowRuntimeHint = extractRuntimeHint(workflowAgent?.runtimeConfig);
|
||||||
// Signal to skills running in this step (e.g. compound-engineering ce-plan /
|
// 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
|
// 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).
|
// 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(
|
const workflowModelDetails = formatModelMarkerDetails(
|
||||||
describeModel(session),
|
describeModel(session),
|
||||||
|
|||||||
@@ -451,9 +451,18 @@ export async function runGraphCustomNode(
|
|||||||
// sets FUSION_HEADLESS=1 only when this is explicitly true.
|
// sets FUSION_HEADLESS=1 only when this is explicitly true.
|
||||||
const unattended = deps.graphUnattendedRuns.has(live.id);
|
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"
|
let outcome: WorkflowStepOutcome = mode === "script"
|
||||||
? await deps.executeScriptWorkflowStep(live, step, worktreePath, settings, nodeEnv)
|
? 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:
|
* FNXC:WorkflowReviewFindings 2026-08-05-06:29:
|
||||||
* Script nodes retain their exit-code verdict semantics, but an explicitly classified review
|
* Script nodes retain their exit-code verdict semantics, but an explicitly classified review
|
||||||
|
|||||||
Reference in New Issue
Block a user