Route workflow stages through task-scoped durable role agents. - Persist normalized multi-role agents and workflow principal fences with migrations. - Route planning, execution, review, and merge workflow nodes through authorized permanent principals with capacity leasing and recovery. - Retire ephemeral workflow-stage workers and expose role-aware agent configuration, workflow editing, and documentation. - Preserve lifecycle-column ratchet coverage by centralizing workflow-role classification rather than adding test exemptions. Files changed: .changeset/fn-8764-workflow-role-agents.md | 7 + CONCEPTS.md | 3 + docs/agents.md | 6 + docs/architecture.md | 6 + docs/cli-reference.md | 2 + docs/dashboard-guide.md | 4 + docs/settings-reference.md | 6 +- docs/storage.md | 2 + docs/workflow-steps.md | 6 + .../src/__tests__/extension-agent-update.test.ts | 11 +- packages/cli/src/__tests__/extension.test.ts | 18 +- packages/cli/src/extension.ts | 41 +- .../core/src/__tests__/agent-permissions.test.ts | 12 + .../core/src/__tests__/agent-role-policy.test.ts | 7 + packages/core/src/__tests__/agent-roles.test.ts | 21 + .../legacy-column-collection-gating-ledger.test.ts | 19 +- .../src/__tests__/postgres/schema-applier.test.ts | 16 +- .../core/src/__tests__/settings-parity.test.ts | 9 +- .../workflow-agent-node-classification.test.ts | 25 + .../src/__tests__/workflow-work-item-cas.test.ts | 38 ++ packages/core/src/agents/agent-permissions.ts | 11 +- packages/core/src/agents/agent-role-policy.ts | 39 +- packages/core/src/agents/agent-store.ts | 190 ++++++- .../core/src/async-stores/async-agent-store.ts | 6 + packages/core/src/config/settings-schema.ts | 5 +- packages/core/src/index.gate.ts | 2 +- packages/core/src/index.ts | 7 +- .../0045_fn_8764_multi_role_workflow_agents.sql | 20 + .../0046_fn_8764_workflow_principal_fence.sql | 49 ++ packages/core/src/postgres/schema-applier.ts | 22 +- packages/core/src/postgres/schema/project.ts | 21 + packages/core/src/store.ts | 2 +- .../task-store/async/async-workflow-workitems.ts | 49 +- packages/core/src/task-store/row-types.ts | 4 + packages/core/src/task-store/settings-helpers.ts | 16 +- packages/core/src/task-store/settings-ops-2.ts | 13 +- packages/core/src/task-store/settings-ops.ts | 16 +- packages/core/src/task-store/task-row-mappers.ts | 4 + .../src/task-store/workflow-task-create-ops.ts | 6 +- .../src/task-store/workflow-workitems-ops-2.ts | 25 +- packages/core/src/types.ts | 2 + packages/core/src/types/agents/agents.ts | 45 +- packages/core/src/types/merge/merge-queue.ts | 17 + packages/core/src/types/settings/settings-scope.ts | 9 +- packages/core/src/workflows/workflow-ir-types.ts | 58 +++ packages/core/src/workflows/workflow-ir.ts | 19 + .../dashboard/app/components/AgentDetailView.css | 14 + .../dashboard/app/components/AgentDetailView.tsx | 34 +- .../dashboard/app/components/NewAgentDialog.tsx | 28 +- .../app/components/WorkflowNodeEditor.tsx | 19 + .../__tests__/AgentDetailView.core.test.tsx | 4 +- .../app/components/__tests__/AgentsView.test.tsx | 2 +- .../__tests__/SettingsModal.general.test.tsx | 86 --- .../__tests__/SettingsModal.test-harness.tsx | 1 - .../components/agent-presets/agentCreatePayload.ts | 9 +- .../app/components/settings/section-keys.ts | 1 - .../settings/sections/GeneralSection.tsx | 8 - .../settings-default-descriptions.test.tsx | 1 - .../app/components/workflow-flow-mapping.ts | 7 + packages/dashboard/src/mission-routes.ts | 26 +- .../src/routes/__tests__/agent-core-routes.test.ts | 23 +- .../src/routes/register-agent-core-routes.ts | 42 +- ...gister-agent-import-export-generation-routes.ts | 21 - .../engine/src/__tests__/agent-action-gate.test.ts | 33 ++ .../engine/src/__tests__/agent-assignment.test.ts | 370 ------------- .../src/__tests__/ephemeral-worker-manager.test.ts | 575 --------------------- ...ecutor-ephemeral-disabled-dispatch-gate.test.ts | 223 -------- .../__tests__/executor-fast-mode-workflows.test.ts | 58 +++ .../engine/src/__tests__/log-severity-manifest.ts | 1 - .../__tests__/log-severity-spam-contract.test.ts | 3 - .../resolved-read-with-literal-filter.test.ts | 4 - .../__tests__/scheduler-ephemeral-toggle.test.ts | 175 ------- .../__tests__/scheduler-workflow-cutover.test.ts | 19 - .../src/__tests__/workflow-agent-capacity.test.ts | 47 ++ .../src/__tests__/workflow-agent-routing.test.ts | 137 +++++ .../src/__tests__/workflow-graph-foreach.test.ts | 15 + .../__tests__/workflow-graph-task-runner.test.ts | 73 +++ .../src/__tests__/workflow-task-runtime.test.ts | 95 ++++ .../src/__tests__/workflow-work-scheduler.test.ts | 20 + packages/engine/src/agents/agent-action-gate.ts | 64 +++ packages/engine/src/agents/agent-assignment.ts | 135 ----- packages/engine/src/agents/agent-reflection.ts | 1 + .../engine/src/agents/ephemeral-worker-manager.ts | 429 --------------- .../engine/src/agents/workflow-agent-capacity.ts | 113 ++++ .../engine/src/agents/workflow-agent-router.ts | 185 +++++++ packages/engine/src/execution/reviewer.ts | 26 +- packages/engine/src/executor.ts | 501 +++++++++++++++--- packages/engine/src/index.ts | 1 - packages/engine/src/merger.ts | 20 +- packages/engine/src/pi.ts | 11 + packages/engine/src/runtimes/in-process-runtime.ts | 37 -- packages/engine/src/scheduler.ts | 114 +--- packages/engine/src/triage.ts | 196 ++++++- .../src/workflows/workflow-graph-executor.ts | 109 +++- .../engine/src/workflows/workflow-graph-loop.ts | 13 +- .../src/workflows/workflow-graph-task-runner.ts | 12 + .../engine/src/workflows/workflow-task-runtime.ts | 125 ++++- .../src/workflows/workflow-work-scheduler.ts | 8 +- 98 files changed, 2722 insertions(+), 2468 deletions(-) Fusion-Task-Id: FN-8764 Fusion-Task-Lineage: 5527fccb-342d-46f6-8108-bbf89142efec Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
316 lines
13 KiB
TypeScript
316 lines
13 KiB
TypeScript
import type { WorkflowIrEdge, WorkflowIrNode, WorkflowLoopConfig, WorkflowOptionalGroupConfig } from "@fusion/core";
|
|
import { WorkflowIrError } from "@fusion/core";
|
|
|
|
import type { WorkflowNodeOutcome, WorkflowNodeResult } from "./workflow-graph-executor.js";
|
|
|
|
const DEFAULT_MAX_ITERATIONS = 3;
|
|
const MAX_ITERATIONS_CAP = 50;
|
|
const DEFAULT_TIMEOUT_MS = 300_000;
|
|
const MAX_TIMEOUT_MS = 3_600_000;
|
|
|
|
interface LoopConfig {
|
|
template: { nodes: WorkflowIrNode[]; edges: WorkflowIrEdge[] };
|
|
exitWhen: WorkflowLoopConfig["exitWhen"];
|
|
exitRegex?: RegExp;
|
|
maxIterations: number;
|
|
timeoutMs: number;
|
|
}
|
|
|
|
export interface LoopEnvironment {
|
|
context: Record<string, unknown>;
|
|
runTemplateNode: (
|
|
node: WorkflowIrNode,
|
|
signal?: AbortSignal,
|
|
contextOverride?: Record<string, unknown>,
|
|
) => Promise<WorkflowNodeResult>;
|
|
shouldTraverseEdge: (edge: WorkflowIrEdge, source: WorkflowNodeResult) => boolean;
|
|
signal?: AbortSignal;
|
|
now?: () => number;
|
|
}
|
|
|
|
export interface LoopRunResult {
|
|
outcome: WorkflowNodeOutcome;
|
|
value?: string;
|
|
visitedNodeIds: string[];
|
|
}
|
|
|
|
function resolveLoopConfig(node: WorkflowIrNode): LoopConfig {
|
|
const cfg = (node.config ?? {}) as Partial<WorkflowLoopConfig>;
|
|
if (!cfg.template || !Array.isArray(cfg.template.nodes) || !Array.isArray(cfg.template.edges)) {
|
|
throw new WorkflowIrError(`loop node '${node.id}' has no template subgraph`);
|
|
}
|
|
if (!cfg.exitWhen) {
|
|
throw new WorkflowIrError(`loop node '${node.id}' has no exitWhen condition`);
|
|
}
|
|
const maxIterations =
|
|
typeof cfg.maxIterations === "number" && Number.isFinite(cfg.maxIterations)
|
|
? Math.max(1, Math.min(MAX_ITERATIONS_CAP, Math.floor(cfg.maxIterations)))
|
|
: DEFAULT_MAX_ITERATIONS;
|
|
const timeoutMs =
|
|
typeof cfg.timeoutMs === "number" && Number.isFinite(cfg.timeoutMs)
|
|
? Math.max(1, Math.min(MAX_TIMEOUT_MS, Math.floor(cfg.timeoutMs)))
|
|
: DEFAULT_TIMEOUT_MS;
|
|
const exitRegex =
|
|
cfg.exitWhen.type === "output-matches" ? new RegExp(cfg.exitWhen.pattern, cfg.exitWhen.flags) : undefined;
|
|
return {
|
|
template: cfg.template,
|
|
exitWhen: cfg.exitWhen,
|
|
exitRegex,
|
|
maxIterations,
|
|
timeoutMs,
|
|
};
|
|
}
|
|
|
|
function buildOutgoing(edges: WorkflowIrEdge[]): Map<string, WorkflowIrEdge[]> {
|
|
const outgoing = new Map<string, WorkflowIrEdge[]>();
|
|
for (const edge of edges) {
|
|
const list = outgoing.get(edge.from) ?? [];
|
|
list.push(edge);
|
|
outgoing.set(edge.from, list);
|
|
}
|
|
return outgoing;
|
|
}
|
|
|
|
function findTemplateEntry(nodes: WorkflowIrNode[], edges: WorkflowIrEdge[], loopId: string): WorkflowIrNode {
|
|
const incoming = new Map<string, number>();
|
|
for (const edge of edges) incoming.set(edge.to, (incoming.get(edge.to) ?? 0) + 1);
|
|
const entries = nodes.filter((n) => (incoming.get(n.id) ?? 0) === 0);
|
|
if (entries.length !== 1) {
|
|
throw new WorkflowIrError(`loop node '${loopId}' template must have exactly one entry node`);
|
|
}
|
|
return entries[0];
|
|
}
|
|
|
|
function exitNodeId(nodes: WorkflowIrNode[], edges: WorkflowIrEdge[], loopId: string): string {
|
|
const outgoing = new Map<string, number>();
|
|
for (const edge of edges) outgoing.set(edge.from, (outgoing.get(edge.from) ?? 0) + 1);
|
|
const exits = nodes.filter((n) => (outgoing.get(n.id) ?? 0) === 0);
|
|
if (exits.length !== 1) {
|
|
throw new WorkflowIrError(`loop node '${loopId}' template must have exactly one exit node`);
|
|
}
|
|
return exits[0].id;
|
|
}
|
|
|
|
function matchesExit(config: LoopConfig, value: unknown): boolean {
|
|
const text = typeof value === "string" ? value : value == null ? "" : String(value);
|
|
const condition = config.exitWhen;
|
|
if (condition.type === "output-contains") {
|
|
return text.includes(condition.value);
|
|
}
|
|
return (config.exitRegex ?? new RegExp(condition.pattern, condition.flags)).test(text);
|
|
}
|
|
|
|
function publishIterationContext(
|
|
target: Record<string, unknown>,
|
|
iterationContext: Record<string, unknown>,
|
|
): void {
|
|
const { ["loop:active"]: _active, ...publicContext } = iterationContext;
|
|
Object.assign(target, publicContext);
|
|
}
|
|
|
|
export async function runLoop(
|
|
loopNode: WorkflowIrNode,
|
|
env: LoopEnvironment,
|
|
): Promise<LoopRunResult> {
|
|
const config = resolveLoopConfig(loopNode);
|
|
const templateById = new Map(config.template.nodes.map((n) => [n.id, n]));
|
|
const outgoing = buildOutgoing(config.template.edges);
|
|
const entry = findTemplateEntry(config.template.nodes, config.template.edges, loopNode.id);
|
|
const defaultExitNodeId = exitNodeId(config.template.nodes, config.template.edges, loopNode.id);
|
|
const sourceNodeId = config.exitWhen.nodeId ?? defaultExitNodeId;
|
|
const now = env.now ?? (() => Date.now());
|
|
const deadline = now() + config.timeoutMs;
|
|
const visitedNodeIds: string[] = [];
|
|
const iterationSummaries: Array<{ iteration: number; outcome: string; value?: string }> = [];
|
|
|
|
for (let iteration = 1; iteration <= config.maxIterations; iteration++) {
|
|
if (env.signal?.aborted) {
|
|
return { outcome: "failure", value: "aborted", visitedNodeIds };
|
|
}
|
|
if (now() >= deadline) {
|
|
env.context[`node:${loopNode.id}:loop`] = {
|
|
iterations: iteration - 1,
|
|
exitReason: "timeout",
|
|
history: iterationSummaries,
|
|
};
|
|
return { outcome: "failure", value: "loop-timeout", visitedNodeIds };
|
|
}
|
|
|
|
const iterationContext: Record<string, unknown> = {
|
|
...env.context,
|
|
"loop:active": {
|
|
loopNodeId: loopNode.id,
|
|
iteration,
|
|
},
|
|
};
|
|
let current: WorkflowIrNode | undefined = entry;
|
|
let lastResult: WorkflowNodeResult = { outcome: "success" };
|
|
|
|
while (current) {
|
|
if (env.signal?.aborted) {
|
|
return { outcome: "failure", value: "aborted", visitedNodeIds };
|
|
}
|
|
if (now() >= deadline) {
|
|
env.context[`node:${loopNode.id}:loop`] = {
|
|
iterations: iteration - 1,
|
|
exitReason: "timeout",
|
|
history: iterationSummaries,
|
|
};
|
|
return { outcome: "failure", value: "loop-timeout", visitedNodeIds };
|
|
}
|
|
|
|
const materializedId = `${loopNode.id}#${iteration}:${current.id}`;
|
|
visitedNodeIds.push(materializedId);
|
|
lastResult = await env.runTemplateNode(current, env.signal, iterationContext);
|
|
if (lastResult.contextPatch) Object.assign(iterationContext, lastResult.contextPatch);
|
|
iterationContext[`node:${current.id}:outcome`] = lastResult.outcome;
|
|
if (lastResult.value !== undefined) iterationContext[`node:${current.id}:value`] = lastResult.value;
|
|
if (lastResult.outcome === "failure") {
|
|
publishIterationContext(env.context, iterationContext);
|
|
env.context[`node:${loopNode.id}:loop`] = {
|
|
iterations: iteration,
|
|
exitReason: "node-failure",
|
|
history: iterationSummaries,
|
|
};
|
|
return { outcome: "failure", value: lastResult.value, visitedNodeIds };
|
|
}
|
|
|
|
const edges: WorkflowIrEdge[] = outgoing.get(current.id) ?? [];
|
|
const matching: WorkflowIrEdge[] = edges.filter((edge: WorkflowIrEdge) =>
|
|
env.shouldTraverseEdge(edge, lastResult),
|
|
);
|
|
current = matching.length > 0 ? templateById.get(matching[0].to) : undefined;
|
|
}
|
|
|
|
const sourceValue = iterationContext[`node:${sourceNodeId}:value`];
|
|
const finalValue = sourceValue ?? lastResult.value;
|
|
iterationSummaries.push({
|
|
iteration,
|
|
outcome: lastResult.outcome,
|
|
...(finalValue !== undefined ? { value: String(finalValue) } : {}),
|
|
});
|
|
publishIterationContext(env.context, iterationContext);
|
|
if (matchesExit(config, finalValue)) {
|
|
env.context[`node:${loopNode.id}:loop`] = {
|
|
iterations: iteration,
|
|
exitReason: "matched",
|
|
finalValue,
|
|
history: iterationSummaries,
|
|
};
|
|
return { outcome: "success", visitedNodeIds };
|
|
}
|
|
}
|
|
|
|
env.context[`node:${loopNode.id}:loop`] = {
|
|
iterations: config.maxIterations,
|
|
exitReason: "iteration-exhausted",
|
|
history: iterationSummaries,
|
|
};
|
|
return { outcome: "failure", value: "loop-iteration-exhausted", visitedNodeIds };
|
|
}
|
|
|
|
/*
|
|
FNXC:WorkflowOptionalGroup 2026-06-21-14:05:
|
|
An enabled `optional-group` runs its `template` subgraph EXACTLY ONCE (single pass — no iteration, no rework budget; rework edges are validation-forbidden inside the template). This reuses the loop's template-walk primitives (`buildOutgoing`, `findTemplateEntry`, `shouldTraverseEdge`) but caps the walk at one pass. The disabled/bypass decision lives in the executor branch (read from per-task `enabledWorkflowSteps`); this helper only runs the body when enabled.
|
|
A template-node failure surfaces as the group's outcome so the group's `failure`/`outcome:` edges route, mirroring `runLoop`'s node-failure short-circuit.
|
|
*/
|
|
export interface OptionalGroupEnvironment {
|
|
context: Record<string, unknown>;
|
|
runTemplateNode: (
|
|
node: WorkflowIrNode,
|
|
signal?: AbortSignal,
|
|
contextOverride?: Record<string, unknown>,
|
|
) => Promise<WorkflowNodeResult>;
|
|
shouldTraverseEdge: (edge: WorkflowIrEdge, source: WorkflowNodeResult) => boolean;
|
|
signal?: AbortSignal;
|
|
}
|
|
|
|
export interface OptionalGroupRunResult {
|
|
outcome: WorkflowNodeOutcome;
|
|
value?: string;
|
|
visitedNodeIds: string[];
|
|
/*
|
|
* FNXC:WorkflowStepResults 2026-06-25-12:00:
|
|
* The exit (last-run) inner template node's result, surfaced so the executor can
|
|
* record the group's outcome into `task.workflowStepResults` (KTD-1/KTD-2, plan
|
|
* U2). For a prompt/gate inner node `value` carries the structured verdict
|
|
* (APPROVE / APPROVE_WITH_NOTES / REVISE); `outcome` distinguishes a gate REVISE
|
|
* (failure) from an advisory REVISE (success). Notes/output are not surfaced by
|
|
* the WorkflowNodeResult contract, so they are recorded only when present.
|
|
*/
|
|
exitStepRecord?: WorkflowNodeResult;
|
|
}
|
|
|
|
function resolveOptionalGroupTemplate(
|
|
node: WorkflowIrNode,
|
|
): { nodes: WorkflowIrNode[]; edges: WorkflowIrEdge[] } {
|
|
const cfg = (node.config ?? {}) as Partial<WorkflowOptionalGroupConfig>;
|
|
if (!cfg.template || !Array.isArray(cfg.template.nodes) || !Array.isArray(cfg.template.edges)) {
|
|
throw new WorkflowIrError(`optional-group node '${node.id}' has no template subgraph`);
|
|
}
|
|
return cfg.template;
|
|
}
|
|
|
|
/**
|
|
* Walk an enabled optional-group's template subgraph once. Mirrors a single
|
|
* loop iteration: entry → follow matching edges → stop at the template exit (no
|
|
* outgoing matching edge). Materialized visited ids use a `<groupId>::<templateNodeId>`
|
|
* scheme so they are distinguishable from top-level ids and parseable back to the
|
|
* template node. The group's own outcome is the last template node's outcome.
|
|
*/
|
|
export async function runOptionalGroup(
|
|
groupNode: WorkflowIrNode,
|
|
env: OptionalGroupEnvironment,
|
|
): Promise<OptionalGroupRunResult> {
|
|
const template = resolveOptionalGroupTemplate(groupNode);
|
|
const templateById = new Map(template.nodes.map((n) => [n.id, n]));
|
|
const outgoing = buildOutgoing(template.edges);
|
|
const entry = findTemplateEntry(template.nodes, template.edges, groupNode.id);
|
|
const visitedNodeIds: string[] = [];
|
|
|
|
/*
|
|
* FNXC:WorkflowAgentRouting 2026-08-07-05:37:
|
|
* Keep the enclosing optional-group's materialized identity while its reusable
|
|
* template runs. The graph executor extends this marker for each child so a
|
|
* reviewer override cannot leak between separate group invocations.
|
|
*/
|
|
const groupContext: Record<string, unknown> = {
|
|
...env.context,
|
|
"optional-group:active": typeof env.context["workflow:node-instance-id"] === "string"
|
|
? env.context["workflow:node-instance-id"]
|
|
: groupNode.id,
|
|
};
|
|
let current: WorkflowIrNode | undefined = entry;
|
|
let lastResult: WorkflowNodeResult = { outcome: "success" };
|
|
|
|
while (current) {
|
|
if (env.signal?.aborted) {
|
|
return { outcome: "failure", value: "aborted", visitedNodeIds };
|
|
}
|
|
|
|
const materializedId = `${groupNode.id}::${current.id}`;
|
|
visitedNodeIds.push(materializedId);
|
|
lastResult = await env.runTemplateNode(current, env.signal, groupContext);
|
|
if (lastResult.contextPatch) Object.assign(groupContext, lastResult.contextPatch);
|
|
groupContext[`node:${current.id}:outcome`] = lastResult.outcome;
|
|
if (lastResult.value !== undefined) groupContext[`node:${current.id}:value`] = lastResult.value;
|
|
|
|
if (lastResult.outcome === "failure") {
|
|
// Publish accumulated template context, then surface the failure as the
|
|
// group's outcome so its failure/outcome: edges route.
|
|
Object.assign(env.context, groupContext);
|
|
return { outcome: "failure", value: lastResult.value, visitedNodeIds, exitStepRecord: lastResult };
|
|
}
|
|
|
|
const edges: WorkflowIrEdge[] = outgoing.get(current.id) ?? [];
|
|
const matching: WorkflowIrEdge[] = edges.filter((edge: WorkflowIrEdge) =>
|
|
env.shouldTraverseEdge(edge, lastResult),
|
|
);
|
|
current = matching.length > 0 ? templateById.get(matching[0].to) : undefined;
|
|
}
|
|
|
|
// Single pass complete: publish the template's context onto the shared context.
|
|
Object.assign(env.context, groupContext);
|
|
return { outcome: lastResult.outcome, value: lastResult.value, visitedNodeIds, exitStepRecord: lastResult };
|
|
}
|