feat(engine): execute and step-execute sessions run as the column agent
U4: seams stamp the governing node id into run context; a per-seam resolveSeamColumnAgent feeds the core resolver and threads the effective agent through resolveExecutorSessionModel, runtime hints, persona, memory tools, and StepSessionExecutor attribution. Characterization tests pin the no-binding path byte-identical. Gating/deferral principal moves in U5.
This commit is contained in:
7
.changeset/workflow-column-agent-assignment.md
Normal file
7
.changeset/workflow-column-agent-assignment.md
Normal file
@@ -0,0 +1,7 @@
|
|||||||
|
---
|
||||||
|
"@runfusion/fusion": minor
|
||||||
|
---
|
||||||
|
|
||||||
|
Add per-column agent assignment for workflow columns, behind the combined `experimentalFeatures.workflowColumns` + `experimentalFeatures.workflowGraphExecutor` flags.
|
||||||
|
|
||||||
|
A workflow column can now name a permanent agent from the registry plus a mode — `defer` (the column agent is the default for work in that column that carries no agent/model settings of its own) or `override` (the column agent supersedes node- and task-level agent/model settings). The binding applies to all session-running work attributable to the column's nodes: custom prompt/gate/script nodes, the execute seam's coding session, and step-execute sessions. Precedence is resolved by one shared `@fusion/core` resolver (`resolveColumnAgentBinding` + `resolveEffectiveAgent`) consumed by every reader, with defer/override expressed as explicit named rules and defer granularity all-or-nothing (an own agent identity OR a complete `modelProvider`+`modelId` pair suppresses the column agent). The binding keys off the node's declared IR column; foreach template nodes inherit the enclosing foreach node's column. A missing/deleted agent at resolution time logs and falls back to normal resolution — a live session is never aborted. The built-in default workflow carries no column agents and stays byte-identical (parity oracle); with either flag off, column agents are inert.
|
||||||
@@ -0,0 +1,398 @@
|
|||||||
|
// Column-agent coding seams: execute + step-execute sessions (plan U4,
|
||||||
|
// R2/R3/R4/R8, KTD-2/KTD-3/KTD-5/KTD-6).
|
||||||
|
//
|
||||||
|
// The graph EXECUTE seam (single coding session) and STEP-EXECUTE seam
|
||||||
|
// (StepSessionExecutor per-step sessions) must run as the column agent when the
|
||||||
|
// governing seam node's DECLARED column carries a binding. Session identity =
|
||||||
|
// model + persona + attribution (gating/heartbeat/restart are U5, untouched here).
|
||||||
|
//
|
||||||
|
// Harness: mirrors executor-step-session.test.ts / executor-column-agent-custom-
|
||||||
|
// node.test.ts — a real TaskExecutor over a mock store with `createFnAgent`
|
||||||
|
// (the outermost session-spawn boundary) mocked, plus the entirely-mocked
|
||||||
|
// StepSessionExecutor from executor-test-helpers so the step-session branch's
|
||||||
|
// constructor options are observable.
|
||||||
|
//
|
||||||
|
// The two per-run seam slots the executor reads — `graphSeamGoverningNodeId` and
|
||||||
|
// `graphColumnAgentResolver` — are normally stamped by the graph seam wiring
|
||||||
|
// (createPromptLikeHandler → execute/stepExecute seams). We seed them directly and
|
||||||
|
// drive `runImplementationPhase` (the exact call the execute seam makes, which
|
||||||
|
// registers a completion interceptor so graph routing is skipped) so the session
|
||||||
|
// build runs the production resolution path with no scripted session layer.
|
||||||
|
|
||||||
|
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||||
|
import "./executor-test-helpers.js";
|
||||||
|
import { TaskExecutor } from "../executor.js";
|
||||||
|
import {
|
||||||
|
createMockStore,
|
||||||
|
mockedCreateFnAgent,
|
||||||
|
mockedStepSessionExecutor,
|
||||||
|
mockExecuteAll,
|
||||||
|
resetExecutorMocks,
|
||||||
|
} from "./executor-test-helpers.js";
|
||||||
|
import type { WorkflowColumnAgent } from "@fusion/core";
|
||||||
|
|
||||||
|
// The mocked resolveExecutorSessionModel (executor-test-helpers) reads
|
||||||
|
// `runtimeConfig.model` in "provider/modelId" form, so the column agent advertises
|
||||||
|
// its model that way; the assigned agent advertises a different one so we can prove
|
||||||
|
// which one reached the session.
|
||||||
|
function makeColumnAgent(overrides: Record<string, unknown> = {}) {
|
||||||
|
return {
|
||||||
|
id: "agent-col",
|
||||||
|
name: "Senior Reviewer",
|
||||||
|
soul: "I am the senior reviewer.",
|
||||||
|
instructionsText: "Always be thorough.",
|
||||||
|
memory: undefined,
|
||||||
|
runtimeConfig: { model: "anthropic/claude-col", runtimeHint: "col-hint" },
|
||||||
|
...overrides,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function makeAssignedAgent(overrides: Record<string, unknown> = {}) {
|
||||||
|
return {
|
||||||
|
id: "agent-Y",
|
||||||
|
name: "Assigned Agent",
|
||||||
|
soul: "I am the assigned agent.",
|
||||||
|
instructionsText: "Assigned persona.",
|
||||||
|
memory: undefined,
|
||||||
|
runtimeConfig: { model: "openai/gpt-assigned", runtimeHint: "assigned-hint" },
|
||||||
|
...overrides,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A mock fn agent that immediately calls fn_task_done so execute() completes. */
|
||||||
|
function installTaskDoneAgent() {
|
||||||
|
mockedCreateFnAgent.mockImplementation((async (opts: any) => {
|
||||||
|
const tools = opts.customTools || [];
|
||||||
|
return {
|
||||||
|
session: {
|
||||||
|
prompt: vi.fn().mockImplementation(async () => {
|
||||||
|
const done = tools.find((t: any) => t.name === "fn_task_done");
|
||||||
|
if (done) await done.execute("tool-1", {});
|
||||||
|
}),
|
||||||
|
dispose: vi.fn(),
|
||||||
|
subscribe: vi.fn(),
|
||||||
|
on: vi.fn(),
|
||||||
|
sessionManager: { getLeafId: vi.fn().mockReturnValue("leaf-1") },
|
||||||
|
state: {},
|
||||||
|
},
|
||||||
|
};
|
||||||
|
}) as any);
|
||||||
|
}
|
||||||
|
|
||||||
|
function makeExecutor(store: ReturnType<typeof createMockStore>, agentsById: Record<string, unknown>) {
|
||||||
|
const agentStore = {
|
||||||
|
getAgent: vi.fn(async (id: string) => agentsById[id] ?? null),
|
||||||
|
};
|
||||||
|
const executor = new TaskExecutor(store as any, "/tmp/test", { agentStore } as any);
|
||||||
|
return { executor, agentStore };
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Seed the per-run column-agent seam slots the executor reads at session-build
|
||||||
|
* time, then drive the implementation phase the way the execute seam does.
|
||||||
|
*/
|
||||||
|
async function runExecuteSeam(
|
||||||
|
executor: TaskExecutor,
|
||||||
|
task: any,
|
||||||
|
governingNodeId: string,
|
||||||
|
binding: WorkflowColumnAgent | undefined,
|
||||||
|
) {
|
||||||
|
(executor as any).graphSeamGoverningNodeId.set(task.id, governingNodeId);
|
||||||
|
(executor as any).graphColumnAgentResolver.set(task.id, (nodeId: string) =>
|
||||||
|
nodeId === governingNodeId ? binding : undefined,
|
||||||
|
);
|
||||||
|
return (executor as any).runImplementationPhase(task);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Force the step-session physics path and seed the seam slots, then run. */
|
||||||
|
async function runStepSessionSeam(
|
||||||
|
executor: TaskExecutor,
|
||||||
|
task: any,
|
||||||
|
governingNodeId: string,
|
||||||
|
binding: WorkflowColumnAgent | undefined,
|
||||||
|
) {
|
||||||
|
(executor as any).graphStepSessionPinned.add(task.id);
|
||||||
|
(executor as any).graphSeamGoverningNodeId.set(task.id, governingNodeId);
|
||||||
|
(executor as any).graphColumnAgentResolver.set(task.id, (nodeId: string) =>
|
||||||
|
nodeId === governingNodeId ? binding : undefined,
|
||||||
|
);
|
||||||
|
return (executor as any).runImplementationPhase(task);
|
||||||
|
}
|
||||||
|
|
||||||
|
function singleSessionTask(overrides: Record<string, unknown> = {}) {
|
||||||
|
return {
|
||||||
|
id: "FN-001",
|
||||||
|
title: "Test",
|
||||||
|
description: "Test task",
|
||||||
|
column: "in-progress",
|
||||||
|
dependencies: [],
|
||||||
|
steps: [{ name: "Implement", status: "in-progress" }],
|
||||||
|
currentStep: 0,
|
||||||
|
log: [],
|
||||||
|
prompt: "# test\n## Steps\n### Step 0: Implement\n- [ ] implement",
|
||||||
|
createdAt: new Date().toISOString(),
|
||||||
|
updatedAt: new Date().toISOString(),
|
||||||
|
...overrides,
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
function lastFnAgentOpts() {
|
||||||
|
const calls = mockedCreateFnAgent.mock.calls;
|
||||||
|
return calls[calls.length - 1]?.[0] as any;
|
||||||
|
}
|
||||||
|
|
||||||
|
function lastStepExecutorOpts() {
|
||||||
|
const calls = mockedStepSessionExecutor.mock.calls;
|
||||||
|
return calls[calls.length - 1]?.[0] as any;
|
||||||
|
}
|
||||||
|
|
||||||
|
function loggedLines(store: ReturnType<typeof createMockStore>): string[] {
|
||||||
|
return store.logEntry.mock.calls.map((call: any[]) => String(call[1] ?? ""));
|
||||||
|
}
|
||||||
|
|
||||||
|
const OVERRIDE_COL: WorkflowColumnAgent = { agentId: "agent-col", mode: "override" };
|
||||||
|
const DEFER_COL: WorkflowColumnAgent = { agentId: "agent-col", mode: "defer" };
|
||||||
|
|
||||||
|
describe("column-agent coding seams (plan U4)", () => {
|
||||||
|
beforeEach(() => {
|
||||||
|
resetExecutorMocks();
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── Characterization (pre-substitution behavior) ──────────────────────────
|
||||||
|
// These pin the assignedAgentId-driven session identity that exists today and
|
||||||
|
// MUST stay byte-identical on the no-binding path after substitution.
|
||||||
|
|
||||||
|
describe("characterization: no binding → assignedAgentId session identity unchanged", () => {
|
||||||
|
it("execute seam: session model/persona built from the assigned agent, no column-agent log", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({ assignedAgentId: "agent-Y" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor } = makeExecutor(store, { "agent-Y": makeAssignedAgent() });
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
// No governing node / no binding seeded → legacy path.
|
||||||
|
await (executor as any).runImplementationPhase(task);
|
||||||
|
|
||||||
|
const opts = lastFnAgentOpts();
|
||||||
|
// Model resolved from the ASSIGNED agent's runtimeConfig.model.
|
||||||
|
expect(opts.defaultProvider).toBe("openai");
|
||||||
|
expect(opts.defaultModelId).toBe("gpt-assigned");
|
||||||
|
// No column-agent adoption logged.
|
||||||
|
expect(loggedLines(store).some((l) => l.includes("running as column agent"))).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("step session: attribution falls back to assignedAgentId; no effectiveAgentId override", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({ assignedAgentId: "agent-Y" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor } = makeExecutor(store, { "agent-Y": makeAssignedAgent() });
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
(executor as any).graphStepSessionPinned.add(task.id);
|
||||||
|
await (executor as any).runImplementationPhase(task);
|
||||||
|
|
||||||
|
const opts = lastStepExecutorOpts();
|
||||||
|
// No column agent governs → no attribution override (StepSessionExecutor
|
||||||
|
// falls back to taskDetail.assignedAgentId ?? "executor").
|
||||||
|
expect(opts.effectiveAgentId).toBeUndefined();
|
||||||
|
// Model precedence input is the assigned agent's runtimeConfig.
|
||||||
|
expect(opts.assignedAgentRuntimeConfig).toEqual(makeAssignedAgent().runtimeConfig);
|
||||||
|
expect(loggedLines(store).some((l) => l.includes("running as column agent"))).toBe(false);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── Execute seam (single coding session) ──────────────────────────────────
|
||||||
|
|
||||||
|
describe("execute seam", () => {
|
||||||
|
it("override column, task assigned to Y → session uses column agent X's model/persona/identity + audit", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({ assignedAgentId: "agent-Y" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor, agentStore } = makeExecutor(store, {
|
||||||
|
"agent-Y": makeAssignedAgent(),
|
||||||
|
"agent-col": makeColumnAgent(),
|
||||||
|
});
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
await runExecuteSeam(executor, task, "execute-node", OVERRIDE_COL);
|
||||||
|
|
||||||
|
const opts = lastFnAgentOpts();
|
||||||
|
// Column agent X's model supersedes the assigned agent Y's.
|
||||||
|
expect(opts.defaultProvider).toBe("anthropic");
|
||||||
|
expect(opts.defaultModelId).toBe("claude-col");
|
||||||
|
// Persona: column agent's soul + instructionsText reach the session system
|
||||||
|
// prompt layers (KTD-6 typed fields).
|
||||||
|
const promptText = JSON.stringify(opts.systemPromptLayers ?? "") + (opts.systemPrompt ?? "");
|
||||||
|
expect(promptText).toContain("I am the senior reviewer.");
|
||||||
|
expect(promptText).toContain("Always be thorough.");
|
||||||
|
// The column agent was fetched (identity), not just the assigned agent.
|
||||||
|
expect(agentStore.getAgent).toHaveBeenCalledWith("agent-col");
|
||||||
|
// Audit names the substitution + mode.
|
||||||
|
expect(
|
||||||
|
loggedLines(store).some(
|
||||||
|
(l) => l.includes("running as column agent 'agent-col' (override)") && l.includes("execute-node"),
|
||||||
|
),
|
||||||
|
).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("defer column, task with complete modelProvider/modelId → task settings win", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
// Task carries a complete own model pair → defer must yield own settings.
|
||||||
|
const task = singleSessionTask({ modelProvider: "task-prov", modelId: "task-model" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor, agentStore } = makeExecutor(store, { "agent-col": makeColumnAgent() });
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
await runExecuteSeam(executor, task, "execute-node", DEFER_COL);
|
||||||
|
|
||||||
|
const opts = lastFnAgentOpts();
|
||||||
|
// The task's own complete pair wins (mocked resolver: no agent runtimeConfig
|
||||||
|
// model, falls through to the task pair).
|
||||||
|
expect(opts.defaultProvider).toBe("task-prov");
|
||||||
|
expect(opts.defaultModelId).toBe("task-model");
|
||||||
|
// Column agent never fetched/adopted.
|
||||||
|
expect(agentStore.getAgent).not.toHaveBeenCalledWith("agent-col");
|
||||||
|
expect(loggedLines(store).some((l) => l.includes("running as column agent"))).toBe(false);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("defer column, bare task (no own settings) → column agent adopted", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask(); // no assignedAgentId, no model pair
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor, agentStore } = makeExecutor(store, { "agent-col": makeColumnAgent() });
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
await runExecuteSeam(executor, task, "execute-node", DEFER_COL);
|
||||||
|
|
||||||
|
const opts = lastFnAgentOpts();
|
||||||
|
expect(opts.defaultProvider).toBe("anthropic");
|
||||||
|
expect(opts.defaultModelId).toBe("claude-col");
|
||||||
|
expect(agentStore.getAgent).toHaveBeenCalledWith("agent-col");
|
||||||
|
expect(
|
||||||
|
loggedLines(store).some((l) => l.includes("running as column agent 'agent-col' (defer)")),
|
||||||
|
).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("column agent missing from registry at seam time → fallback to assignedAgentId path, logged, run proceeds", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({ assignedAgentId: "agent-Y" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
// Column agent absent from the registry; assigned agent present.
|
||||||
|
const { executor } = makeExecutor(store, { "agent-Y": makeAssignedAgent() });
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
await runExecuteSeam(executor, task, "execute-node", OVERRIDE_COL);
|
||||||
|
|
||||||
|
const opts = lastFnAgentOpts();
|
||||||
|
// Fell back to the assigned agent's model.
|
||||||
|
expect(opts.defaultProvider).toBe("openai");
|
||||||
|
expect(opts.defaultModelId).toBe("gpt-assigned");
|
||||||
|
// Fallback audited; no adoption claim.
|
||||||
|
expect(
|
||||||
|
loggedLines(store).some(
|
||||||
|
(l) => l.includes("column agent 'agent-col' not found") && l.includes("falling back"),
|
||||||
|
),
|
||||||
|
).toBe(true);
|
||||||
|
expect(loggedLines(store).some((l) => l.includes("running as column agent"))).toBe(false);
|
||||||
|
// Run still proceeded: a session was built and the task done tool fired
|
||||||
|
// (the missing column agent never aborted the session — R8).
|
||||||
|
expect(mockedCreateFnAgent).toHaveBeenCalled();
|
||||||
|
});
|
||||||
|
|
||||||
|
it("integration: the column agent's executor model reaches createResolvedAgentSession options end-to-end", async () => {
|
||||||
|
// Per the plugin-skills learning — prove with the REAL resolution layers
|
||||||
|
// (only the outermost createFnAgent/session-spawn boundary is mocked).
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({ assignedAgentId: "agent-Y" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor } = makeExecutor(store, {
|
||||||
|
"agent-Y": makeAssignedAgent(),
|
||||||
|
"agent-col": makeColumnAgent({ runtimeConfig: { model: "anthropic/claude-e2e", runtimeHint: "e2e-hint" } }),
|
||||||
|
});
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
await runExecuteSeam(executor, task, "execute-node", OVERRIDE_COL);
|
||||||
|
|
||||||
|
const opts = lastFnAgentOpts();
|
||||||
|
expect(opts.defaultProvider).toBe("anthropic");
|
||||||
|
expect(opts.defaultModelId).toBe("claude-e2e");
|
||||||
|
// Runtime hint also follows the column agent end-to-end.
|
||||||
|
expect(opts.runtimeHint).toBe("e2e-hint");
|
||||||
|
});
|
||||||
|
});
|
||||||
|
|
||||||
|
// ── Step-execute seam (StepSessionExecutor per-step sessions) ─────────────
|
||||||
|
|
||||||
|
describe("step-execute seam", () => {
|
||||||
|
it("foreach instance node inherits the foreach's bound column → instance session carries column agent identity (attribution asserted)", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({ assignedAgentId: "agent-Y" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor, agentStore } = makeExecutor(store, {
|
||||||
|
"agent-Y": makeAssignedAgent(),
|
||||||
|
"agent-col": makeColumnAgent(),
|
||||||
|
});
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
// Governing node is the foreach INSTANCE id; the resolver (which the real
|
||||||
|
// core resolver implements via template inheritance) returns the foreach's
|
||||||
|
// bound column binding for that instance id.
|
||||||
|
const instanceNodeId = "foreach-1#0:step-exec";
|
||||||
|
await runStepSessionSeam(executor, task, instanceNodeId, OVERRIDE_COL);
|
||||||
|
|
||||||
|
const opts = lastStepExecutorOpts();
|
||||||
|
// Attribution: the per-step session is attributed to the column agent.
|
||||||
|
expect(opts.effectiveAgentId).toBe("agent-col");
|
||||||
|
// Model precedence input is the column agent's runtimeConfig (not the
|
||||||
|
// assigned agent's).
|
||||||
|
expect(opts.assignedAgentRuntimeConfig).toEqual(makeColumnAgent().runtimeConfig);
|
||||||
|
expect(agentStore.getAgent).toHaveBeenCalledWith("agent-col");
|
||||||
|
expect(mockExecuteAll).toHaveBeenCalled();
|
||||||
|
expect(
|
||||||
|
loggedLines(store).some((l) => l.includes("running as column agent 'agent-col' (override)")),
|
||||||
|
).toBe(true);
|
||||||
|
});
|
||||||
|
|
||||||
|
it("defer column with task own complete model pair → step session keeps assigned-agent attribution", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({
|
||||||
|
assignedAgentId: "agent-Y",
|
||||||
|
modelProvider: "task-prov",
|
||||||
|
modelId: "task-model",
|
||||||
|
});
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor, agentStore } = makeExecutor(store, {
|
||||||
|
"agent-Y": makeAssignedAgent(),
|
||||||
|
"agent-col": makeColumnAgent(),
|
||||||
|
});
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
await runStepSessionSeam(executor, task, "foreach-1#0:step-exec", DEFER_COL);
|
||||||
|
|
||||||
|
const opts = lastStepExecutorOpts();
|
||||||
|
// Own settings (complete model pair) suppress the defer column agent.
|
||||||
|
expect(opts.effectiveAgentId).toBeUndefined();
|
||||||
|
expect(opts.assignedAgentRuntimeConfig).toEqual(makeAssignedAgent().runtimeConfig);
|
||||||
|
expect(agentStore.getAgent).not.toHaveBeenCalledWith("agent-col");
|
||||||
|
});
|
||||||
|
|
||||||
|
it("column agent missing from registry at step-execute seam → fallback to assigned-agent attribution, logged", async () => {
|
||||||
|
const store = createMockStore();
|
||||||
|
const task = singleSessionTask({ assignedAgentId: "agent-Y" });
|
||||||
|
store.getTask.mockResolvedValue(task as any);
|
||||||
|
const { executor } = makeExecutor(store, { "agent-Y": makeAssignedAgent() });
|
||||||
|
installTaskDoneAgent();
|
||||||
|
|
||||||
|
await runStepSessionSeam(executor, task, "foreach-1#0:step-exec", OVERRIDE_COL);
|
||||||
|
|
||||||
|
const opts = lastStepExecutorOpts();
|
||||||
|
expect(opts.effectiveAgentId).toBeUndefined();
|
||||||
|
expect(opts.assignedAgentRuntimeConfig).toEqual(makeAssignedAgent().runtimeConfig);
|
||||||
|
expect(
|
||||||
|
loggedLines(store).some(
|
||||||
|
(l) => l.includes("column agent 'agent-col' not found") && l.includes("falling back"),
|
||||||
|
),
|
||||||
|
).toBe(true);
|
||||||
|
});
|
||||||
|
});
|
||||||
|
});
|
||||||
@@ -28,6 +28,7 @@ import type {
|
|||||||
import { observeWorkflowParity, WORKFLOW_INTERPRETER_DUAL_OBSERVE_FLAG } from "./workflow-parity-observer.js";
|
import { observeWorkflowParity, WORKFLOW_INTERPRETER_DUAL_OBSERVE_FLAG } from "./workflow-parity-observer.js";
|
||||||
import {
|
import {
|
||||||
FOREACH_ACTIVE_CONTEXT_KEY,
|
FOREACH_ACTIVE_CONTEXT_KEY,
|
||||||
|
SEAM_GOVERNING_NODE_CONTEXT_KEY,
|
||||||
type ForeachActiveContext,
|
type ForeachActiveContext,
|
||||||
type WorkflowLegacySeams,
|
type WorkflowLegacySeams,
|
||||||
} from "./workflow-node-handlers.js";
|
} from "./workflow-node-handlers.js";
|
||||||
@@ -3259,6 +3260,22 @@ export class TaskExecutor {
|
|||||||
return `${taskId}:${instanceId}`;
|
return `${taskId}:${instanceId}`;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Column-agent seam wiring (column-agent plan U4, R2/R3/R4). Per-run binding
|
||||||
|
* resolver keyed by task id: maps a governing node id to its column-agent
|
||||||
|
* binding (if any), computed once per run in maybeExecuteWorkflowGraph from the
|
||||||
|
* resolved IR. The execute / step-execute seams consume it to decide whether the
|
||||||
|
* coding/step session runs as a column agent. Cleared in the run's finally. */
|
||||||
|
private graphColumnAgentResolver = new Map<string, (nodeId: string) => WorkflowColumnAgent | undefined>();
|
||||||
|
|
||||||
|
/** Column-agent seam wiring (column-agent plan U4). The governing graph node id
|
||||||
|
* for the implementation pass currently in flight for a task — the execute-seam
|
||||||
|
* prompt node's id (execute seam), or the foreach instance node id (step-execute
|
||||||
|
* seam, which the core resolver maps through template inheritance). Stamped by
|
||||||
|
* the seam from the reserved {@link SEAM_GOVERNING_NODE_CONTEXT_KEY} context key
|
||||||
|
* right before it drives the implementation phase, read inside execute()'s
|
||||||
|
* session build, and cleared by the seam afterward. Keyed by task id. */
|
||||||
|
private graphSeamGoverningNodeId = new Map<string, string>();
|
||||||
|
|
||||||
/** Tasks currently being orchestrated by the graph runner. Process-wide for
|
/** Tasks currently being orchestrated by the graph runner. Process-wide for
|
||||||
* the same reason as executingTaskLock (FN-4811): duplicate execute()
|
* the same reason as executingTaskLock (FN-4811): duplicate execute()
|
||||||
* invocations can arrive from different TaskExecutor instances in one
|
* invocations can arrive from different TaskExecutor instances in one
|
||||||
@@ -3334,6 +3351,11 @@ export class TaskExecutor {
|
|||||||
}
|
}
|
||||||
const resolveBindingForNode = (nodeId: string): WorkflowColumnAgent | undefined =>
|
const resolveBindingForNode = (nodeId: string): WorkflowColumnAgent | undefined =>
|
||||||
columnAgentIr ? resolveColumnAgentBinding(columnAgentIr, nodeId) : undefined;
|
columnAgentIr ? resolveColumnAgentBinding(columnAgentIr, nodeId) : undefined;
|
||||||
|
// Column-agent seam wiring (U4): expose the same per-run resolver to the
|
||||||
|
// execute / step-execute seams (which key off a governing node id stamped
|
||||||
|
// into context), so the coding/step session runs as the column agent under
|
||||||
|
// the SAME binding lookup the custom-node seam uses (KTD-2 single resolver).
|
||||||
|
this.graphColumnAgentResolver.set(task.id, resolveBindingForNode);
|
||||||
|
|
||||||
const runner = new WorkflowGraphTaskRunner({
|
const runner = new WorkflowGraphTaskRunner({
|
||||||
store: this.store,
|
store: this.store,
|
||||||
@@ -3401,6 +3423,10 @@ export class TaskExecutor {
|
|||||||
// Clear per-run step-inversion pins (KTD-8: pinned only for the run's life).
|
// Clear per-run step-inversion pins (KTD-8: pinned only for the run's life).
|
||||||
this.graphStepSessionPinned.delete(task.id);
|
this.graphStepSessionPinned.delete(task.id);
|
||||||
this.graphStepRunOnce.delete(task.id);
|
this.graphStepRunOnce.delete(task.id);
|
||||||
|
// Clear per-run column-agent seam wiring (U4): the resolver and any dangling
|
||||||
|
// governing-node-id are scoped to this run only.
|
||||||
|
this.graphColumnAgentResolver.delete(task.id);
|
||||||
|
this.graphSeamGoverningNodeId.delete(task.id);
|
||||||
// Per-instance keys: clear every instance slot owned by this task.
|
// Per-instance keys: clear every instance slot owned by this task.
|
||||||
const ctxPrefix = `${task.id}:`;
|
const ctxPrefix = `${task.id}:`;
|
||||||
for (const key of this.graphStepActiveContext.keys()) {
|
for (const key of this.graphStepActiveContext.keys()) {
|
||||||
@@ -4157,8 +4183,22 @@ export class TaskExecutor {
|
|||||||
// so planning is a no-op for already-specified tasks. Custom planning
|
// so planning is a no-op for already-specified tasks. Custom planning
|
||||||
// behavior is expressed as a custom prompt node before the execute seam.
|
// behavior is expressed as a custom prompt node before the execute seam.
|
||||||
planning: async () => ({ outcome: "success", value: "pre-specified" }),
|
planning: async () => ({ outcome: "success", value: "pre-specified" }),
|
||||||
execute: async (seamTask) => {
|
execute: async (seamTask, context) => {
|
||||||
const result = await this.runImplementationPhase(seamTask);
|
// Column-agent seam wiring (U4, R4): record the governing node id (the
|
||||||
|
// execute-seam prompt node, stamped into context by createPromptLikeHandler)
|
||||||
|
// so execute()'s session build can resolve the column-agent binding for the
|
||||||
|
// node's DECLARED column. Cleared after the pass so a later seam without a
|
||||||
|
// binding cannot inherit a stale node id.
|
||||||
|
const governingNodeId = context?.[SEAM_GOVERNING_NODE_CONTEXT_KEY];
|
||||||
|
if (typeof governingNodeId === "string") {
|
||||||
|
this.graphSeamGoverningNodeId.set(seamTask.id, governingNodeId);
|
||||||
|
}
|
||||||
|
let result: { taskDone: boolean; modifiedFiles: string[] };
|
||||||
|
try {
|
||||||
|
result = await this.runImplementationPhase(seamTask);
|
||||||
|
} finally {
|
||||||
|
this.graphSeamGoverningNodeId.delete(seamTask.id);
|
||||||
|
}
|
||||||
if (result.taskDone) {
|
if (result.taskDone) {
|
||||||
return { outcome: "success", value: "implemented" };
|
return { outcome: "success", value: "implemented" };
|
||||||
}
|
}
|
||||||
@@ -4236,6 +4276,17 @@ export class TaskExecutor {
|
|||||||
// Stamp the active instance so `runGraphTaskStep` can honor
|
// Stamp the active instance so `runGraphTaskStep` can honor
|
||||||
// `deferDoneToReview` when judging a non-terminal step (FIX 3).
|
// `deferDoneToReview` when judging a non-terminal step (FIX 3).
|
||||||
this.graphStepActiveContext.set(this.graphActiveContextKey(seamTask.id, active.instanceId), active);
|
this.graphStepActiveContext.set(this.graphActiveContextKey(seamTask.id, active.instanceId), active);
|
||||||
|
// Column-agent seam wiring (U4, R4): record the governing node id — the
|
||||||
|
// foreach INSTANCE node id (`<foreachId>#<i>:<templateNodeId>`) stamped into
|
||||||
|
// context by createPromptLikeHandler — so the (once-per-run, KTD-2/KTD-8)
|
||||||
|
// step-session implementation pass resolves the column-agent binding for the
|
||||||
|
// step-execute node's effective column (template-node column, else inherited
|
||||||
|
// foreach column). All instances share the same template node and thus the
|
||||||
|
// same binding, so the first instance to drive the pass sets it correctly.
|
||||||
|
const stepGoverningNodeId = context[SEAM_GOVERNING_NODE_CONTEXT_KEY];
|
||||||
|
if (typeof stepGoverningNodeId === "string" && !this.graphSeamGoverningNodeId.has(seamTask.id)) {
|
||||||
|
this.graphSeamGoverningNodeId.set(seamTask.id, stepGoverningNodeId);
|
||||||
|
}
|
||||||
const result = await runTaskStep(
|
const result = await runTaskStep(
|
||||||
{
|
{
|
||||||
store: this.store,
|
store: this.store,
|
||||||
@@ -4570,6 +4621,75 @@ export class TaskExecutor {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Resolve the effective COLUMN AGENT governing the coding/step session currently
|
||||||
|
* being built for a task (column-agent plan U4, R2/R3/R4/R8).
|
||||||
|
*
|
||||||
|
* Reads the governing node id stamped by the active seam ({@link
|
||||||
|
* graphSeamGoverningNodeId}) and the per-run binding resolver ({@link
|
||||||
|
* graphColumnAgentResolver}), both scoped to a graph-owned run. Feeds the task's
|
||||||
|
* OWN settings (`assignedAgentId` + complete `modelProvider`/`modelId` pair) into
|
||||||
|
* the shared core resolver (`resolveEffectiveAgent`, KTD-2/KTD-5) so defer/override
|
||||||
|
* precedence is never reimplemented here. When the verdict is `column-agent`,
|
||||||
|
* fetches the full Agent best-effort and audits the adoption; on a missing/deleted
|
||||||
|
* agent it logs and returns undefined so the caller falls back to the
|
||||||
|
* `assignedAgentId` path (R8). Returns undefined for the legacy/no-binding path so
|
||||||
|
* the session build is byte-identical (characterization parity).
|
||||||
|
*
|
||||||
|
* Exposes the resolved Agent object (not just an id) so U5 can consume the same
|
||||||
|
* effective principal for gating/heartbeat/restart without re-resolving.
|
||||||
|
*/
|
||||||
|
private async resolveSeamColumnAgent(
|
||||||
|
task: Task,
|
||||||
|
detail: TaskDetail,
|
||||||
|
): Promise<{ agent: Agent; mode: WorkflowColumnAgent["mode"] | undefined } | undefined> {
|
||||||
|
const governingNodeId = this.graphSeamGoverningNodeId.get(task.id);
|
||||||
|
const resolveBinding = this.graphColumnAgentResolver.get(task.id);
|
||||||
|
if (!governingNodeId || !resolveBinding) return undefined;
|
||||||
|
|
||||||
|
const binding = resolveBinding(governingNodeId);
|
||||||
|
if (!binding) return undefined;
|
||||||
|
|
||||||
|
// The task's OWN settings: its assigned agent identity and a COMPLETE model
|
||||||
|
// pair (an incomplete pair does not count — KTD-5, mirrors
|
||||||
|
// resolveExecutorSessionModel's both-present rule).
|
||||||
|
const ownAgentId = typeof detail.assignedAgentId === "string" && detail.assignedAgentId.trim()
|
||||||
|
? detail.assignedAgentId.trim()
|
||||||
|
: undefined;
|
||||||
|
const ownModelComplete = Boolean(detail.modelProvider && detail.modelId);
|
||||||
|
const effective = resolveEffectiveAgent({
|
||||||
|
binding,
|
||||||
|
ownAgentId,
|
||||||
|
ownModelProvider: ownModelComplete ? detail.modelProvider : undefined,
|
||||||
|
ownModelId: ownModelComplete ? detail.modelId : undefined,
|
||||||
|
});
|
||||||
|
if (effective.source !== "column-agent") return undefined;
|
||||||
|
|
||||||
|
// Column agent governs: fetch the full Agent (best-effort, R8 fallback).
|
||||||
|
let agent: Agent | null = null;
|
||||||
|
try {
|
||||||
|
agent = (await this.options.agentStore?.getAgent(effective.agentId)) ?? null;
|
||||||
|
} catch {
|
||||||
|
agent = null;
|
||||||
|
}
|
||||||
|
if (!agent) {
|
||||||
|
await this.store.logEntry(
|
||||||
|
task.id,
|
||||||
|
`Workflow seam node '${governingNodeId}': column agent '${effective.agentId}' not found — falling back to assigned-agent resolution`,
|
||||||
|
undefined,
|
||||||
|
this.getRunContextFor(task.id),
|
||||||
|
);
|
||||||
|
return undefined;
|
||||||
|
}
|
||||||
|
await this.store.logEntry(
|
||||||
|
task.id,
|
||||||
|
`Workflow seam node '${governingNodeId}': running as column agent '${effective.agentId}' (${binding.mode})`,
|
||||||
|
undefined,
|
||||||
|
this.getRunContextFor(task.id),
|
||||||
|
);
|
||||||
|
return { agent, mode: binding.mode };
|
||||||
|
}
|
||||||
|
|
||||||
/** Run a custom (non-seam) graph node on the proven WorkflowStep machinery.
|
/** Run a custom (non-seam) graph node on the proven WorkflowStep machinery.
|
||||||
*
|
*
|
||||||
* `columnBinding` (plan U3) is the agent binding governing this node's
|
* `columnBinding` (plan U3) is the agent binding governing this node's
|
||||||
@@ -5289,7 +5409,17 @@ export class TaskExecutor {
|
|||||||
const stepSessionAgent = detail.assignedAgentId && this.options.agentStore
|
const stepSessionAgent = detail.assignedAgentId && this.options.agentStore
|
||||||
? await this.options.agentStore.getAgent(detail.assignedAgentId).catch(() => null)
|
? await this.options.agentStore.getAgent(detail.assignedAgentId).catch(() => null)
|
||||||
: null;
|
: null;
|
||||||
const stepSessionRuntimeHint = extractRuntimeHint(stepSessionAgent?.runtimeConfig);
|
|
||||||
|
// Column-agent SESSION IDENTITY (U4, R2/R3/R4/R8): when the governing
|
||||||
|
// step-execute node's declared column binds an agent that supersedes the
|
||||||
|
// task's assigned agent, the per-step session's MODEL, runtime hint, and
|
||||||
|
// attribution adopt the column agent. The core resolver decides defer vs
|
||||||
|
// override (KTD-2); a missing agent logs + falls back (R8). Gating contexts
|
||||||
|
// still key off the ASSIGNED agent here — principal substitution for
|
||||||
|
// gating/heartbeat is U5 (kept out of this unit deliberately).
|
||||||
|
const stepColumnAgent = await this.resolveSeamColumnAgent(task, detail);
|
||||||
|
const stepIdentityAgent = stepColumnAgent?.agent ?? stepSessionAgent;
|
||||||
|
const stepSessionRuntimeHint = extractRuntimeHint(stepIdentityAgent?.runtimeConfig);
|
||||||
|
|
||||||
let accumulatedStepTokenUsage = detail.tokenUsage;
|
let accumulatedStepTokenUsage = detail.tokenUsage;
|
||||||
const tokenUsageRecordedSteps = new Set<number>();
|
const tokenUsageRecordedSteps = new Set<number>();
|
||||||
@@ -5304,7 +5434,10 @@ export class TaskExecutor {
|
|||||||
stuckTaskDetector: this.options.stuckTaskDetector,
|
stuckTaskDetector: this.options.stuckTaskDetector,
|
||||||
pluginRunner: this.options.pluginRunner,
|
pluginRunner: this.options.pluginRunner,
|
||||||
runtimeHint: stepSessionRuntimeHint,
|
runtimeHint: stepSessionRuntimeHint,
|
||||||
assignedAgentRuntimeConfig: (stepSessionAgent?.runtimeConfig ?? undefined) as Record<string, unknown> | undefined,
|
assignedAgentRuntimeConfig: (stepIdentityAgent?.runtimeConfig ?? undefined) as Record<string, unknown> | undefined,
|
||||||
|
// Attribute the per-step run auditor to the column agent when it governs
|
||||||
|
// (U4); absent → StepSessionExecutor falls back to assignedAgentId.
|
||||||
|
effectiveAgentId: stepColumnAgent?.agent.id,
|
||||||
actionGateContext: this.buildActionGateContext(task.id, stepSessionAgent, settings.defaultAgentPermissionPolicy),
|
actionGateContext: this.buildActionGateContext(task.id, stepSessionAgent, settings.defaultAgentPermissionPolicy),
|
||||||
permanentAgentGating: this.buildPermanentAgentGatingContext(task.id, stepSessionAgent, settings.defaultAgentPermissionPolicy),
|
permanentAgentGating: this.buildPermanentAgentGatingContext(task.id, stepSessionAgent, settings.defaultAgentPermissionPolicy),
|
||||||
// Pass skill selection context from the main executor session
|
// Pass skill selection context from the main executor session
|
||||||
@@ -5784,7 +5917,18 @@ export class TaskExecutor {
|
|||||||
const assignedAgent = assignedAgentId && this.options.agentStore
|
const assignedAgent = assignedAgentId && this.options.agentStore
|
||||||
? await this.options.agentStore.getAgent(assignedAgentId).catch(() => null)
|
? await this.options.agentStore.getAgent(assignedAgentId).catch(() => null)
|
||||||
: null;
|
: null;
|
||||||
const executorRuntimeHint = extractRuntimeHint(assignedAgent?.runtimeConfig);
|
|
||||||
|
// Column-agent SESSION IDENTITY (U4, R2/R3/R4/R8): when the governing execute
|
||||||
|
// seam node's declared column binds an agent that supersedes the task's
|
||||||
|
// assigned agent, the coding session's MODEL, runtime hint, persona, and
|
||||||
|
// memory tools adopt the column agent. The core resolver decides defer vs
|
||||||
|
// override (KTD-2); a missing agent logs + falls back (R8). No binding →
|
||||||
|
// `columnAgentSeam` is undefined and every line below is byte-identical to the
|
||||||
|
// assigned-agent path (characterization parity). Gating contexts still key off
|
||||||
|
// the ASSIGNED agent — principal substitution for gating/heartbeat is U5.
|
||||||
|
const columnAgentSeam = await this.resolveSeamColumnAgent(task, detail);
|
||||||
|
const identityAgent = columnAgentSeam?.agent ?? assignedAgent;
|
||||||
|
const executorRuntimeHint = extractRuntimeHint(identityAgent?.runtimeConfig);
|
||||||
|
|
||||||
// Log fast mode status
|
// Log fast mode status
|
||||||
if (executionMode === "fast") {
|
if (executionMode === "fast") {
|
||||||
@@ -5831,11 +5975,11 @@ export class TaskExecutor {
|
|||||||
})
|
})
|
||||||
: []),
|
: []),
|
||||||
createWebFetchTool(),
|
createWebFetchTool(),
|
||||||
...createMemoryTools(this.rootDir, settings, assignedAgent ? {
|
...createMemoryTools(this.rootDir, settings, identityAgent ? {
|
||||||
agentMemory: {
|
agentMemory: {
|
||||||
agentId: assignedAgent.id,
|
agentId: identityAgent.id,
|
||||||
agentName: assignedAgent.name,
|
agentName: identityAgent.name,
|
||||||
memory: assignedAgent.memory,
|
memory: identityAgent.memory,
|
||||||
},
|
},
|
||||||
} : undefined),
|
} : undefined),
|
||||||
// Conditionally add agent self-reflection when enabled and task has an assigned agent.
|
// Conditionally add agent self-reflection when enabled and task has an assigned agent.
|
||||||
@@ -5890,11 +6034,14 @@ export class TaskExecutor {
|
|||||||
// 3. Global execution lane pair (executionGlobalProvider + executionGlobalModelId)
|
// 3. Global execution lane pair (executionGlobalProvider + executionGlobalModelId)
|
||||||
// 4. Project default override pair (defaultProviderOverride + defaultModelIdOverride)
|
// 4. Project default override pair (defaultProviderOverride + defaultModelIdOverride)
|
||||||
// 5. Global default pair (defaultProvider + defaultModelId)
|
// 5. Global default pair (defaultProvider + defaultModelId)
|
||||||
|
// Column-agent session identity (U4): the model precedence input is the
|
||||||
|
// EFFECTIVE identity agent's runtimeConfig (column agent when it governs,
|
||||||
|
// else the assigned agent — byte-identical no-binding path).
|
||||||
const { provider: executorProvider, modelId: executorModelId } = resolveExecutorSessionModel(
|
const { provider: executorProvider, modelId: executorModelId } = resolveExecutorSessionModel(
|
||||||
detail.modelProvider,
|
detail.modelProvider,
|
||||||
detail.modelId,
|
detail.modelId,
|
||||||
settings,
|
settings,
|
||||||
(assignedAgent?.runtimeConfig ?? undefined) as Record<string, unknown> | undefined,
|
(identityAgent?.runtimeConfig ?? undefined) as Record<string, unknown> | undefined,
|
||||||
);
|
);
|
||||||
const executorFallbackProvider = settings.fallbackProvider;
|
const executorFallbackProvider = settings.fallbackProvider;
|
||||||
const executorFallbackModelId = settings.fallbackModelId;
|
const executorFallbackModelId = settings.fallbackModelId;
|
||||||
@@ -5928,8 +6075,15 @@ export class TaskExecutor {
|
|||||||
|
|
||||||
executorLog.log(`${task.id}: creating agent session (provider=${executorProvider ?? "default"}, model=${executorModelId ?? "default"}, resuming=${isResuming})`);
|
executorLog.log(`${task.id}: creating agent session (provider=${executorProvider ?? "default"}, model=${executorModelId ?? "default"}, resuming=${isResuming})`);
|
||||||
|
|
||||||
// Resolve per-agent custom instructions for the executor role
|
// Resolve per-agent custom instructions for the executor role.
|
||||||
const executorInstructions = await this.resolveInstructionsForRole("executor", settings);
|
// Column-agent session identity (U4, R3/KTD-6): when a column agent governs,
|
||||||
|
// its TYPED persona (soul/instructionsText, via buildAgentPersona — the same
|
||||||
|
// source the custom-node path uses) supersedes the role-resolved executor
|
||||||
|
// instructions, so the coding session speaks AS the column agent. No binding
|
||||||
|
// → role instructions unchanged (characterization parity).
|
||||||
|
const columnAgentPersona = columnAgentSeam ? this.buildAgentPersona(columnAgentSeam.agent) : undefined;
|
||||||
|
const executorInstructions = columnAgentPersona
|
||||||
|
?? (await this.resolveInstructionsForRole("executor", settings));
|
||||||
|
|
||||||
// Build structured layers for cross-session prompt caching.
|
// Build structured layers for cross-session prompt caching.
|
||||||
const executorPluginContributions = buildPluginPromptSection(
|
const executorPluginContributions = buildPluginPromptSection(
|
||||||
|
|||||||
@@ -125,6 +125,18 @@ export interface StepSessionExecutorOptions {
|
|||||||
permanentAgentGating?: PermanentAgentGatingContext;
|
permanentAgentGating?: PermanentAgentGatingContext;
|
||||||
/** Task-scoped environment injected into non-git subprocesses. */
|
/** Task-scoped environment injected into non-git subprocesses. */
|
||||||
taskEnv?: NodeJS.ProcessEnv;
|
taskEnv?: NodeJS.ProcessEnv;
|
||||||
|
/**
|
||||||
|
* Column-agent identity override for session attribution (column-agent plan U4,
|
||||||
|
* R2/R3/R4). When the governing foreach/step-execute node's declared column
|
||||||
|
* binds an agent that supersedes the task's `assignedAgentId` (override, or
|
||||||
|
* defer with no own settings), the executor passes the column agent's id here so
|
||||||
|
* the per-step run auditor attributes the session to who actually ran — not
|
||||||
|
* `taskDetail.assignedAgentId`. Absent → attribution falls back to
|
||||||
|
* `taskDetail.assignedAgentId ?? "executor"` (byte-identical legacy path). The
|
||||||
|
* column agent's MODEL flows separately via {@link assignedAgentRuntimeConfig}
|
||||||
|
* (the executor swaps it to the column agent's `runtimeConfig` at the seam).
|
||||||
|
*/
|
||||||
|
effectiveAgentId?: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
// ── File Scope Extraction ─────────────────────────────────────────────
|
// ── File Scope Extraction ─────────────────────────────────────────────
|
||||||
@@ -1018,7 +1030,10 @@ Follow instructions precisely and avoid unrelated changes.`,
|
|||||||
defaultThinkingLevel: taskDetail.thinkingLevel ?? settings.defaultThinkingLevel,
|
defaultThinkingLevel: taskDetail.thinkingLevel ?? settings.defaultThinkingLevel,
|
||||||
runAuditor: createRunAuditor(this.store, {
|
runAuditor: createRunAuditor(this.store, {
|
||||||
runId: generateSyntheticRunId("workflow-step", taskDetail.id),
|
runId: generateSyntheticRunId("workflow-step", taskDetail.id),
|
||||||
agentId: taskDetail.assignedAgentId ?? "executor",
|
// Column-agent attribution (U4): the effective column agent is the
|
||||||
|
// principal that actually ran when the seam node's column governs;
|
||||||
|
// fall back to the task's assigned agent (legacy, byte-identical).
|
||||||
|
agentId: this.options.effectiveAgentId ?? taskDetail.assignedAgentId ?? "executor",
|
||||||
taskId: taskDetail.id,
|
taskId: taskDetail.id,
|
||||||
taskLineageId: taskDetail.lineageId,
|
taskLineageId: taskDetail.lineageId,
|
||||||
phase: "execute",
|
phase: "execute",
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
import { WorkflowIrError, getStepParser } from "@fusion/core";
|
import { WorkflowIrError, getStepParser, instanceNodeId } from "@fusion/core";
|
||||||
import type { TaskDetail, TaskStep, WorkflowIrNode } from "@fusion/core";
|
import type { TaskDetail, TaskStep, WorkflowIrNode } from "@fusion/core";
|
||||||
|
|
||||||
import type { WorkflowNodeHandler, WorkflowNodeResult } from "./workflow-graph-executor.js";
|
import type { WorkflowNodeHandler, WorkflowNodeResult } from "./workflow-graph-executor.js";
|
||||||
@@ -71,6 +71,19 @@ export interface StepReviewSeamResult {
|
|||||||
* which step they operate on and the per-instance baseline/checkpoint state. */
|
* which step they operate on and the per-instance baseline/checkpoint state. */
|
||||||
export const FOREACH_ACTIVE_CONTEXT_KEY = "foreach:active";
|
export const FOREACH_ACTIVE_CONTEXT_KEY = "foreach:active";
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Reserved context key carrying the GOVERNING graph node id into the legacy
|
||||||
|
* coding seams (column-agent plan U4, R4). The execute seam reads the seam node's
|
||||||
|
* own id; the step-execute seam reads the foreach INSTANCE node id
|
||||||
|
* (`<foreachId>#<i>:<templateNodeId>`) so the core column-agent resolver can map
|
||||||
|
* it through template inheritance to the governing column's binding. The seam
|
||||||
|
* stamps it into a per-run executor slot before driving the implementation pass,
|
||||||
|
* so the binding the session runs under keys off the node's DECLARED IR column
|
||||||
|
* — never the task's current board lane. Custom (non-seam) nodes never use this:
|
||||||
|
* runGraphCustomNode receives its binding directly as a parameter (U3).
|
||||||
|
*/
|
||||||
|
export const SEAM_GOVERNING_NODE_CONTEXT_KEY = "workflow:seam-governing-node-id";
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Reserved context marker set by the split sub-walk (`runSplitJoin`) for the
|
* Reserved context marker set by the split sub-walk (`runSplitJoin`) for the
|
||||||
* duration of its branches' execution and cleared at the join (KTD-4, U5). A
|
* duration of its branches' execution and cleared at the join (KTD-4, U5). A
|
||||||
@@ -179,9 +192,24 @@ export function createPromptLikeHandler(
|
|||||||
// succeed — that would merge a task with no step work done.
|
// succeed — that would merge a task with no step work done.
|
||||||
return { outcome: "failure", value: "step-execute-unwired" };
|
return { outcome: "failure", value: "step-execute-unwired" };
|
||||||
}
|
}
|
||||||
|
// Column-agent seam wiring (U4, R4): the GOVERNING node for a step-execute
|
||||||
|
// session is the foreach INSTANCE node id, so the core resolver can map it
|
||||||
|
// through template inheritance to the enclosing foreach's bound column (or
|
||||||
|
// the template node's own column when it declares one). The template node id
|
||||||
|
// is THIS node's id; the foreach node id + step index come from the active
|
||||||
|
// instance context. Stamped so the seam threads it into the session build.
|
||||||
|
context.context[SEAM_GOVERNING_NODE_CONTEXT_KEY] = instanceNodeId(
|
||||||
|
active.foreachNodeId,
|
||||||
|
active.stepIndex,
|
||||||
|
node.id,
|
||||||
|
);
|
||||||
return seams.stepExecute(context.task, context.context);
|
return seams.stepExecute(context.task, context.context);
|
||||||
}
|
}
|
||||||
if (seam) {
|
if (seam) {
|
||||||
|
// Column-agent seam wiring (U4, R4): for the execute seam the governing node
|
||||||
|
// IS the seam node, so its declared column drives the binding. (Other seams
|
||||||
|
// — planning/review/merge/schedule — stamp it too; only execute reads it.)
|
||||||
|
context.context[SEAM_GOVERNING_NODE_CONTEXT_KEY] = node.id;
|
||||||
return seams[seam]!(context.task, context.context);
|
return seams[seam]!(context.task, context.context);
|
||||||
}
|
}
|
||||||
if (!runCustomNode) {
|
if (!runCustomNode) {
|
||||||
|
|||||||
Reference in New Issue
Block a user