mirror of
https://github.com/Sea-Haven-Industries/shoc-backend.git
synced 2026-09-30 07:13:12 +00:00
The webhook and reconciliation saves now stage the same pending-uplift cancellation, with its own sync audit row, as the board cancel. The rule and the write live in one data-layer helper so the paths cannot drift.
400 lines
15 KiB
C#
400 lines
15 KiB
C#
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<ProcurementWorkOrder>(
|
|
new[] { WorkOrder("2", "new") },
|
|
"next"),
|
|
new ProcurementPage<ProcurementWorkOrder>(
|
|
new[] { WorkOrder("1", "cancelled") },
|
|
null)
|
|
},
|
|
new Dictionary<string, Queue<ProcurementPage<ProcurementWorkOrderComment>>>
|
|
{
|
|
["2"] = new Queue<ProcurementPage<ProcurementWorkOrderComment>>(new[]
|
|
{
|
|
new ProcurementPage<ProcurementWorkOrderComment>(
|
|
new[] { Comment("2", "c2") },
|
|
"comments-next"),
|
|
new ProcurementPage<ProcurementWorkOrderComment>(
|
|
Array.Empty<ProcurementWorkOrderComment>(),
|
|
null)
|
|
}),
|
|
["1"] = new Queue<ProcurementPage<ProcurementWorkOrderComment>>(new[]
|
|
{
|
|
new ProcurementPage<ProcurementWorkOrderComment>(
|
|
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<ProcurementWorkOrder>(
|
|
Array.Empty<ProcurementWorkOrder>(),
|
|
"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<ProcurementWorkOrder>(
|
|
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<ProcurementWorkOrder>(
|
|
new[]
|
|
{
|
|
new ProcurementWorkOrder
|
|
{
|
|
WorkOrderId = "42",
|
|
Description = "Legacy mock"
|
|
}
|
|
},
|
|
null)
|
|
},
|
|
new Dictionary<string, Queue<ProcurementPage<ProcurementWorkOrderComment>>>
|
|
{
|
|
["42"] = new Queue<ProcurementPage<ProcurementWorkOrderComment>>(new[]
|
|
{
|
|
new ProcurementPage<ProcurementWorkOrderComment>(
|
|
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<ApplicationDbContext>()
|
|
.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<ApplicationDbContext>()
|
|
.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);
|
|
}
|
|
|
|
internal static WorkOrderReconciliationService CreateService(
|
|
IReadOnlyList<ProcurementWorkOrder> workOrders,
|
|
IWorkOrderWebhookDataService data) =>
|
|
Create(
|
|
new MockClient(
|
|
new[] { new ProcurementPage<ProcurementWorkOrder>(workOrders, null) },
|
|
new()),
|
|
data,
|
|
new RecordingJobs());
|
|
|
|
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<WorkOrderReconciliationService>.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<ProcurementPage<ProcurementWorkOrder>> _workOrders;
|
|
private readonly Dictionary<string, Queue<ProcurementPage<ProcurementWorkOrderComment>>> _comments;
|
|
|
|
public MockClient(
|
|
IEnumerable<ProcurementPage<ProcurementWorkOrder>> workOrders,
|
|
Dictionary<string, Queue<ProcurementPage<ProcurementWorkOrderComment>>> comments)
|
|
{
|
|
_workOrders = new Queue<ProcurementPage<ProcurementWorkOrder>>(workOrders);
|
|
_comments = comments;
|
|
}
|
|
|
|
public Task<ProcurementPage<ProcurementWorkOrder>> GetWorkOrdersAsync(
|
|
string? cursor,
|
|
int limit,
|
|
CancellationToken cancellationToken) =>
|
|
Task.FromResult(_workOrders.Dequeue());
|
|
|
|
public Task<ProcurementPage<ProcurementWorkOrderComment>> GetCommentsAsync(
|
|
string workOrderId,
|
|
string? cursor,
|
|
int limit,
|
|
CancellationToken cancellationToken) =>
|
|
Task.FromResult(_comments.TryGetValue(workOrderId, out var pages)
|
|
? pages.Dequeue()
|
|
: new ProcurementPage<ProcurementWorkOrderComment>(
|
|
Array.Empty<ProcurementWorkOrderComment>(),
|
|
null));
|
|
}
|
|
|
|
private sealed class RecordingWorkOrders : IWorkOrderWebhookDataService
|
|
{
|
|
public List<WorkOrderWebhookMutation> Items { get; } = new();
|
|
|
|
public Task<WorkOrderWebhookPersistenceResult> 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<ReconciliationJobSnapshot> EnqueueAsync(
|
|
string reason,
|
|
DateTimeOffset now,
|
|
CancellationToken cancellationToken) =>
|
|
Task.FromResult(Snapshot("Pending"));
|
|
|
|
public Task<ReconciliationLease?> TryAcquirePendingAsync(
|
|
DateTimeOffset now,
|
|
TimeSpan leaseDuration,
|
|
CancellationToken cancellationToken) =>
|
|
Task.FromResult<ReconciliationLease?>(new(_run, _fence));
|
|
|
|
public Task<ReconciliationJobSnapshot> GetStatusAsync(CancellationToken cancellationToken) =>
|
|
Task.FromResult(Snapshot("Pending"));
|
|
|
|
public Task<bool> 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<WorkOrderReconciliationOptions>
|
|
{
|
|
public TestOptions(WorkOrderReconciliationOptions value) => CurrentValue = value;
|
|
public WorkOrderReconciliationOptions CurrentValue { get; }
|
|
public WorkOrderReconciliationOptions Get(string? name) => CurrentValue;
|
|
public IDisposable? OnChange(
|
|
Action<WorkOrderReconciliationOptions, string?> listener) => null;
|
|
}
|
|
}
|