using System.Security.Claims; using System.Security.Cryptography; using System.Text; using Microsoft.Extensions.Logging; using Microsoft.Extensions.Options; using SeaHaven.DataServices.Interfaces; using SeaHaven.Services.Configuration; using SeaHaven.Services.Helpers; using SeaHaven.Services.Interfaces; namespace SeaHaven.Services.Implementation { public sealed class WorkOrderReconciliationService : IWorkOrderReconciliationService, IWorkOrderReconciliationRunner { private const string Source = "procurement"; private readonly IProcurementWorkOrderClient _client; private readonly IWorkOrderWebhookDataService _workOrders; private readonly IWorkOrderReconciliationDataService _jobs; private readonly IOptionsMonitor _options; private readonly TimeProvider _timeProvider; private readonly ILogger _logger; public WorkOrderReconciliationService( IProcurementWorkOrderClient client, IWorkOrderWebhookDataService workOrders, IWorkOrderReconciliationDataService jobs, IOptionsMonitor options, TimeProvider timeProvider, ILogger logger) { _client = client; _workOrders = workOrders; _jobs = jobs; _options = options; _timeProvider = timeProvider; _logger = logger; } public async Task TriggerAsync( ClaimsPrincipal user, string reason, CancellationToken cancellationToken) { RequireAdmin(user); return await TriggerAsync(reason, cancellationToken); } public async Task GetStatusAsync( ClaimsPrincipal user, CancellationToken cancellationToken) { RequireAdmin(user); var status = await _jobs.GetStatusAsync(cancellationToken); return new WorkOrderReconciliationStatus( status.RunId, status.State, status.Reason, status.RequestedAt, status.StartedAt, status.CompletedAt, status.WorkOrdersProcessed, status.CommentsProcessed, status.ErrorCode); } public async Task TriggerAsync( string reason, CancellationToken cancellationToken) { if (!_options.CurrentValue.Enabled) return new WorkOrderReconciliationTriggerResult(false, Guid.Empty); var before = await _jobs.GetStatusAsync(cancellationToken); var queued = await _jobs.EnqueueAsync(reason, _timeProvider.GetUtcNow(), cancellationToken); return new WorkOrderReconciliationTriggerResult( before.State is not ("Pending" or "Running"), queued.RunId ?? Guid.Empty); } public async Task RunPendingAsync(CancellationToken cancellationToken) { var options = _options.CurrentValue; if (!options.Enabled) return false; var leaseDuration = TimeSpan.FromSeconds(options.LeaseSeconds); var lease = await _jobs.TryAcquirePendingAsync( _timeProvider.GetUtcNow(), leaseDuration, cancellationToken); if (lease == null) return false; try { var (workOrders, comments) = await ReconcileAsync( lease, options, leaseDuration, cancellationToken); await _jobs.CompleteAsync( lease.RunId, lease.FenceToken, _timeProvider.GetUtcNow(), workOrders, comments, cancellationToken); return true; } catch (OperationCanceledException) when (cancellationToken.IsCancellationRequested) { throw; } catch (Exception ex) { _logger.LogError( ex, "Procurement work-order reconciliation run {RunId} failed.", lease.RunId); await _jobs.FailAsync( lease.RunId, lease.FenceToken, _timeProvider.GetUtcNow(), "reconciliation_failed", cancellationToken); return false; } } private static void RequireAdmin(ClaimsPrincipal user) { if (user is null || !(user.Identity?.IsAuthenticated ?? false) || !user.IsInRole("Admin")) { throw new UnauthorizedAccessException(); } } private async Task<(int WorkOrders, int Comments)> ReconcileAsync( ReconciliationLease lease, WorkOrderReconciliationOptions options, TimeSpan leaseDuration, CancellationToken cancellationToken) { var cursor = (string?)null; var seenCursors = new HashSet(StringComparer.Ordinal); var workOrderCount = 0; var commentCount = 0; for (var pageNumber = 0; pageNumber < options.MaxPages; pageNumber++) { var page = await _client.GetWorkOrdersAsync(cursor, options.PageSize, cancellationToken); foreach (var item in page.Items) { await _workOrders.ApplyAsync(Map(item), cancellationToken); workOrderCount++; commentCount += await ReconcileCommentsAsync( item.WorkOrderId, lease, options, leaseDuration, cancellationToken); } await EnsureLeaseAsync(lease, leaseDuration, cancellationToken); cursor = ValidateNextCursor(page.NextCursor, options.MaxCursorLength, seenCursors); if (cursor == null) return (workOrderCount, commentCount); } throw new InvalidOperationException("Procurement work-order page limit exceeded."); } private async Task ReconcileCommentsAsync( string workOrderId, ReconciliationLease lease, WorkOrderReconciliationOptions options, TimeSpan leaseDuration, CancellationToken cancellationToken) { string? cursor = null; var seenCursors = new HashSet(StringComparer.Ordinal); var count = 0; for (var pageNumber = 0; pageNumber < options.MaxPages; pageNumber++) { var page = await _client.GetCommentsAsync( workOrderId, cursor, options.PageSize, cancellationToken); foreach (var comment in page.Items) { await _workOrders.ApplyAsync(Map(comment), cancellationToken); count++; } await EnsureLeaseAsync(lease, leaseDuration, cancellationToken); cursor = ValidateNextCursor(page.NextCursor, options.MaxCursorLength, seenCursors); if (cursor == null) return count; } throw new InvalidOperationException("Procurement comment page limit exceeded."); } private async Task EnsureLeaseAsync( ReconciliationLease lease, TimeSpan leaseDuration, CancellationToken cancellationToken) { if (!await _jobs.RenewLeaseAsync( lease.RunId, lease.FenceToken, _timeProvider.GetUtcNow(), leaseDuration, cancellationToken)) { throw new InvalidOperationException("Reconciliation lease was lost."); } } private WorkOrderWebhookMutation Map(ProcurementWorkOrder item) { var updatedAt = item.UpdatedAt ?? item.CreatedAt ?? DateTimeOffset.UnixEpoch; var hash = WorkOrderExternalVersion.Compute( item.WorkOrderId, item.WoStatus, item.Description, item.Customer, item.SiteCode, item.Building, item.Address, item.Severity, item.Priority, item.AssignedTo, Format(item.DateReported), Format(item.ScheduledStart), Format(item.DueDate), item.RecordType, Format(item.CreatedAt), item.SourceEmailS3Key, null, null, null, null); return new WorkOrderWebhookMutation { DeliveryId = ReceiptId("work-order", item.WorkOrderId, updatedAt, hash), EventType = "reconciliation.work_order", OccurredAt = item.UpdatedAt ?? item.CreatedAt ?? updatedAt, UpdatedAt = updatedAt, ProcessedAt = _timeProvider.GetUtcNow(), BodySha256 = hash, VersionHash = hash, ExternalWorkOrderId = item.WorkOrderId, WorkerOrderNumber = item.WorkOrderId, Source = Source, IsStateEvent = true, IsCancelled = item.WoStatus == "cancelled" || item.RecordType == "cancellation", Description = item.Description, Status = WorkOrderIngestFieldMapper.MapStatus(item.WoStatus), Severity = item.Severity, Priority = item.Priority ?? WorkOrderIngestFieldMapper.MapSeverityToPriority(item.Severity), AssignedTo = item.AssignedTo, RecordType = item.RecordType, SourceEmailS3Key = item.SourceEmailS3Key, Customer = item.Customer, SiteCode = item.SiteCode, Building = item.Building, Address = item.Address, DueDate = item.DueDate?.UtcDateTime, DateReported = item.DateReported?.UtcDateTime, ScheduledStart = item.ScheduledStart?.UtcDateTime, CreatedAt = item.CreatedAt?.UtcDateTime }; } private WorkOrderWebhookMutation Map(ProcurementWorkOrderComment item) { var updatedAt = item.IngestedAt ?? item.CreatedAt ?? DateTimeOffset.UnixEpoch; var hash = WorkOrderExternalVersion.Compute( item.WorkOrderId, item.CommentId, item.RecordType, item.Commenter, item.Text, Format(item.CreatedAt), Format(item.IngestedAt), item.SourceEmailS3Key); return new WorkOrderWebhookMutation { DeliveryId = ReceiptId("comment", item.CommentId, updatedAt, hash), EventType = "reconciliation.comment", OccurredAt = item.CreatedAt ?? updatedAt, UpdatedAt = updatedAt, ProcessedAt = _timeProvider.GetUtcNow(), BodySha256 = hash, VersionHash = hash, ExternalWorkOrderId = item.WorkOrderId, WorkerOrderNumber = item.WorkOrderId, Source = Source, IsStateEvent = false, CommentId = item.CommentId, CommentText = item.Text, Commenter = item.Commenter, CommentType = item.RecordType, SourceEmailS3Key = item.SourceEmailS3Key }; } private static string? ValidateNextCursor( string? cursor, int maxLength, HashSet seen) { if (cursor == null) return null; if (cursor.Length == 0 || cursor.Length > maxLength || !seen.Add(cursor)) throw new InvalidOperationException("Procurement API returned an invalid cursor."); return cursor; } private static string ReceiptId( string kind, string externalId, DateTimeOffset updatedAt, string hash) { var raw = $"{kind}:{externalId}:{updatedAt.ToUniversalTime():O}:{hash}"; var digest = Convert.ToHexString(SHA256.HashData(Encoding.UTF8.GetBytes(raw))) .ToLowerInvariant(); return $"reconcile:{digest}"; } private static string? Format(DateTimeOffset? value) => value?.ToUniversalTime().ToString("O"); } }