using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Configuration; using Microsoft.Extensions.Logging.Abstractions; using w4c_workflows.Data; using w4c_workflows.Models; using w4c_workflows.Services.Execution; using w4c_workflows.Services.Runs; using Xunit; namespace w4c_workflows.Tests; [Collection("WorkflowsPostgres")] public class TaskDispatcherTests { private readonly WorkflowsPostgresFixture _fixture; public TaskDispatcherTests(WorkflowsPostgresFixture fixture) { _fixture = fixture; } private static TaskDispatcher NewDispatcher(WorkflowsDbContext db, FakeJobQueue jobs) => new(db, jobs, new ConfigurationBuilder().Build(), NullLogger.Instance); [Fact] public async Task Dispatch_creates_task_run_and_enqueues_job() { var tenantId = "t" + Guid.NewGuid().ToString("N")[..12]; await using var db = _fixture.CreateContext(); var workflow = WorkflowDataHelpers.CompileAndSave(db, tenantId, WorkflowDataHelpers.LinearYaml); var run = WorkflowDataHelpers.AddPendingRun(db, workflow, "{\"input\":1}"); var jobs = new FakeJobQueue(); var dispatcher = NewDispatcher(db, jobs); var taskRunId = await dispatcher.DispatchAsync( run, WorkflowDataHelpers.TaskByKey(workflow, "root"), workflow.Path, run.InputJson, 1, default); // The TaskRun is persisted before the job is enqueued. var taskRun = await db.TaskRuns.SingleAsync(t => t.Id == taskRunId); Assert.Equal(run.Id, taskRun.RunId); Assert.Equal(TaskRunStatus.Running, taskRun.Status); Assert.Equal(1, taskRun.Attempt); Assert.Equal("{\"input\":1}", taskRun.InputJson); var (tenant, fields) = Assert.Single(jobs.Enqueued); Assert.Equal(tenantId, tenant); Assert.Equal(TaskInvocation.TypeValue, fields["type"]); Assert.Equal(run.Id.ToString(), fields["run_id"]); Assert.Equal("root", fields["task_key"]); Assert.Equal("root.sh", fields["entry_file"]); Assert.EndsWith("workflows", fields["working_dir"]); } }