using Data.SeaHavenIndustries; using SeaHaven.DataServices.Interfaces; using SeaHaven.Services.DTOs; using SeaHaven.Services.Interfaces; namespace SeaHaven.Services.Implementation { public class SyncService : ISyncService { private readonly ISyncDataService _data; private readonly ISyncExternalSource _external; private readonly IWorkOrderAccountResolver _accountResolver; public SyncService( ISyncDataService data, ISyncExternalSource external, IWorkOrderAccountResolver accountResolver) { _data = data; _external = external; _accountResolver = accountResolver; } public async Task SyncWorkOrdersAsync(CancellationToken cancellationToken) { var synced = 0; var created = 0; var updated = 0; var lastInternalWO = await _data.GetLastInternalWONumberAsync(cancellationToken); int nextInternal = 10000001; if (lastInternalWO != null && int.TryParse(lastInternalWO, out var parsed)) nextInternal = parsed + 1; await foreach (var item in _external.ScanWorkOrdersAsync(cancellationToken)) { var externalId = Get(item, "work_order_id"); if (string.IsNullOrEmpty(externalId)) continue; var existing = await _data.GetWorkOrderByExternalIdAsync(externalId, cancellationToken); var locationId = await _data.ResolveLocationAsync( Get(item, "site_code"), Get(item, "building"), Get(item, "address"), cancellationToken); if (existing == null) { var customer = Get(item, "customer"); var accountId = await _accountResolver.TryResolveFromCustomerAsync(customer, cancellationToken); if (accountId is not int resolvedAccountId) continue; var wo = new WorkOrder { InternalWONumber = (nextInternal++).ToString("D8"), ExternalWorkOrderId = externalId, WorkerOrderNumber = externalId, WorkerOrderTitle = Get(item, "description"), Description = Get(item, "description"), Status = MapStatus(Get(item, "wo_status")), Priority = MapSeverityToPriority(Get(item, "severity")), Severity = Get(item, "severity"), Customer = customer, AccountId = resolvedAccountId, SiteCode = Get(item, "site_code"), Building = Get(item, "building"), LocationId = locationId, DueDate = ParseDate(Get(item, "due_date")), DateReported = ParseDate(Get(item, "date_reported")), ScheduledStart = ParseDate(Get(item, "scheduled_start")), AssignTo = null, SourceEmailS3Key = Get(item, "source_email_s3_key"), CreatedDate = ParseDate(Get(item, "created_at")) ?? DateTime.UtcNow, istemplate = false }; _data.EnqueueWorkOrder(wo); created++; } else { existing.WorkerOrderTitle = Get(item, "description") ?? existing.WorkerOrderTitle; existing.Description = Get(item, "description") ?? existing.Description; existing.Status = MapStatus(Get(item, "wo_status")) ?? existing.Status; existing.Priority = MapSeverityToPriority(Get(item, "severity")) ?? existing.Priority; existing.Severity = Get(item, "severity") ?? existing.Severity; existing.SiteCode = Get(item, "site_code") ?? existing.SiteCode; existing.Building = Get(item, "building") ?? existing.Building; existing.LocationId = locationId ?? existing.LocationId; existing.DueDate = ParseDate(Get(item, "due_date")) ?? existing.DueDate; existing.DateReported = ParseDate(Get(item, "date_reported")) ?? existing.DateReported; existing.ScheduledStart = ParseDate(Get(item, "scheduled_start")) ?? existing.ScheduledStart; existing.SourceEmailS3Key = Get(item, "source_email_s3_key") ?? existing.SourceEmailS3Key; updated++; } synced++; if (synced % 100 == 0) await _data.SaveChangesAsync(cancellationToken); } await _data.SaveChangesAsync(cancellationToken); return new SyncWorkOrdersResult { Synced = synced, Created = created, Updated = updated, LocationsCreated = await _data.GetLocationCountAsync(cancellationToken) }; } public async Task SyncCommentsAsync(CancellationToken cancellationToken) { var synced = 0; var created = 0; var skipped = 0; await foreach (var item in _external.ScanCommentsAsync(cancellationToken)) { var externalCommentId = Get(item, "comment_id"); var externalWoId = Get(item, "work_order_id"); if (string.IsNullOrEmpty(externalCommentId) || string.IsNullOrEmpty(externalWoId)) continue; var alreadyExists = await _data.CommentExistsByExternalIdAsync(externalCommentId, cancellationToken); if (alreadyExists) { skipped++; synced++; continue; } var workOrder = await _data.GetWorkOrderByExternalIdAsync(externalWoId, cancellationToken); if (workOrder == null) { skipped++; continue; } var comment = new Comments { ExternalCommentId = externalCommentId, WorkerOrderId = workOrder.Id, Commenttext = Get(item, "text"), Commenter = Get(item, "commenter"), RecordType = Get(item, "record_type"), CommentType = "customer", CreatedDate = ParseDate(Get(item, "created_at")) ?? DateTime.UtcNow, }; _data.EnqueueComment(comment); created++; synced++; if (synced % 100 == 0) await _data.SaveChangesAsync(cancellationToken); } await _data.SaveChangesAsync(cancellationToken); return new SyncCommentsResult { Synced = synced, Created = created, Skipped = skipped }; } public async Task BackfillInternalWONumbersAsync(CancellationToken cancellationToken) { var wosMissing = await _data.GetWorkOrdersMissingInternalNumberAsync(cancellationToken); if (wosMissing.Count == 0) return new BackfillResult { Updated = 0 }; var lastWO = await _data.GetLastInternalWONumberAsync(cancellationToken); int next = 10000001; if (lastWO != null && int.TryParse(lastWO, out var lastNum)) next = lastNum + 1; foreach (var wo in wosMissing) { wo.InternalWONumber = next.ToString("D8"); next++; } await _data.SaveChangesAsync(cancellationToken); return new BackfillResult { Updated = wosMissing.Count }; } public async Task BackfillDispatchNumbersAsync(CancellationToken cancellationToken) { var missing = await _data.GetDispatchesMissingNumberAsync(cancellationToken); if (missing.Count == 0) return new BackfillResult { Updated = 0 }; int next = 1; var last = await _data.GetLastDispatchNumberAsync(cancellationToken); if (last != null && last.StartsWith("DSP-") && int.TryParse(last.Substring(4), out var n)) next = n + 1; foreach (var d in missing) { d.DispatchNumber = $"DSP-{next:D5}"; next++; } await _data.SaveChangesAsync(cancellationToken); return new BackfillResult { Updated = missing.Count }; } public async Task SyncVendorRepliesAsync(CancellationToken cancellationToken) { var synced = 0; var created = 0; var skipped = 0; await foreach (var item in _external.ScanVendorRepliesAsync(cancellationToken)) { var replyId = Get(item, "reply_id"); var woNumber = Get(item, "internal_wo_number"); var dispatchNumber = Get(item, "dispatch_number"); if (string.IsNullOrEmpty(replyId) || string.IsNullOrEmpty(woNumber)) continue; var alreadyExists = await _data.CommentExistsByExternalIdAsync(replyId, cancellationToken); if (alreadyExists) { skipped++; synced++; continue; } var workOrder = await _data.GetWorkOrderByInternalNumberAsync(woNumber, cancellationToken); if (workOrder == null) { skipped++; continue; } int? dispatchId = null; if (!string.IsNullOrWhiteSpace(dispatchNumber)) { var dispatch = await _data.GetDispatchByNumberAndWorkOrderIdAsync(dispatchNumber, workOrder.Id, cancellationToken); dispatchId = dispatch?.Id; } var comment = new Comments { ExternalCommentId = replyId, WorkerOrderId = workOrder.Id, DispatchId = dispatchId, Commenttext = Get(item, "reply_text"), Commenter = Get(item, "sender_email"), CommentType = "vendor", RecordType = "vendor_reply", CreatedDate = ParseDate(Get(item, "received_at")) ?? DateTime.UtcNow, }; _data.EnqueueComment(comment); created++; synced++; if (synced % 50 == 0) await _data.SaveChangesAsync(cancellationToken); } await _data.SaveChangesAsync(cancellationToken); return new SyncCommentsResult { Synced = synced, Created = created, Skipped = skipped }; } public async Task BackfillCommentTypesAsync(CancellationToken cancellationToken) { var updated = await _data.ExecuteBackfillCommentTypesAsync(cancellationToken); return new BackfillResult { Updated = updated }; } public async Task SyncAllAsync(CancellationToken cancellationToken) { var woResult = await SyncWorkOrdersAsync(cancellationToken); var commentResult = await SyncCommentsAsync(cancellationToken); return new SyncAllResult { WorkOrders = woResult, Comments = commentResult }; } private static string? Get(Dictionary item, string key) { return item.TryGetValue(key, out var val) ? val : null; } private static DateTime? ParseDate(string? value) { if (string.IsNullOrWhiteSpace(value)) return null; if (DateTime.TryParse(value, out var dt)) return dt; return null; } private static string? MapStatus(string? dynamoStatus) { return dynamoStatus switch { "new" => "Open", "assigned" => "Open", "in_progress" => "In Progress", "on_hold" => "On Hold", "completed" => "Done", "cancelled" => "Cancelled", "unknown" => "Open", _ => "Open" }; } private static string? MapSeverityToPriority(string? severity) { if (string.IsNullOrWhiteSpace(severity)) return null; return $"Sev {severity}"; } } }