using System.Text.Json.Nodes; using w4c_workflows.Models.Nodes; using w4c_workflows.Services.Nodes; using w4c_workflows.Services.Nodes.Executors; using w4c_workflows.Services.Nodes.Interpolation; using Xunit; namespace w4c_workflows.Tests; public class NodeGraphRunnerTests { // ------------------------------------------------------------------ fixtures private sealed class FailingNodeExecutor : INodeExecutor { public string Type => "test.fail"; public Task RunAsync(NodeExecutionContext context, CancellationToken ct) => Task.FromResult(NodeExecutionOutcome.Failed("boom", "boom_code", "exploded")); } private static NodeGraphRunner Runner(params INodeExecutor[] extra) { var executors = new List { new NoOpNodeExecutor(), new SetNodeExecutor(), new IfNodeExecutor(), }; executors.AddRange(extra); return new NodeGraphRunner(new NodeExecutorRegistry(executors), new NodeParameterInterpolator()); } private static NodeGraph Compile(string yaml, NodeBlueprintCatalog? catalog = null) { var result = new NodeGraphCompiler(catalog ?? NodeTestData.CoreCatalog()).Compile(NodeTestData.Parse(yaml)); Assert.True(result.Success, string.Join("\n", result.Errors)); return result.Graph!; } private static FlowItem Item(string json) => FlowItem.FromJson(JsonNode.Parse(json)!.AsObject()); private static NodeBlueprintCatalog FailCatalog(params NodePort[] outputs) => NodeTestData.Catalog(NodeTestData.NodeBlueprint("test.fail", outputs: outputs.ToList())); // ------------------------------------------------------------------ tests [Fact] public async Task Linear_graph_shapes_items_with_per_item_interpolation() { var yaml = """ name: set-demo tasks: - id: pass node: { type: core.noop } next: shape - id: shape node: { type: core.set } parameters: mode: manual keepOnlySet: true fields: - name: greeting value: "=Hello {{ $json.name }}" """; var result = await Runner().RunAsync(Compile(yaml), new[] { Item("""{"name":"Ada"}""") }); Assert.True(result.Succeeded, result.Failure?.Message); Assert.Equal(new[] { "pass", "shape" }, result.ExecutionOrder); Assert.Single(result.OutputOf("pass")); var shaped = result.OutputOf("shape"); Assert.Single(shaped); Assert.Equal("Hello Ada", shaped[0].Json["greeting"]!.GetValue()); Assert.Single(shaped[0].Json); // keepOnlySet dropped the input field } [Fact] public async Task If_routes_each_item_to_the_matching_branch() { var yaml = """ name: branch tasks: - id: decide node: { type: core.if } parameters: left: "={{ $json.status }}" operator: equals right: ok edges: - { output: 0, to: yes } - { output: 1, to: no } - id: yes node: { type: core.noop } - id: no node: { type: core.noop } """; var result = await Runner().RunAsync( Compile(yaml), new[] { Item("""{"status":"ok"}"""), Item("""{"status":"bad"}""") }); Assert.True(result.Succeeded, result.Failure?.Message); Assert.Single(result.OutputOf("yes")); Assert.Single(result.OutputOf("no")); Assert.Equal("ok", result.OutputOf("yes")[0].Json["status"]!.GetValue()); Assert.Equal("bad", result.OutputOf("no")[0].Json["status"]!.GetValue()); } [Fact] public async Task Node_reference_reads_an_upstream_nodes_output() { var yaml = """ name: ref tasks: - id: source node: { type: core.set } parameters: mode: json jsonOutput: { id: 7 } next: echo - id: echo node: { type: core.set } parameters: mode: manual keepOnlySet: true fields: - name: copied value: '={{ $("source").id }}' """; var result = await Runner().RunAsync(Compile(yaml), Array.Empty()); Assert.True(result.Succeeded, result.Failure?.Message); var echoed = result.OutputOf("echo"); Assert.Single(echoed); Assert.Equal(7, echoed[0].Json["copied"]!.GetValue()); } [Fact] public async Task Produced_items_carry_provenance() { var yaml = """ name: prov tasks: - id: pass node: { type: core.noop } next: shape - id: shape node: { type: core.set } parameters: mode: manual fields: - name: seen value: "true" """; var result = await Runner().RunAsync(Compile(yaml), new[] { Item("""{"id":1}""") }); var passOrigin = result.OutputOf("pass")[0].Origin; Assert.NotNull(passOrigin); Assert.Equal("pass", passOrigin!.NodeId); Assert.Equal(0, passOrigin.OutputIndex); // The set node passes the input item through, so it keeps the upstream origin. Assert.Equal("pass", result.OutputOf("shape")[0].Origin!.NodeId); } [Fact] public async Task A_failure_is_routed_to_the_error_output() { var yaml = """ name: error-branch tasks: - id: work node: { type: test.fail } onError: recover - id: recover node: { type: core.noop } """; var catalog = FailCatalog(NodePort.Main, NodePort.Error); var result = await Runner(new FailingNodeExecutor()) .RunAsync(Compile(yaml, catalog), new[] { Item("""{"x":1}""") }); Assert.True(result.Succeeded, result.Failure?.Message); var recovered = result.OutputOf("recover"); Assert.Single(recovered); Assert.Equal(1, recovered[0].Json["x"]!.GetValue()); Assert.Equal("boom", recovered[0].Json["error"]!["message"]!.GetValue()); Assert.Equal("boom_code", recovered[0].Json["error"]!["code"]!.GetValue()); } [Fact] public async Task A_failure_without_an_error_output_fails_the_run() { var yaml = """ name: hard-fail tasks: - id: work node: { type: test.fail } """; var result = await Runner(new FailingNodeExecutor()) .RunAsync(Compile(yaml, FailCatalog(NodePort.Main)), new[] { Item("{}") }); Assert.False(result.Succeeded); Assert.Equal("boom", result.Failure!.Message); } [Fact] public async Task A_failure_with_an_unconnected_error_output_fails_the_run() { // The blueprint declares an error port, but no branch is wired to it. // Merely having the port must not swallow the failure (P1-10 depends on // the failure reaching WorkflowRun.Error); only a real error branch does. var yaml = """ name: dangling-error-port tasks: - id: work node: { type: test.fail } """; var result = await Runner(new FailingNodeExecutor()) .RunAsync(Compile(yaml, FailCatalog(NodePort.Main, NodePort.Error)), new[] { Item("{}") }); Assert.False(result.Succeeded); Assert.Equal("boom", result.Failure!.Message); } [Fact] public async Task ContinueOnFail_succeeds_without_an_error_output() { var yaml = """ name: soft-fail tasks: - id: work node: { type: test.fail } continueOnFail: true """; var result = await Runner(new FailingNodeExecutor()) .RunAsync(Compile(yaml, FailCatalog(NodePort.Main)), new[] { Item("{}") }); Assert.True(result.Succeeded, result.Failure?.Message); Assert.Empty(result.OutputOf("work")); } [Fact] public async Task A_missing_executor_fails_the_run() { var catalog = NodeTestData.Catalog(NodeTestData.NodeBlueprint("test.noexec")); var yaml = """ name: no-executor tasks: - id: work node: { type: test.noexec } """; var result = await Runner().RunAsync(Compile(yaml, catalog), new[] { Item("{}") }); Assert.False(result.Succeeded); Assert.Equal("missing_executor", result.Failure!.Code); } [Fact] public async Task Fan_out_sends_the_same_items_to_every_target() { var yaml = """ name: fan-out tasks: - id: source node: { type: core.noop } edges: - { output: 0, to: left } - { output: 0, to: right } - id: left node: { type: core.noop } - id: right node: { type: core.noop } """; var result = await Runner().RunAsync(Compile(yaml), new[] { Item("""{"n":1}""") }); Assert.True(result.Succeeded, result.Failure?.Message); Assert.Single(result.OutputOf("left")); Assert.Single(result.OutputOf("right")); } [Fact] public async Task Filter_routes_matching_and_non_matching_items() { var yaml = """ name: filter-demo tasks: - id: f node: { type: core.filter } parameters: left: "={{ $json.n }}" operator: greater right: 1 edges: - { output: 0, to: kept } - { output: 1, to: dropped } - id: kept node: { type: core.noop } - id: dropped node: { type: core.noop } """; var result = await Runner(new FilterNodeExecutor()) .RunAsync(Compile(yaml), new[] { Item("""{"n":0}"""), Item("""{"n":2}""") }); Assert.True(result.Succeeded, result.Failure?.Message); Assert.Equal(2, Assert.Single(result.OutputOf("kept")).Json["n"]!.GetValue()); Assert.Equal(0, Assert.Single(result.OutputOf("dropped")).Json["n"]!.GetValue()); } [Fact] public async Task Switch_routes_each_item_to_its_rule_output() { var yaml = """ name: switch-demo tasks: - id: route node: { type: core.switch } parameters: rules: - { left: "={{ $json.kind }}", operator: equals, right: a } - { left: "={{ $json.kind }}", operator: equals, right: b } fallbackEnabled: true edges: - { output: 0, to: a } - { output: 1, to: b } - { output: 4, to: other } - id: a node: { type: core.noop } - id: b node: { type: core.noop } - id: other node: { type: core.noop } """; var result = await Runner(new SwitchNodeExecutor()).RunAsync( Compile(yaml), new[] { Item("""{"kind":"a"}"""), Item("""{"kind":"b"}"""), Item("""{"kind":"c"}""") }); Assert.True(result.Succeeded, result.Failure?.Message); Assert.Equal("a", Assert.Single(result.OutputOf("a")).Json["kind"]!.GetValue()); Assert.Equal("b", Assert.Single(result.OutputOf("b")).Json["kind"]!.GetValue()); Assert.Equal("c", Assert.Single(result.OutputOf("other")).Json["kind"]!.GetValue()); } [Fact] public async Task A_loop_node_iterates_its_body_then_emits_done() { var yaml = """ name: loop-demo tasks: - id: source node: { type: core.noop } edges: - { output: 0, to: loop } - id: loop node: { type: core.splitInBatches } parameters: { batchSize: 2 } edges: - { output: 0, to: body } - { output: 1, to: done } - id: body node: { type: core.noop } edges: - { output: 0, to: loop } - id: done node: { type: core.noop } """; var seed = Enumerable.Range(1, 5).Select(i => Item($$"""{"n":{{i}}}""")).ToArray(); var result = await Runner(new SplitInBatchesNodeExecutor()).RunAsync(Compile(yaml), seed); Assert.True(result.Succeeded, result.Failure?.Message); // 5 items in batches of 2 → 3 batch iterations + 1 final (done) invocation. Assert.Equal(4, result.ExecutionOrder.Count(id => id == "loop")); Assert.Equal(3, result.ExecutionOrder.Count(id => id == "body")); Assert.Equal(5, result.OutputOf("done").Count); // original items reach the exit } /// /// P1-3: an entry loop node (zero non-loop-back incoming) must be re-triggered /// by its own loop-back edge. Previously readiness compared `arrived == /// incoming` and an entry loop node (incoming 0) silently ran a single batch. /// [Fact] public async Task An_entry_loop_node_iterates_each_batch() { var yaml = """ name: entry-loop tasks: - id: loop node: { type: core.splitInBatches } parameters: { batchSize: 2 } edges: - { output: 0, to: body } - { output: 1, to: done } - id: body node: { type: core.noop } edges: - { output: 0, to: loop } - id: done node: { type: core.noop } """; var seed = Enumerable.Range(1, 5).Select(i => Item($$"""{"n":{{i}}}""")).ToArray(); var result = await Runner(new SplitInBatchesNodeExecutor()).RunAsync(Compile(yaml), seed); Assert.True(result.Succeeded, result.Failure?.Message); Assert.Equal(4, result.ExecutionOrder.Count(id => id == "loop")); Assert.Equal(3, result.ExecutionOrder.Count(id => id == "body")); Assert.Equal(5, result.OutputOf("done").Count); } /// Records the run iteration each invocation observed. private sealed class RecordingRunIndexExecutor : INodeExecutor { public List RunIndices { get; } = new(); public List InterpolatedRunIndices { get; } = new(); public string Type => "test.record"; public Task RunAsync(NodeExecutionContext context, CancellationToken ct) { RunIndices.Add(context.RunIndex); InterpolatedRunIndices.Add( context.Parameters.TryGetPropertyValue("seen", out var value) && value != null ? value.GetValue() : -1); return Task.FromResult(NodeExecutionOutcome.Single(context.Input(0).ToList())); } } /// /// The run iteration index is real, not the hard-coded 0 it used to be: a node /// re-triggered by its loop-back edge sees an incrementing /// and $runIndex resolves to it. /// [Fact] public async Task Loop_body_sees_an_incrementing_run_index() { var recording = new RecordingRunIndexExecutor(); var loop = new NodeGraphNode { Id = "loop", Blueprint = NodeTestData.CoreCatalog().Get("core.splitInBatches")!, Parameters = new JsonObject { ["batchSize"] = 2 }, }; var body = new NodeGraphNode { Id = "body", Blueprint = NodeTestData.NodeBlueprint("test.record"), Parameters = new JsonObject { ["seen"] = "={{ $runIndex }}" }, }; var graph = new NodeGraph { Nodes = new[] { loop, body }, Edges = new[] { new NodeGraphEdge { FromNodeId = "loop", FromOutput = 0, ToNodeId = "body", ToInput = 0 }, new NodeGraphEdge { FromNodeId = "body", FromOutput = 0, ToNodeId = "loop", ToInput = 0, IsLoopBack = true }, }, EntryNodeId = "loop", }; var seed = Enumerable.Range(1, 5).Select(i => Item($$"""{"n":{{i}}}""")).ToArray(); var result = await Runner(new SplitInBatchesNodeExecutor(), recording).RunAsync(graph, seed); Assert.True(result.Succeeded, result.Failure?.Message); // 5 items in batches of 2 → the body runs on iterations 0, 1 and 2 (the // fourth loop invocation is the empty "done" pass that emits no batch). Assert.Equal(new[] { 0, 1, 2 }, recording.RunIndices); Assert.Equal(new[] { 0, 1, 2 }, recording.InterpolatedRunIndices); } }