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>
589 lines
22 KiB
TypeScript
589 lines
22 KiB
TypeScript
import { describe, expect, it, vi } from "vitest";
|
|
import { EventEmitter } from "node:events";
|
|
import type { Settings, TaskDetail, WorkflowDefinition, WorkflowIr } from "@fusion/core";
|
|
import { createTaskStoreForTest, PG_AVAILABLE } from "../../../core/src/__test-utils__/pg-test-harness.js";
|
|
|
|
import { NotificationService } from "../notification/notification-service.js";
|
|
import { WorkflowGraphTaskRunner, type WorkflowGraphRunnerStore } from "../workflows/workflow-graph-task-runner.js";
|
|
import type { WorkflowNodeResult } from "../workflows/workflow-graph-executor.js";
|
|
|
|
const task = { id: "FN-9001" } as TaskDetail;
|
|
const flagOn = { experimentalFeatures: { workflowGraphExecutor: true } } as unknown as Pick<
|
|
Settings,
|
|
"experimentalFeatures"
|
|
>;
|
|
const flagOff = { experimentalFeatures: { workflowGraphExecutor: false } } as unknown as Pick<
|
|
Settings,
|
|
"experimentalFeatures"
|
|
>;
|
|
|
|
/** start → lint(custom) → execute → review → merge → notify(custom) → end, with seam failure edges to end. */
|
|
function fullLifecycleIr(): WorkflowIr {
|
|
return {
|
|
version: "v1",
|
|
name: "full",
|
|
nodes: [
|
|
{ id: "start", kind: "start" },
|
|
{ id: "lint", kind: "prompt", config: { prompt: "lint it" } },
|
|
{ id: "execute", kind: "prompt", config: { seam: "execute" } },
|
|
{ id: "review", kind: "prompt", config: { seam: "review" } },
|
|
{ id: "merge", kind: "prompt", config: { seam: "merge" } },
|
|
{ id: "notify", kind: "script", config: { scriptName: "notify" } },
|
|
{ id: "zend", kind: "end" },
|
|
],
|
|
edges: [
|
|
{ from: "start", to: "lint" },
|
|
{ from: "lint", to: "execute", condition: "success" },
|
|
{ from: "execute", to: "review", condition: "success" },
|
|
{ from: "review", to: "merge", condition: "success" },
|
|
{ from: "merge", to: "notify", condition: "success" },
|
|
{ from: "notify", to: "zend", condition: "success" },
|
|
{ from: "execute", to: "zend", condition: "failure" },
|
|
{ from: "review", to: "zend", condition: "failure" },
|
|
{ from: "merge", to: "zend", condition: "failure" },
|
|
],
|
|
};
|
|
}
|
|
|
|
function definition(ir: WorkflowIr): WorkflowDefinition {
|
|
return {
|
|
id: "WF-001",
|
|
name: "Full lifecycle",
|
|
description: "",
|
|
kind: "workflow",
|
|
ir,
|
|
layout: {},
|
|
createdAt: "2026-06-03T00:00:00.000Z",
|
|
updatedAt: "2026-06-03T00:00:00.000Z",
|
|
};
|
|
}
|
|
|
|
function storeWith(def: WorkflowDefinition | undefined, workflowId = "WF-001"): WorkflowGraphRunnerStore {
|
|
return {
|
|
getTaskWorkflowSelection: () => (def ? { workflowId, stepIds: [] } : undefined),
|
|
getWorkflowDefinition: async () => def,
|
|
};
|
|
}
|
|
|
|
function recordingSeams(calls: string[], overrides: Partial<Record<string, WorkflowNodeResult>> = {}) {
|
|
const seam = (name: string) => async (): Promise<WorkflowNodeResult> => {
|
|
calls.push(name);
|
|
return overrides[name] ?? { outcome: "success" };
|
|
};
|
|
return {
|
|
planning: seam("planning"),
|
|
execute: seam("execute"),
|
|
workflowStep: seam("workflow-step"),
|
|
review: seam("review"),
|
|
merge: seam("merge"),
|
|
schedule: seam("schedule"),
|
|
};
|
|
}
|
|
|
|
describe("WorkflowGraphTaskRunner (CU-U2)", () => {
|
|
it("runs the full lifecycle in graph order: custom → execute → review → merge → custom", async () => {
|
|
const calls: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return { outcome: "success" };
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
|
|
expect(result.disposition).toBe("completed");
|
|
expect(calls).toEqual(["custom:lint", "execute", "review", "merge", "custom:notify"]);
|
|
expect(result.visitedNodeIds).toEqual(["start", "lint", "execute", "review", "merge", "notify"]);
|
|
});
|
|
|
|
it("projects agent-generated completion summaries from workflow nodes onto the task", async () => {
|
|
const ir: WorkflowIr = {
|
|
version: "v1",
|
|
name: "summary",
|
|
nodes: [
|
|
{ id: "start", kind: "start" },
|
|
{ id: "summary", kind: "prompt", config: { prompt: "summarize", summaryTarget: "task" } },
|
|
{ id: "zend", kind: "end" },
|
|
],
|
|
edges: [
|
|
{ from: "start", to: "summary" },
|
|
{ from: "summary", to: "zend", condition: "success" },
|
|
],
|
|
};
|
|
const projections: Array<{ taskId: string; summary?: string }> = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(ir)),
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({
|
|
outcome: "success",
|
|
contextPatch: { summary: "Implemented the workflow and verified the result." },
|
|
}),
|
|
publishTaskProjection: async (taskId, patch) => {
|
|
projections.push({ taskId, summary: patch.summary });
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
|
|
expect(result.disposition).toBe("completed");
|
|
expect(projections).toEqual([
|
|
{ taskId: "FN-9001", summary: "Implemented the workflow and verified the result." },
|
|
]);
|
|
});
|
|
|
|
it("selected workflows reaching the merge seam produce the canonical merged notification once", async () => {
|
|
const emitter = new EventEmitter();
|
|
const graphTask = {
|
|
...task,
|
|
title: "Workflow merge",
|
|
description: "Graph path",
|
|
column: "done",
|
|
mergeDetails: { mergeConfirmed: true },
|
|
} as TaskDetail;
|
|
const store = Object.assign(emitter, {
|
|
getSettings: vi.fn(async () => ({ ntfyEnabled: true, ntfyTopic: "topic" }) as Settings),
|
|
getTask: vi.fn(async (_id: string) => graphTask),
|
|
getTaskWorkflowSelection: () => ({ workflowId: "WF-001", stepIds: [] }),
|
|
getWorkflowDefinition: async () => definition(fullLifecycleIr()),
|
|
}) as unknown as EventEmitter & WorkflowGraphRunnerStore & {
|
|
getSettings: () => Promise<Settings>;
|
|
getTask: (id: string) => Promise<TaskDetail>;
|
|
};
|
|
const sendNotification = vi.fn(async () => ({ success: true, providerId: "mock" }));
|
|
const service = new NotificationService(store as any);
|
|
service.registerProvider({ getProviderId: () => "mock", isEventSupported: () => true, sendNotification });
|
|
await service.start();
|
|
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store,
|
|
seams: {
|
|
...recordingSeams([]),
|
|
merge: async () => {
|
|
store.emit("task:moved", { task: graphTask, from: "in-review", to: "done" });
|
|
store.emit("task:merged", {
|
|
task: graphTask,
|
|
branch: "fusion/fn-9001",
|
|
merged: true,
|
|
worktreeRemoved: false,
|
|
branchDeleted: false,
|
|
});
|
|
return { outcome: "success", value: "merged" };
|
|
},
|
|
},
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
});
|
|
|
|
const result = await runner.run(graphTask, flagOn);
|
|
|
|
expect(result.disposition).toBe("completed");
|
|
await vi.waitFor(() => {
|
|
expect(sendNotification).toHaveBeenCalledTimes(1);
|
|
});
|
|
expect(sendNotification).toHaveBeenCalledWith(
|
|
"merged",
|
|
expect.objectContaining({ taskId: "FN-9001", taskTitle: "Workflow merge", event: "merged" }),
|
|
);
|
|
await service.stop();
|
|
});
|
|
|
|
it("a failing seam terminates the run as failed without running later nodes", async () => {
|
|
const calls: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams(calls, { review: { outcome: "failure", value: "REVISE" } }),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return { outcome: "success" };
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
|
|
expect(result.disposition).toBe("failed");
|
|
expect(calls).toEqual(["custom:lint", "execute", "review"]);
|
|
expect(calls).not.toContain("merge");
|
|
});
|
|
|
|
it("a failing custom gate before execute blocks the whole pipeline", async () => {
|
|
const calls: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return node.id === "lint" ? { outcome: "failure", value: "lint-failed" } : { outcome: "success" };
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
|
|
expect(result.disposition).toBe("failed");
|
|
expect(calls).toEqual(["custom:lint"]);
|
|
});
|
|
|
|
it("ignores stale workflowGraphExecutor=false and still runs the graph", async () => {
|
|
const calls: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return { outcome: "success" };
|
|
},
|
|
});
|
|
const result = await runner.run(task, flagOff);
|
|
expect(result.disposition).toBe("completed");
|
|
expect(calls).toEqual(["custom:lint", "execute", "review", "merge", "custom:notify"]);
|
|
expect(result.visitedNodeIds).toEqual(["start", "lint", "execute", "review", "merge", "notify"]);
|
|
});
|
|
|
|
it("falls back when the task has no workflow selection", async () => {
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(undefined),
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
});
|
|
const result = await runner.run(task, flagOn);
|
|
expect(result).toMatchObject({ disposition: "fell-back", reason: "no-selection" });
|
|
});
|
|
|
|
it("falls back when the selected workflow no longer exists", async () => {
|
|
const store: WorkflowGraphRunnerStore = {
|
|
getTaskWorkflowSelection: () => ({ workflowId: "WF-404", stepIds: [] }),
|
|
getWorkflowDefinition: async () => undefined,
|
|
};
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store,
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
});
|
|
const result = await runner.run(task, flagOn);
|
|
expect(result.disposition).toBe("fell-back");
|
|
expect(result.reason).toMatch(/workflow-missing/);
|
|
});
|
|
|
|
// FNXC:PgMigrationQuarantine 2026-07-17-18:15: FN-8258 exercises persisted workflow selection against the PostgreSQL AsyncDataLayer rather than removed SQLite inMemoryDb construction.
|
|
(PG_AVAILABLE ? it : it.skip)("persists a valid workflow through the store and launches it through the graph runner", async () => {
|
|
const harness = await createTaskStoreForTest({ prefix: "fusion_workflow_graph_runner" });
|
|
const store = harness.store;
|
|
try {
|
|
const invalidIr: WorkflowIr = {
|
|
version: "v2",
|
|
name: "invalid-save",
|
|
columns: [{ id: "todo", name: "Todo", traits: [] }],
|
|
nodes: [
|
|
{ id: "start", kind: "start", column: "todo" },
|
|
{ id: "dup", kind: "prompt", column: "todo" },
|
|
{ id: "dup", kind: "script", column: "todo" },
|
|
{ id: "end", kind: "end", column: "todo" },
|
|
],
|
|
edges: [
|
|
{ from: "start", to: "dup" },
|
|
{ from: "dup", to: "end" },
|
|
],
|
|
};
|
|
await expect(store.createWorkflowDefinition({ name: "Invalid", ir: invalidIr })).rejects.toThrow(
|
|
/Workflow IR has duplicate node id 'dup'/,
|
|
);
|
|
|
|
const workflow = await store.createWorkflowDefinition({ name: "Valid", ir: fullLifecycleIr() });
|
|
const persisted = await store.getWorkflowDefinition(workflow.id);
|
|
expect(persisted?.id).toBe(workflow.id);
|
|
const savedTask = await store.createTask({ description: "save run", enabledWorkflowSteps: [] });
|
|
await store.selectTaskWorkflow(savedTask.id, workflow.id);
|
|
|
|
const calls: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store,
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return { outcome: "success" };
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(savedTask, flagOn);
|
|
|
|
expect(result.disposition).toBe("completed");
|
|
expect(calls).toEqual(["custom:lint", "execute", "review", "merge", "custom:notify"]);
|
|
} finally {
|
|
await harness.teardown();
|
|
}
|
|
});
|
|
|
|
it("resolves built-in workflow selections without requiring the store to return a definition", async () => {
|
|
const calls: string[] = [];
|
|
const getWorkflowDefinition = vi.fn(async () => undefined);
|
|
const store: WorkflowGraphRunnerStore = {
|
|
getTaskWorkflowSelection: () => ({ workflowId: "builtin:legacy-coding", stepIds: [] }),
|
|
getWorkflowDefinition,
|
|
};
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store,
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
|
|
expect(result.disposition).toBe("completed");
|
|
// FNXC:WorkflowBuiltins 2026-06-28-23:29:
|
|
// Use Legacy coding here because default Coding is now stepwise and requires
|
|
// parse/foreach task-step context. This still proves built-in registry fallback
|
|
// works when the store intentionally does not return a persisted definition.
|
|
expect(calls).toEqual(["planning", "execute", "review", "merge"]);
|
|
expect(result.reason).toBeUndefined();
|
|
expect(getWorkflowDefinition).not.toHaveBeenCalled();
|
|
});
|
|
|
|
it("fails closed with invalid-ir before any side-effect seam when resolved IR is malformed", async () => {
|
|
const badIr: WorkflowIr = {
|
|
version: "v2",
|
|
name: "bad",
|
|
columns: [{ id: "c", name: "C", traits: [] }],
|
|
nodes: [
|
|
{ id: "start", kind: "start", column: "c" },
|
|
{ id: "dup", kind: "prompt", column: "c" },
|
|
{ id: "dup", kind: "script", column: "c" },
|
|
{ id: "end", kind: "end", column: "c" },
|
|
],
|
|
edges: [
|
|
{ from: "start", to: "dup" },
|
|
{ from: "dup", to: "end" },
|
|
],
|
|
};
|
|
const calls: string[] = [];
|
|
const events: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(badIr)),
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return { outcome: "success" };
|
|
},
|
|
onEvent: (e) => events.push(`${e.type}:${e.detail}`),
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
|
|
expect(result.disposition).toBe("failed");
|
|
expect(result.outcome).toBe("failure");
|
|
expect(result.reason).toMatch(/invalid-ir: Workflow IR has duplicate node id 'dup'/);
|
|
expect(result.visitedNodeIds).toEqual([]);
|
|
expect(calls).toEqual([]);
|
|
expect(events.some((event) => event.includes("terminal:invalid-ir"))).toBe(true);
|
|
});
|
|
|
|
it("fails closed with invalid-ir before any custom node when resolved IR has a dangling edge", async () => {
|
|
const badIr: WorkflowIr = {
|
|
version: "v2",
|
|
name: "bad-edge",
|
|
columns: [{ id: "c", name: "C", traits: [] }],
|
|
nodes: [
|
|
{ id: "start", kind: "start", column: "c" },
|
|
{ id: "a", kind: "prompt", column: "c" },
|
|
{ id: "end", kind: "end", column: "c" },
|
|
],
|
|
edges: [
|
|
{ from: "start", to: "a" },
|
|
{ from: "a", to: "ghost" },
|
|
],
|
|
};
|
|
const calls: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(badIr)),
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return { outcome: "success" };
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
|
|
expect(result.disposition).toBe("failed");
|
|
expect(result.outcome).toBe("failure");
|
|
expect(result.reason).toMatch(/invalid-ir: Workflow edge 'a' -> 'ghost' references undefined node 'ghost'/);
|
|
expect(result.visitedNodeIds).toEqual([]);
|
|
expect(calls).toEqual([]);
|
|
});
|
|
|
|
it("a custom-node failure AFTER side effects terminates as failed, not fell-back", async () => {
|
|
const calls: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
throw new Error(`custom boom: ${node.id}`);
|
|
},
|
|
});
|
|
const result = await runner.run(task, flagOn);
|
|
expect(calls).toEqual([]);
|
|
expect(result.visitedNodeIds).toEqual(["start", "lint"]);
|
|
expect(result.disposition).toBe("failed");
|
|
expect(result.reason).toBeUndefined();
|
|
});
|
|
|
|
it("exposes node outcomes in the shared context for downstream consumers", async () => {
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({ outcome: "success", value: "APPROVE" }),
|
|
});
|
|
const result = await runner.run(task, flagOn);
|
|
expect(result.context?.["node:lint:outcome"]).toBe("success");
|
|
expect(result.context?.["node:lint:value"]).toBe("APPROVE");
|
|
});
|
|
|
|
it("runs principal admission before an executable handler and honors its fail-closed result", async () => {
|
|
const calls: string[] = [];
|
|
const admitted: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams(calls),
|
|
runCustomNode: async (node) => {
|
|
calls.push(`custom:${node.id}`);
|
|
return { outcome: "success" };
|
|
},
|
|
beforeNodeExecution: (node, _task, context) => {
|
|
if (node.id === "lint") {
|
|
admitted.push(node.id);
|
|
context["workflow:principal-agent-id"] = "planner";
|
|
return { outcome: "failure", value: "workflow-principal-named-principal-unavailable:triage" };
|
|
}
|
|
return undefined;
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
expect(admitted).toEqual(["lint"]);
|
|
expect(calls).toEqual([]);
|
|
expect(result.disposition).toBe("suspended");
|
|
expect(result.suspension).toMatchObject({ reason: "capacity", nodeId: "lint" });
|
|
expect(result.context?.["workflow:principal-agent-id"]).toBe("planner");
|
|
});
|
|
|
|
it("restores a direct-resume principal fence before node admission", async () => {
|
|
const observed: Record<string, unknown>[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
beforeNodeExecution: (node, _task, context) => {
|
|
if (node.id === "lint") observed.push({ ...context });
|
|
},
|
|
});
|
|
|
|
await runner.run(task, flagOn, "lint", {
|
|
"workflow:work-item-id": "work-item-1",
|
|
"workflow:principal-agent-id": "reviewer-1",
|
|
"workflow:principal-role": "reviewer",
|
|
"workflow:principal-authority": "review-node-override",
|
|
"workflow:node-instance-id": "lint",
|
|
});
|
|
|
|
expect(observed).toEqual([expect.objectContaining({
|
|
"workflow:work-item-id": "work-item-1",
|
|
"workflow:principal-agent-id": "reviewer-1",
|
|
"workflow:principal-role": "reviewer",
|
|
"workflow:principal-authority": "review-node-override",
|
|
"workflow:node-instance-id": "lint",
|
|
})]);
|
|
});
|
|
|
|
it("releases a principal reservation after its node handler settles", async () => { const released: string[] = [];
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
beforeNodeExecution: (node, _task, context) => {
|
|
if (node.id === "execute") {
|
|
context["workflow:release-principal"] = () => released.push(node.id);
|
|
}
|
|
},
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
expect(result.disposition).toBe("completed");
|
|
expect(released).toEqual(["execute"]);
|
|
});
|
|
|
|
it("onEvent diagnostics failures never affect the run", async () => {
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fullLifecycleIr())),
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
onEvent: () => {
|
|
throw new Error("diagnostics boom");
|
|
},
|
|
});
|
|
const result = await runner.run(task, flagOn);
|
|
expect(result.disposition).toBe("completed");
|
|
});
|
|
|
|
// #1407/#1412: the runner forwards its injected branchPersistence into the
|
|
// WorkflowGraphExecutor, which writes per-branch state and prunes stale runs.
|
|
// Uses a real in-memory persistence whose method shape matches the store-
|
|
// backed adapter the production executor builds (saveBranchState /
|
|
// loadBranchStates / clearStaleBranchStates) — no mock of a nonexistent API.
|
|
function fanoutIr(): WorkflowIr {
|
|
return {
|
|
version: "v1",
|
|
name: "fanout",
|
|
nodes: [
|
|
{ id: "start", kind: "start" },
|
|
{ id: "split", kind: "split" },
|
|
{ id: "a", kind: "prompt", config: { prompt: "a" } },
|
|
{ id: "b", kind: "prompt", config: { prompt: "b" } },
|
|
{ id: "join", kind: "join", config: { mode: "all" } },
|
|
{ id: "zend", kind: "end" },
|
|
],
|
|
edges: [
|
|
{ from: "start", to: "split" },
|
|
{ from: "split", to: "a" },
|
|
{ from: "split", to: "b" },
|
|
{ from: "a", to: "join" },
|
|
{ from: "b", to: "join" },
|
|
{ from: "join", to: "zend", condition: "success" },
|
|
],
|
|
};
|
|
}
|
|
|
|
it("forwards branchPersistence to the executor: writes branch state and prunes stale runs", async () => {
|
|
const saved: Array<{ branchId: string; currentNodeId: string; status: string }> = [];
|
|
const pruneCalls: Array<{ taskId: string; keepRunId: string }> = [];
|
|
const persistence = {
|
|
saveBranchState: (s: { branchId: string; currentNodeId: string; status: string }) => {
|
|
saved.push({ branchId: s.branchId, currentNodeId: s.currentNodeId, status: s.status });
|
|
},
|
|
loadBranchStates: () => [],
|
|
clearStaleBranchStates: (taskId: string, keepRunId: string) => {
|
|
pruneCalls.push({ taskId, keepRunId });
|
|
},
|
|
};
|
|
|
|
const runner = new WorkflowGraphTaskRunner({
|
|
store: storeWith(definition(fanoutIr())),
|
|
seams: recordingSeams([]),
|
|
runCustomNode: async () => ({ outcome: "success" }),
|
|
branchPersistence: persistence,
|
|
});
|
|
|
|
const result = await runner.run(task, flagOn);
|
|
expect(result.disposition).toBe("completed");
|
|
|
|
// Both branches persisted, and each reached "completed" at the join.
|
|
expect(saved.some((s) => s.branchId === "a")).toBe(true);
|
|
expect(saved.some((s) => s.branchId === "b")).toBe(true);
|
|
expect(saved.some((s) => s.status === "completed")).toBe(true);
|
|
|
|
// Prune ran (on start AND completion) keyed by the runner's runId.
|
|
expect(pruneCalls.length).toBeGreaterThanOrEqual(2);
|
|
expect(pruneCalls.every((c) => c.taskId === task.id)).toBe(true);
|
|
expect(pruneCalls.every((c) => c.keepRunId === `${task.id}:WF-001`)).toBe(true);
|
|
});
|
|
});
|