using Data.SeaHavenIndustries; using Data.SeaHavenIndustries.Enums; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Logging.Abstractions; using Microsoft.Extensions.Options; using SeaHaven.DataServices.Implementation; using SeaHaven.DataServices.Interfaces; using SeaHaven.Services.Configuration; using SeaHaven.Services.Implementation; using SeaHaven.Services.Interfaces; namespace SeaHavenIndustries.Tests; public sealed class WorkOrderReconciliationTests { [Fact] public async Task Full_scan_paginates_work_orders_and_each_comment_collection() { var client = new MockClient( new[] { new ProcurementPage( new[] { WorkOrder("2", "new") }, "next"), new ProcurementPage( new[] { WorkOrder("1", "cancelled") }, null) }, new Dictionary>> { ["2"] = new Queue>(new[] { new ProcurementPage( new[] { Comment("2", "c2") }, "comments-next"), new ProcurementPage( Array.Empty(), null) }), ["1"] = new Queue>(new[] { new ProcurementPage( new[] { Comment("1", "c1") }, null) }) }); var mutations = new RecordingWorkOrders(); var jobs = new RecordingJobs(); var service = Create(client, mutations, jobs); Assert.True(await service.RunPendingAsync(CancellationToken.None)); Assert.Equal(4, mutations.Items.Count); Assert.Contains(mutations.Items, m => m.ExternalWorkOrderId == "1" && m.IsCancelled); Assert.Contains(mutations.Items, m => m.ExternalWorkOrderId == "1" && m.LifecycleStatus == LifecycleStatus.Canceled); Assert.Contains(mutations.Items, m => m.CommentId == "c2"); Assert.Equal(2, jobs.CompletedWorkOrders); Assert.Equal(2, jobs.CompletedComments); Assert.True(jobs.RenewCount >= 3); } [Fact] public async Task Repeated_cursor_fails_bounded_run_instead_of_looping() { var page = new ProcurementPage( Array.Empty(), "same"); var service = Create( new MockClient(new[] { page, page }, new()), new RecordingWorkOrders(), new RecordingJobs()); Assert.False(await service.RunPendingAsync(CancellationToken.None)); } [Fact] public async Task Lost_lease_after_fetch_prevents_work_order_persistence() { var mutations = new RecordingWorkOrders(); var jobs = new RecordingJobs { LoseLeaseOnRenewal = true }; var service = Create( new MockClient( new[] { new ProcurementPage( new[] { WorkOrder("1", "new") }, null) }, new()), mutations, jobs); Assert.False(await service.RunPendingAsync(CancellationToken.None)); Assert.Empty(mutations.Items); } [Fact] public async Task Legacy_records_without_timestamps_use_stable_epoch_freshness() { var client = new MockClient( new[] { new ProcurementPage( new[] { new ProcurementWorkOrder { WorkOrderId = "42", Description = "Legacy mock" } }, null) }, new Dictionary>> { ["42"] = new Queue>(new[] { new ProcurementPage( new[] { new ProcurementWorkOrderComment { WorkOrderId = "42", CommentId = "legacy-comment", Text = "Legacy comment" } }, null) }) }); var mutations = new RecordingWorkOrders(); var service = Create(client, mutations, new RecordingJobs()); Assert.True(await service.RunPendingAsync(CancellationToken.None)); Assert.Equal(2, mutations.Items.Count); Assert.All(mutations.Items, mutation => Assert.Equal(DateTimeOffset.UnixEpoch, mutation.UpdatedAt)); Assert.All(mutations.Items, mutation => Assert.StartsWith("reconcile:", mutation.DeliveryId)); } [Fact] public async Task Durable_job_survives_scope_restart_and_rejects_stale_fence() { var database = $"reconciliation-{Guid.NewGuid()}"; var options = new DbContextOptionsBuilder() .UseInMemoryDatabase(database) .Options; Guid runId; Guid fence; await using (var context = new ApplicationDbContext(options)) { var idProperty = context.Model .FindEntityType(typeof(WorkOrderReconciliationJob))! .FindProperty(nameof(WorkOrderReconciliationJob.Id))!; Assert.Equal( Microsoft.EntityFrameworkCore.Metadata.ValueGenerated.Never, idProperty.ValueGenerated); var data = new WorkOrderReconciliationDataService(context); var queued = await data.EnqueueAsync( "admin", DateTimeOffset.UtcNow, CancellationToken.None); runId = queued.RunId!.Value; } await using (var context = new ApplicationDbContext(options)) { var data = new WorkOrderReconciliationDataService(context); var lease = await data.TryAcquirePendingAsync( DateTimeOffset.UtcNow, TimeSpan.FromMinutes(2), CancellationToken.None); fence = lease!.FenceToken; Assert.Equal(runId, lease.RunId); await data.CompleteAsync( runId, Guid.NewGuid(), DateTimeOffset.UtcNow, 99, 99, CancellationToken.None); Assert.Equal("Running", (await data.GetStatusAsync(CancellationToken.None)).State); await data.CompleteAsync( runId, fence, DateTimeOffset.UtcNow, 2, 3, CancellationToken.None); var status = await data.GetStatusAsync(CancellationToken.None); Assert.Equal("Succeeded", status.State); Assert.Equal(2, status.WorkOrdersProcessed); Assert.Equal(3, status.CommentsProcessed); } } [Fact] public async Task Equal_timestamp_converges_to_lexicographically_greater_version() { var options = new DbContextOptionsBuilder() .UseInMemoryDatabase($"convergence-{Guid.NewGuid()}") .Options; await using var context = new ApplicationDbContext(options); context.Accounts.Add(new Accounts { Id = 1, Name = "Recon Customer", IsDeleted = false }); await context.SaveChangesAsync(); var data = new WorkOrderWebhookDataService(context); var time = DateTimeOffset.Parse("2026-07-24T12:00:00Z"); await data.ApplyAsync(Mutation("d1", "aaa", "older-tie", time), CancellationToken.None); await data.ApplyAsync(Mutation("d2", "zzz", "winner", time), CancellationToken.None); await data.ApplyAsync(Mutation("d3", "aaa", "loser-replay", time), CancellationToken.None); Assert.Equal("winner", (await context.workOrders.SingleAsync()).WorkerOrderTitle); } private static WorkOrderReconciliationService Create( IProcurementWorkOrderClient client, IWorkOrderWebhookDataService workOrders, IWorkOrderReconciliationDataService jobs) => new( client, workOrders, jobs, new TestOptions(new WorkOrderReconciliationOptions { Enabled = true, MaxPages = 10, PageSize = 100, LeaseSeconds = 120 }), TimeProvider.System, NullLogger.Instance); private static ProcurementWorkOrder WorkOrder(string id, string status) => new() { WorkOrderId = id, WoStatus = status, Description = $"Mock {id}", UpdatedAt = DateTimeOffset.Parse("2026-07-24T12:00:00Z") }; private static ProcurementWorkOrderComment Comment(string workOrderId, string commentId) => new() { WorkOrderId = workOrderId, CommentId = commentId, Text = "Mock comment", IngestedAt = DateTimeOffset.Parse("2026-07-24T12:01:00Z") }; private static WorkOrderWebhookMutation Mutation( string delivery, string versionHash, string title, DateTimeOffset updatedAt) => new() { DeliveryId = delivery, EventType = "work_order.updated", OccurredAt = updatedAt, UpdatedAt = updatedAt, ProcessedAt = updatedAt, BodySha256 = delivery, VersionHash = versionHash, ExternalWorkOrderId = "123", WorkerOrderNumber = "123", Source = "procurement", IsStateEvent = true, Title = title, Customer = "Recon Customer" }; private sealed class MockClient : IProcurementWorkOrderClient { private readonly Queue> _workOrders; private readonly Dictionary>> _comments; public MockClient( IEnumerable> workOrders, Dictionary>> comments) { _workOrders = new Queue>(workOrders); _comments = comments; } public Task> GetWorkOrdersAsync( string? cursor, int limit, CancellationToken cancellationToken) => Task.FromResult(_workOrders.Dequeue()); public Task> GetCommentsAsync( string workOrderId, string? cursor, int limit, CancellationToken cancellationToken) => Task.FromResult(_comments.TryGetValue(workOrderId, out var pages) ? pages.Dequeue() : new ProcurementPage( Array.Empty(), null)); } private sealed class RecordingWorkOrders : IWorkOrderWebhookDataService { public List Items { get; } = new(); public Task ApplyAsync( WorkOrderWebhookMutation mutation, CancellationToken cancellationToken) { Items.Add(mutation); return Task.FromResult(new WorkOrderWebhookPersistenceResult( WorkOrderWebhookPersistenceStatus.Applied)); } } private sealed class RecordingJobs : IWorkOrderReconciliationDataService { private readonly Guid _run = Guid.NewGuid(); private readonly Guid _fence = Guid.NewGuid(); public int CompletedWorkOrders { get; private set; } public int CompletedComments { get; private set; } public int RenewCount { get; private set; } public bool LoseLeaseOnRenewal { get; init; } public Task EnqueueAsync( string reason, DateTimeOffset now, CancellationToken cancellationToken) => Task.FromResult(Snapshot("Pending")); public Task TryAcquirePendingAsync( DateTimeOffset now, TimeSpan leaseDuration, CancellationToken cancellationToken) => Task.FromResult(new(_run, _fence)); public Task GetStatusAsync(CancellationToken cancellationToken) => Task.FromResult(Snapshot("Pending")); public Task RenewLeaseAsync( Guid runId, Guid fenceToken, DateTimeOffset now, TimeSpan leaseDuration, CancellationToken cancellationToken) { RenewCount++; return Task.FromResult(!LoseLeaseOnRenewal); } public Task CompleteAsync( Guid runId, Guid fenceToken, DateTimeOffset completedAt, int workOrders, int comments, CancellationToken cancellationToken) { CompletedWorkOrders = workOrders; CompletedComments = comments; return Task.CompletedTask; } public Task FailAsync( Guid runId, Guid fenceToken, DateTimeOffset completedAt, string errorCode, CancellationToken cancellationToken) => Task.CompletedTask; private ReconciliationJobSnapshot Snapshot(string state) => new(_run, state, "test", null, null, null, 0, 0, null); } private sealed class TestOptions : IOptionsMonitor { public TestOptions(WorkOrderReconciliationOptions value) => CurrentValue = value; public WorkOrderReconciliationOptions CurrentValue { get; } public WorkOrderReconciliationOptions Get(string? name) => CurrentValue; public IDisposable? OnChange( Action listener) => null; } }