using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Logging.Abstractions; using w4c_workflows.Data; using w4c_workflows.Models; using w4c_workflows.Models.Credentials; using w4c_workflows.Models.Nodes; using w4c_workflows.Services; using w4c_workflows.Services.Credentials; using w4c_workflows.Services.Execution; using w4c_workflows.Services.Messaging; using w4c_workflows.Services.Nodes; using w4c_workflows.Services.Nodes.Executors; using w4c_workflows.Services.Nodes.Interpolation; using Xunit; namespace w4c_workflows.Tests; /// /// S2 worker half: a graph.run job is claimed by the worker and drives the /// whole node graph through , recording one /// per node and owning the terminal run status. /// [Collection("WorkflowsPostgres")] public class GraphJobExecutorTests { private readonly WorkflowsPostgresFixture _fixture; public GraphJobExecutorTests(WorkflowsPostgresFixture fixture) { _fixture = fixture; } private static NodeWorkflowRunner NewRunner( WorkflowsDbContext db, IEnumerable? extraExecutors = null, CredentialVault? vault = null) { var catalog = new NodeBlueprintCatalog(NodeBlueprintCatalog.LoadEmbedded()); var executors = new List { new NoOpNodeExecutor(), new SetNodeExecutor(), new IfNodeExecutor(), new SplitInBatchesNodeExecutor(), new ExecuteWorkflowNodeExecutor(), new CodeNodeExecutor(new RuntimeRegistry(new LanguageRegistry(), Array.Empty())), }; if (extraExecutors != null) executors.AddRange(extraExecutors); var graphRunner = new NodeGraphRunner(new NodeExecutorRegistry(executors), new NodeParameterInterpolator()); return new NodeWorkflowRunner( catalog, graphRunner, vault ?? new CredentialVault(new ReversibleTestCipher(), new CredentialTypeCatalog())); } private static (GraphJobExecutor Executor, FakeJobQueue Jobs) NewExecutor( WorkflowsDbContext db, IEnumerable? extraExecutors = null, CredentialVault? vault = null) { var nodeRunner = NewRunner(db, extraExecutors, vault); var jobs = new FakeJobQueue(); var services = new ServiceCollection(); services.AddSingleton(db); var provider = services.BuildServiceProvider(); var executor = new GraphJobExecutor( provider.GetRequiredService(), jobs, nodeRunner, NullLogger.Instance); return (executor, jobs); } private static StreamMessage GraphJob(WorkflowRun run, string? workingDir = null) => new("1-0", GraphRunMessage.ToFields(run, workingDir, 1)); private static string Tenant() => "t" + Guid.NewGuid().ToString("N")[..12]; [Fact] public async Task Executes_the_graph_and_records_a_task_run_per_node() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeYaml, "workflows/node-exec.yaml"); var run = WorkflowDataHelpers.AddPendingRun(db, workflow); var (executor, jobs) = NewExecutor(db); await executor.ExecuteAsync(tenantId, GraphJob(run), default); var reloaded = await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id); Assert.Equal(RunStatus.Succeeded, reloaded.Status); Assert.NotNull(reloaded.OutputJson); Assert.Contains("copied", reloaded.OutputJson); var taskRuns = await db.TaskRuns.Where(t => t.RunId == run.Id).ToListAsync(); Assert.Equal(2, taskRuns.Count); Assert.All(taskRuns, t => Assert.Equal(TaskRunStatus.Succeeded, t.Status)); Assert.Contains("1-0", jobs.Acked); } [Fact] public async Task Run_is_marked_running_before_the_graph_executes() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeYaml, "workflows/node-running.yaml"); var run = WorkflowDataHelpers.AddRun(db, workflow, RunStatus.Pending); var (executor, _) = NewExecutor(db); await executor.ExecuteAsync(tenantId, GraphJob(run), default); var reloaded = await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id); Assert.Equal(RunStatus.Succeeded, reloaded.Status); Assert.NotNull(reloaded.StartedAt); Assert.NotNull(reloaded.FinishedAt); } [Fact] public async Task Node_without_an_executor_fails_the_run() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeNoExecutorYaml, "workflows/node-fail.yaml"); var run = WorkflowDataHelpers.AddPendingRun(db, workflow); var (executor, jobs) = NewExecutor(db); await executor.ExecuteAsync(tenantId, GraphJob(run), default); var reloaded = await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id); Assert.Equal(RunStatus.Failed, reloaded.Status); Assert.Contains("no executor", reloaded.Error); Assert.Contains("1-0", jobs.Acked); } [Fact] public async Task Loop_runs_and_records_each_iteration() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); const string yaml = """ name: loop-run 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 workflow = WorkflowDataHelpers.CompileAndSave(db, tenantId, yaml, "workflows/loop.yaml"); Assert.Contains(workflow.TaskEdges, e => e.IsLoopBack); var run = WorkflowDataHelpers.AddPendingRun( db, workflow, """[{"n":1},{"n":2},{"n":3},{"n":4},{"n":5}]"""); var (executor, _) = NewExecutor(db); await executor.ExecuteAsync(tenantId, GraphJob(run), default); var reloaded = await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id); Assert.Equal(RunStatus.Succeeded, reloaded.Status); var loopTask = workflow.Tasks.Single(t => t.Key == "loop"); var iterations = await db.TaskRuns.CountAsync(t => t.RunId == run.Id && t.TaskId == loopTask.Id); Assert.Equal(4, iterations); // 3 batches + the final done invocation } [Fact] public async Task Resolves_declared_credentials_from_the_vault() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var vault = new CredentialVault(new ReversibleTestCipher(), new CredentialTypeCatalog()); await vault.CreateAsync( db, tenantId, "my-api", "httpHeaderAuth", new System.Text.Json.Nodes.JsonObject { ["name"] = "X-API-Key", ["value"] = "abc" }, default); var yaml = """ name: cred-run tasks: - id: fetch node: { type: core.httpRequest } parameters: { url: "https://api.example.com" } credentials: { httpAuth: my-api } """; var workflow = WorkflowDataHelpers.CompileAndSave(db, tenantId, yaml, "workflows/cred.yaml"); var run = WorkflowDataHelpers.AddPendingRun(db, workflow); var capturing = new CapturingHttpExecutor(); var (executor, _) = NewExecutor(db, new[] { capturing }, vault); await executor.ExecuteAsync(tenantId, GraphJob(run), default); Assert.NotNull(capturing.Seen); Assert.Equal("httpHeaderAuth", capturing.Seen!["httpAuth"].Type); Assert.Equal("abc", capturing.Seen!["httpAuth"].Data["value"]!.GetValue()); Assert.Equal(RunStatus.Succeeded, (await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id)).Status); } [Fact] public async Task A_terminal_run_is_skipped_and_acked_not_reexecuted() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeYaml, "workflows/node-terminal.yaml"); var run = WorkflowDataHelpers.AddRun(db, workflow, RunStatus.Succeeded); var (executor, jobs) = NewExecutor(db); await executor.ExecuteAsync(tenantId, GraphJob(run), default); Assert.Equal(RunStatus.Succeeded, (await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id)).Status); Assert.False(await db.TaskRuns.AnyAsync(t => t.RunId == run.Id)); Assert.Contains("1-0", jobs.Acked); } [Fact] public async Task Malformed_job_is_dead_lettered() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var (executor, jobs) = NewExecutor(db); var fields = new Dictionary { ["type"] = GraphRunInvocation.TypeValue, ["tenant_id"] = tenantId, // run_id deliberately missing }; await executor.ExecuteAsync(tenantId, new StreamMessage("9-0", fields), default); var dlq = Assert.Single(jobs.DeadLetteredFor(tenantId)); Assert.Equal("malformed", dlq.Reason); } [Fact] public async Task Tenant_mismatch_is_dead_lettered() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeYaml, "workflows/node-mismatch.yaml"); var run = WorkflowDataHelpers.AddPendingRun(db, workflow); var (executor, jobs) = NewExecutor(db); await executor.ExecuteAsync("other-tenant", GraphJob(run), default); var dlq = Assert.Single(jobs.DeadLetteredFor("other-tenant")); Assert.Equal("tenant_mismatch", dlq.Reason); Assert.Empty(jobs.Acked); } [Fact] public async Task Cancellation_leaves_the_job_unacked_for_redelivery() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeYaml, "workflows/node-cancel.yaml"); var run = WorkflowDataHelpers.AddPendingRun(db, workflow); var (executor, jobs) = NewExecutor(db); using var cts = new CancellationTokenSource(); cts.Cancel(); await Assert.ThrowsAnyAsync( () => executor.ExecuteAsync(tenantId, GraphJob(run), cts.Token)); Assert.Empty(jobs.Acked); } [Fact] public async Task Redelivery_retires_prior_in_flight_rows_and_numbers_the_new_attempt() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeYaml, "workflows/node-redelivery.yaml"); var run = WorkflowDataHelpers.AddRun(db, workflow, RunStatus.Running, startedAt: DateTime.UtcNow); // Simulate a crashed first delivery: the run is still running and the first // node's history row never finished. var firstTask = workflow.Tasks.Single(t => t.Key == "seed"); db.TaskRuns.Add(new TaskRun { Id = Guid.NewGuid(), RunId = run.Id, TaskId = firstTask.Id, Attempt = 1, Status = TaskRunStatus.Running, StartedAt = DateTime.UtcNow.AddSeconds(-5), }); await db.SaveChangesAsync(); var (executor, _) = NewExecutor(db); await executor.ExecuteAsync(tenantId, GraphJob(run), default); // Read from the database, not the change tracker: the retire uses ExecuteUpdate. db.ChangeTracker.Clear(); var rows = await db.TaskRuns.Where(t => t.RunId == run.Id).ToListAsync(); Assert.Single(rows, r => r.Attempt == 1 && r.Status == TaskRunStatus.Dead); Assert.Equal(2, rows.Count(r => r.Attempt == 2)); Assert.All(rows.Where(r => r.Attempt == 2), r => Assert.Equal(TaskRunStatus.Succeeded, r.Status)); Assert.DoesNotContain(rows, r => r.Status == TaskRunStatus.Running); var reloaded = await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id); Assert.Equal(RunStatus.Succeeded, reloaded.Status); } [Fact] public async Task A_run_already_terminalized_is_not_resurrected_by_a_late_worker() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave( db, tenantId, WorkflowDataHelpers.NodeYaml, "workflows/node-timeout-race.yaml"); // The control-plane timeout sweep already failed this run while a worker // was still executing it. var run = WorkflowDataHelpers.AddRun(db, workflow, RunStatus.Failed, startedAt: DateTime.UtcNow.AddMinutes(-10), finishedAt: DateTime.UtcNow, error: "run timed out after 300s"); var runner = NewRunner(db); var outcome = await runner.RunAsync(db, run, workflow, workingDirectory: null, default); Assert.False(outcome.Succeeded); Assert.Equal("the run was already terminal", outcome.Error); Assert.Equal(RunStatus.Failed, run.Status); Assert.Equal("run timed out after 300s", run.Error); } /// Stands in for the HTTP executor and records the credentials it was handed. private sealed class CapturingHttpExecutor : INodeExecutor { public string Type => "core.httpRequest"; public IReadOnlyDictionary? Seen { get; private set; } public Task RunAsync(NodeExecutionContext context, CancellationToken ct) { Seen = context.Credentials; return Task.FromResult(NodeExecutionOutcome.Single(context.Input(0).ToList())); } } /// Fails with the resolved secret embedded in the URL (P1-10 regression). private sealed class LeakyFailExecutor : INodeExecutor { public string Type => "core.httpRequest"; public Task RunAsync(NodeExecutionContext context, CancellationToken ct) { var secret = context.Credentials["httpAuth"].Data["value"]!.GetValue(); return Task.FromResult( NodeExecutionOutcome.Failed($"GET https://api.example.com/{secret}/items failed", "request_failed")); } } [Fact] public async Task Credential_literals_are_redacted_from_run_and_task_errors() { var tenantId = Tenant(); await using var db = _fixture.CreateContext(); var vault = new CredentialVault(new ReversibleTestCipher(), new CredentialTypeCatalog()); const string secret = "ghp_leakedTokenValue123456"; await vault.CreateAsync( db, tenantId, "my-api", "httpHeaderAuth", new System.Text.Json.Nodes.JsonObject { ["name"] = "X-API-Key", ["value"] = secret }, default); var yaml = """ name: cred-leak tasks: - id: fetch node: { type: core.httpRequest } parameters: { url: "https://api.example.com" } credentials: { httpAuth: my-api } """; var workflow = WorkflowDataHelpers.CompileAndSave(db, tenantId, yaml, "workflows/cred-leak.yaml"); var run = WorkflowDataHelpers.AddPendingRun(db, workflow); var (executor, _) = NewExecutor(db, new[] { new LeakyFailExecutor() }, vault); await executor.ExecuteAsync(tenantId, GraphJob(run), default); var reloaded = await db.WorkflowRuns.SingleAsync(r => r.Id == run.Id); Assert.Equal(RunStatus.Failed, reloaded.Status); Assert.DoesNotContain(secret, reloaded.Error); var taskRun = await db.TaskRuns.SingleAsync(t => t.RunId == run.Id); Assert.Equal(TaskRunStatus.Failed, taskRun.Status); Assert.DoesNotContain(secret, taskRun.Error); } }