diff --git a/.changeset/fn-9108-workflow-principal-session-identity.md b/.changeset/fn-9108-workflow-principal-session-identity.md new file mode 100644 index 0000000000..65ca5e12eb --- /dev/null +++ b/.changeset/fn-9108-workflow-principal-session-identity.md @@ -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. diff --git a/packages/engine/src/__tests__/executor-workflow-step-principal.test.ts b/packages/engine/src/__tests__/executor-workflow-step-principal.test.ts new file mode 100644 index 0000000000..36a37b7b64 --- /dev/null +++ b/packages/engine/src/__tests__/executor-workflow-step-principal.test.ts @@ -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 = {}) { + 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 = {}) { + 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" }); + }); +}); diff --git a/packages/engine/src/executor/execute-workflow-step.ts b/packages/engine/src/executor/execute-workflow-step.ts index dddf6770f6..ab3c2e9c45 100644 --- a/packages/engine/src/executor/execute-workflow-step.ts +++ b/packages/engine/src/executor/execute-workflow-step.ts @@ -125,7 +125,7 @@ export async function executeWorkflowStep( worktreePath: string, settings: Settings, taskEnv?: NodeJS.ProcessEnv, - stepOptions?: { unattended?: boolean }, + stepOptions?: { unattended?: boolean; principalAgentId?: string }, ): Promise { 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), diff --git a/packages/engine/src/executor/run-graph-custom-node.ts b/packages/engine/src/executor/run-graph-custom-node.ts index b849fe5850..6b72d80b7a 100644 --- a/packages/engine/src/executor/run-graph-custom-node.ts +++ b/packages/engine/src/executor/run-graph-custom-node.ts @@ -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