using Data.SeaHavenIndustries; using SeaHaven.DataServices.Helpers; using SeaHaven.DataServices.Interfaces; using SeaHaven.Services.DTOs; using SeaHaven.Services.Helpers; using SeaHaven.Services.Interfaces; namespace SeaHaven.Services.Implementation { public class WorkOrderIngestService : IWorkOrderIngestService { private readonly IWorkOrderIngestDataService _ingestData; private readonly ISyncFieldMergePolicy _mergePolicy; private readonly IWorkOrderFieldLockService _fieldLocks; private readonly IWorkOrderAuditService _audit; private readonly IWorkOrderAccountResolver _accountResolver; public WorkOrderIngestService( IWorkOrderIngestDataService ingestData, ISyncFieldMergePolicy mergePolicy, IWorkOrderFieldLockService fieldLocks, IWorkOrderAuditService audit, IWorkOrderAccountResolver accountResolver) { _ingestData = ingestData; _mergePolicy = mergePolicy; _fieldLocks = fieldLocks; _audit = audit; _accountResolver = accountResolver; } public async Task UpsertBatchAsync( IReadOnlyList items, CancellationToken cancellationToken = default) { var result = new WorkOrderIngestBatchResultDto(); if (items.Count == 0) return result; // Existing work orders this batch may cancel hold the per-work-order lock until the // batch commits, so an uplift create on one of them either commits first (and is // cancelled below) or waits and then sees the work order cancelled. var cancellingExternalIds = items .Where(item => !string.IsNullOrWhiteSpace(item.ExternalWorkOrderId) && PendingUpliftCancellation.IsCancelled(null, WorkOrderIngestFieldMapper.MapStatus(item.WoStatus))) .Select(item => item.ExternalWorkOrderId) .Distinct() .ToList(); var lockedWorkOrderIds = cancellingExternalIds.Count == 0 ? Array.Empty() : await _ingestData.GetWorkOrderIdsByExternalIdsAsync(cancellingExternalIds, cancellationToken); await _ingestData.ExecuteTransactionalAsync(lockedWorkOrderIds, async cancellationToken => { var nextSeed = await AllocateNextWoSeedAsync(cancellationToken); foreach (var item in items) { if (string.IsNullOrWhiteSpace(item.ExternalWorkOrderId)) { result.Skipped++; continue; } try { var (upsert, updatedSeed) = await UpsertOneAsync(item, nextSeed, cancellationToken); nextSeed = updatedSeed; result.Results.Add(upsert); if (upsert.Created) result.Created++; else result.Updated++; } catch (Exceptions.WorkOrderBoardValidationException ex) when (ex.Code == "AccountUnresolved") { result.Skipped++; } } await _ingestData.SaveAsync(cancellationToken); }, cancellationToken); return result; } private async Task<(WorkOrderIngestItemResultDto Result, int NextSeed)> UpsertOneAsync( WorkOrderIngestPayloadDto item, int nextSeed, CancellationToken cancellationToken) { var existing = await _ingestData.GetTrackedByExternalIdAsync(item.ExternalWorkOrderId, cancellationToken); var locationId = await ResolveLocationIdAsync( item.SiteCode, item.Building, item.Address, cancellationToken); if (existing == null) { var accountId = await _accountResolver.ResolveForUnauthenticatedCreateAsync( item.Customer, cancellationToken); var (woNumber, updatedSeed) = AllocateInternalWoNumber(nextSeed); var wo = new WorkOrder { InternalWONumber = woNumber, ExternalWorkOrderId = item.ExternalWorkOrderId, WorkerOrderNumber = item.ExternalWorkOrderId, WorkerOrderTitle = item.Description, Description = item.Description, Status = WorkOrderIngestFieldMapper.MapStatus(item.WoStatus) ?? "Open", Priority = WorkOrderIngestFieldMapper.MapSeverityToPriority(item.Severity), Severity = item.Severity, Customer = item.Customer, AccountId = accountId, SiteCode = item.SiteCode, Building = item.Building, LocationId = locationId, DueDate = item.DueDate, DateReported = item.DateReported, ScheduledStart = item.ScheduledStart, SourceEmailS3Key = item.SourceEmailS3Key, CreatedDate = item.CreatedAt ?? DateTime.UtcNow, istemplate = false }; _ingestData.TrackWorkOrder(wo); await _ingestData.SaveAsync(cancellationToken); return (new WorkOrderIngestItemResultDto { ExternalWorkOrderId = item.ExternalWorkOrderId, WorkOrderId = wo.Id, Created = true }, updatedSeed); } var syncContext = new WorkOrderSyncContext { WorkOrder = existing, FieldLocks = _fieldLocks, Audit = _audit }; if (item.Description != null) { await _mergePolicy.TryApplyAsync(syncContext, "Description", item.Description); await _mergePolicy.TryApplyAsync(syncContext, "WorkerOrderTitle", item.Description); } var mappedStatus = WorkOrderIngestFieldMapper.MapStatus(item.WoStatus); if (mappedStatus != null) { var statusBefore = existing.Status; await _mergePolicy.TryApplyAsync(syncContext, "Status", mappedStatus); // A sync cancellation cancels the work order's pending uplifts in the same save. if (PendingUpliftCancellation.AppliesToStatusText(statusBefore, existing.Status)) await _ingestData.StageCancelPendingUpliftsAsync(existing.Id, cancellationToken); } var priority = WorkOrderIngestFieldMapper.MapSeverityToPriority(item.Severity); if (priority != null) await _mergePolicy.TryApplyAsync(syncContext, "Priority", priority); if (item.Severity != null) await _mergePolicy.TryApplyAsync(syncContext, "Severity", item.Severity); if (item.SiteCode != null) await _mergePolicy.TryApplyAsync(syncContext, "SiteCode", item.SiteCode); if (item.Building != null) await _mergePolicy.TryApplyAsync(syncContext, "Building", item.Building); if (locationId.HasValue) await _mergePolicy.TryApplyAsync(syncContext, "LocationId", locationId.Value); if (item.DueDate.HasValue) await _mergePolicy.TryApplyAsync(syncContext, "DueDate", item.DueDate.Value); if (item.Customer != null) existing.Customer = item.Customer; // Sync-owned metadata (not in ShocOwnedFieldSet) existing.DateReported = item.DateReported ?? existing.DateReported; existing.ScheduledStart = item.ScheduledStart ?? existing.ScheduledStart; existing.SourceEmailS3Key = item.SourceEmailS3Key ?? existing.SourceEmailS3Key; return (new WorkOrderIngestItemResultDto { ExternalWorkOrderId = item.ExternalWorkOrderId, WorkOrderId = existing.Id, Created = false }, nextSeed); } private async Task AllocateNextWoSeedAsync(CancellationToken cancellationToken) { var maxId = await _ingestData.GetMaxWorkOrderIdAsync(cancellationToken); return Math.Max(maxId + 1, 1); } private static (string Number, int NextSeed) AllocateInternalWoNumber(int seed) { if (!WorkOrderNumberNormalizer.TryNormalize(seed.ToString(), out var normalized, out _)) normalized = seed.ToString().PadLeft(11, '0'); return (normalized, seed + 1); } private async Task ResolveLocationIdAsync( string? siteCode, string? building, string? address, CancellationToken cancellationToken) { if (string.IsNullOrWhiteSpace(siteCode) && string.IsNullOrWhiteSpace(building)) return null; var matchCode = siteCode ?? building; var existing = await _ingestData.FindLocationAsync(matchCode, cancellationToken); if (existing != null) return existing.Id; var location = new Locations { Name = matchCode, Title = building, Address1 = address, Status = "Active" }; await _ingestData.AddAndSaveLocationAsync(location, cancellationToken); return location.Id; } } }