using Data.SeaHavenIndustries; using Microsoft.EntityFrameworkCore; using SeaHaven.DataServices.Interfaces; namespace SeaHaven.DataServices.Implementation { public sealed class WorkOrderWebhookDataService : IWorkOrderWebhookDataService { private static readonly SemaphoreSlim InMemoryNumberLock = new(1, 1); private readonly ApplicationDbContext _context; public WorkOrderWebhookDataService(ApplicationDbContext context) { _context = context; } public async Task ApplyAsync( WorkOrderWebhookMutation mutation, CancellationToken cancellationToken) { var prior = await _context.WorkOrderWebhookDeliveries .AsNoTracking() .SingleOrDefaultAsync(d => d.DeliveryId == mutation.DeliveryId, cancellationToken); if (prior != null) return ExistingDeliveryResult(prior, mutation.BodySha256); var workOrder = await _context.workOrders .SingleOrDefaultAsync( w => w.ExternalWorkOrderId == mutation.ExternalWorkOrderId, cancellationToken); if (workOrder == null) { var internalNumber = await AllocateInternalNumberAsync(cancellationToken); workOrder = new WorkOrder { ExternalWorkOrderId = mutation.ExternalWorkOrderId, ExternalSource = mutation.Source, WorkerOrderNumber = mutation.WorkerOrderNumber, InternalWONumber = internalNumber.ToString("D11"), WorkerOrderTitle = mutation.Title ?? mutation.Description ?? $"Imported work order {mutation.ExternalWorkOrderId}", Description = mutation.Description, Source = mutation.Source, istemplate = false, CreatedDate = mutation.CreatedAt ?? mutation.ProcessedAt.UtcDateTime }; _context.workOrders.Add(workOrder); } var staleState = mutation.IsStateEvent && !IsNewer( mutation.UpdatedAt, mutation.VersionHash, workOrder.ExternalUpdatedAt, workOrder.ExternalVersionHash); if (mutation.IsStateEvent && !staleState) { var locationKey = mutation.SiteCode ?? mutation.Building; if (!string.IsNullOrWhiteSpace(locationKey)) { var locationIsNewer = await UpsertReceiptAsync( mutation.Source, "location", locationKey, mutation.UpdatedAt, mutation.VersionHash, mutation.ProcessedAt, cancellationToken); var location = await _context.Locations .FirstOrDefaultAsync( l => l.ExternalSource == mutation.Source && l.ExternalLocationId == locationKey, cancellationToken); if (location == null) { location = new Locations { ExternalSource = mutation.Source, ExternalLocationId = locationKey, Name = locationKey, Title = mutation.Building, Address1 = mutation.Address, Status = "Active" }; _context.Locations.Add(location); } else if (locationIsNewer) { location.Name = locationKey; location.Title = mutation.Building; location.Address1 = mutation.Address; location.Status = "Active"; } workOrder.Locations = location; } workOrder.WorkerOrderNumber = mutation.WorkerOrderNumber; workOrder.WorkerOrderTitle = mutation.Title ?? mutation.Description; workOrder.Description = mutation.Description; workOrder.Status = mutation.IsCancelled ? "Cancelled" : mutation.Status; if (mutation.LifecycleStatus.HasValue) workOrder.LifecycleStatus = mutation.LifecycleStatus.Value; workOrder.Severity = mutation.Severity; workOrder.Priority = mutation.Priority; workOrder.ExternalAssignedTo = mutation.AssignedTo; workOrder.ExternalRecordType = mutation.RecordType; workOrder.Customer = mutation.Customer; workOrder.SiteCode = mutation.SiteCode; workOrder.Building = mutation.Building; workOrder.DueDate = mutation.DueDate; workOrder.DateReported = mutation.DateReported; workOrder.ScheduledStart = mutation.ScheduledStart; workOrder.Source = mutation.Source; workOrder.SourceEmailS3Key = mutation.SourceEmailS3Key; workOrder.ExternalSource = mutation.Source; workOrder.istemplate = false; workOrder.ExternalLastOccurredAt = mutation.OccurredAt; workOrder.ExternalUpdatedAt = mutation.UpdatedAt; workOrder.ExternalVersionHash = mutation.VersionHash; } if (mutation.CommentId != null) { var receiptIsNewer = await UpsertReceiptAsync( mutation.Source, "comment", mutation.CommentId, mutation.UpdatedAt, mutation.VersionHash, mutation.ProcessedAt, cancellationToken); var comment = await _context.Comments .SingleOrDefaultAsync( c => c.ExternalSource == mutation.Source && c.ExternalCommentId == mutation.CommentId, cancellationToken); if (comment == null) { comment = new Comments { ExternalSource = mutation.Source, ExternalCommentId = mutation.CommentId, RecordType = "WorkOrder", WorkOrder = workOrder, CreatedDate = mutation.OccurredAt.UtcDateTime, }; _context.Comments.Add(comment); } if (receiptIsNewer) { comment.Commenttext = mutation.CommentText; comment.Commenter = mutation.Commenter; comment.CommentType = mutation.CommentType; comment.ExternalUpdatedAt = mutation.UpdatedAt; comment.ExternalVersionHash = mutation.VersionHash; comment.ExternalSourceEmailS3Key = mutation.SourceEmailS3Key; } } _context.WorkOrderWebhookDeliveries.Add(new WorkOrderWebhookDelivery { DeliveryId = mutation.DeliveryId, EventType = mutation.EventType, OccurredAt = mutation.OccurredAt, ProcessedAt = mutation.ProcessedAt, BodySha256 = mutation.BodySha256 }); try { await _context.SaveChangesAsync(cancellationToken); return new WorkOrderWebhookPersistenceResult( WorkOrderWebhookPersistenceStatus.Applied, staleState); } catch (DbUpdateException) { _context.ChangeTracker.Clear(); var concurrent = await _context.WorkOrderWebhookDeliveries .AsNoTracking() .SingleOrDefaultAsync(d => d.DeliveryId == mutation.DeliveryId, cancellationToken); if (concurrent == null) throw; return ExistingDeliveryResult(concurrent, mutation.BodySha256); } } private async Task AllocateInternalNumberAsync(CancellationToken cancellationToken) { if (_context.Database.IsSqlServer()) { return await _context.Database .SqlQueryRaw( "SELECT NEXT VALUE FOR dbo.WorkOrderInternalNumberSequence AS [Value]") .SingleAsync(cancellationToken); } await InMemoryNumberLock.WaitAsync(cancellationToken); try { var values = await _context.workOrders .AsNoTracking() .Where(w => w.InternalWONumber != null) .Select(w => w.InternalWONumber!) .ToListAsync(cancellationToken); return values .Select(v => long.TryParse(v, out var parsed) ? parsed : 0L) .DefaultIfEmpty() .Max() + 1; } finally { InMemoryNumberLock.Release(); } } private async Task UpsertReceiptAsync( string source, string kind, string externalId, DateTimeOffset updatedAt, string versionHash, DateTimeOffset processedAt, CancellationToken cancellationToken) { var receipt = await _context.WorkOrderExternalReceipts.SingleOrDefaultAsync( r => r.Source == source && r.Kind == kind && r.ExternalId == externalId, cancellationToken); if (receipt == null) { _context.WorkOrderExternalReceipts.Add(new WorkOrderExternalReceipt { Source = source, Kind = kind, ExternalId = externalId, UpdatedAt = updatedAt, VersionHash = versionHash, ProcessedAt = processedAt }); return true; } if (!IsNewer(updatedAt, versionHash, receipt.UpdatedAt, receipt.VersionHash)) return false; receipt.UpdatedAt = updatedAt; receipt.VersionHash = versionHash; receipt.ProcessedAt = processedAt; return true; } private static bool IsNewer( DateTimeOffset candidateUpdatedAt, string candidateHash, DateTimeOffset? currentUpdatedAt, string? currentHash) { if (!currentUpdatedAt.HasValue) return true; var timestampComparison = candidateUpdatedAt.CompareTo(currentUpdatedAt.Value); return timestampComparison > 0 || (timestampComparison == 0 && string.CompareOrdinal(candidateHash, currentHash ?? string.Empty) > 0); } private static WorkOrderWebhookPersistenceResult ExistingDeliveryResult( WorkOrderWebhookDelivery delivery, string bodySha256) { return new WorkOrderWebhookPersistenceResult( string.Equals(delivery.BodySha256, bodySha256, StringComparison.OrdinalIgnoreCase) ? WorkOrderWebhookPersistenceStatus.Duplicate : WorkOrderWebhookPersistenceStatus.HashConflict); } } }