w4c-workflows-api/w4c-workflows-api.Tests/TaskDispatcherTests.cs

55 lines
2.1 KiB
C#

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<TaskDispatcher>.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"]);
}
}