using Data.SeaHavenIndustries; 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; await _ingestData.ExecuteTransactionalAsync(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) await _mergePolicy.TryApplyAsync(syncContext, "Status", mappedStatus); 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; } } }