using Amazon.DynamoDBv2; using Amazon.DynamoDBv2.Model; using Api.SeaHavenIndustries.Options; using Data.SeaHavenIndustries; using Microsoft.AspNetCore.Authorization; using Microsoft.AspNetCore.Mvc; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; namespace Api.SeaHavenIndustries.Controllers { [Authorize(Roles = "Admin")] [ApiController] [Route("api/[controller]")] public class SyncController : Controller { private readonly ApplicationDbContext _db; private readonly AmazonDynamoDBClient _dynamo; private readonly SyncOptions _syncOptions; private const string WO_TABLE = "WorkOrders"; private const string COMMENTS_TABLE = "WorkOrderComments"; private const string VENDOR_REPLIES_TABLE = "VendorReplies"; public SyncController(ApplicationDbContext db, IOptions syncOptions) { _db = db; _syncOptions = syncOptions.Value; _dynamo = new AmazonDynamoDBClient(Amazon.RegionEndpoint.USEast1); } private IActionResult? SyncDisabledResult() { if (_syncOptions.Enabled) return null; return StatusCode(StatusCodes.Status503ServiceUnavailable, new { message = "DynamoDB sync bridge is disabled. Use POST /api/workorders/ingest." }); } [HttpPost("WorkOrders")] public async Task SyncWorkOrders() { if (SyncDisabledResult() is { } disabled) return disabled; var synced = 0; var created = 0; var updated = 0; var locationsCreated = 0; var lastInternalWO = await _db.workOrders .Where(w => w.InternalWONumber != null && w.InternalWONumber != "") .OrderByDescending(w => w.InternalWONumber) .Select(w => w.InternalWONumber) .FirstOrDefaultAsync(); int nextInternal = 10000001; if (lastInternalWO != null && int.TryParse(lastInternalWO, out var parsed)) nextInternal = parsed + 1; Dictionary? lastKey = null; do { var request = new ScanRequest { TableName = WO_TABLE, Limit = 100, ExclusiveStartKey = lastKey }; var response = await _dynamo.ScanAsync(request); foreach (var item in response.Items) { var externalId = GetString(item, "work_order_id"); if (string.IsNullOrEmpty(externalId)) continue; var existing = await _db.workOrders .FirstOrDefaultAsync(w => w.ExternalWorkOrderId == externalId); var (locationId, locationWasCreated) = await ResolveLocationId( GetString(item, "site_code"), GetString(item, "building"), GetString(item, "address")); if (locationWasCreated) locationsCreated++; if (existing == null) { var wo = new WorkOrder { InternalWONumber = (nextInternal++).ToString("D8"), ExternalWorkOrderId = externalId, WorkerOrderNumber = externalId, WorkerOrderTitle = GetString(item, "description"), Description = GetString(item, "description"), Status = MapStatus(GetString(item, "wo_status")), Priority = MapSeverityToPriority(GetString(item, "severity")), Severity = GetString(item, "severity"), Customer = GetString(item, "customer"), SiteCode = GetString(item, "site_code"), Building = GetString(item, "building"), LocationId = locationId, DueDate = ParseDate(GetString(item, "due_date")), DateReported = ParseDate(GetString(item, "date_reported")), ScheduledStart = ParseDate(GetString(item, "scheduled_start")), AssignTo = null, SourceEmailS3Key = GetString(item, "source_email_s3_key"), CreatedDate = ParseDate(GetString(item, "created_at")) ?? DateTime.UtcNow, istemplate = false }; _db.workOrders.Add(wo); created++; } else { existing.WorkerOrderTitle = GetString(item, "description") ?? existing.WorkerOrderTitle; existing.Description = GetString(item, "description") ?? existing.Description; existing.Status = MapStatus(GetString(item, "wo_status")) ?? existing.Status; existing.Priority = MapSeverityToPriority(GetString(item, "severity")) ?? existing.Priority; existing.Severity = GetString(item, "severity") ?? existing.Severity; existing.SiteCode = GetString(item, "site_code") ?? existing.SiteCode; existing.Building = GetString(item, "building") ?? existing.Building; existing.LocationId = locationId ?? existing.LocationId; existing.DueDate = ParseDate(GetString(item, "due_date")) ?? existing.DueDate; existing.DateReported = ParseDate(GetString(item, "date_reported")) ?? existing.DateReported; existing.ScheduledStart = ParseDate(GetString(item, "scheduled_start")) ?? existing.ScheduledStart; existing.SourceEmailS3Key = GetString(item, "source_email_s3_key") ?? existing.SourceEmailS3Key; updated++; } synced++; if (synced % 100 == 0) await _db.SaveChangesAsync(); } lastKey = response.LastEvaluatedKey; } while (lastKey != null && lastKey.Count > 0); await _db.SaveChangesAsync(); return Ok(new { synced, created, updated, locationsCreated = await _db.Locations.CountAsync() }); } [HttpPost("Comments")] public async Task SyncComments() { if (SyncDisabledResult() is { } disabled) return disabled; var synced = 0; var created = 0; var skipped = 0; Dictionary? lastKey = null; do { var request = new ScanRequest { TableName = COMMENTS_TABLE, Limit = 100, ExclusiveStartKey = lastKey }; var response = await _dynamo.ScanAsync(request); foreach (var item in response.Items) { var externalCommentId = GetString(item, "comment_id"); var externalWoId = GetString(item, "work_order_id"); if (string.IsNullOrEmpty(externalCommentId) || string.IsNullOrEmpty(externalWoId)) continue; var alreadyExists = await _db.Comments .AnyAsync(c => c.ExternalCommentId == externalCommentId); if (alreadyExists) { skipped++; synced++; continue; } var workOrder = await _db.workOrders .FirstOrDefaultAsync(w => w.ExternalWorkOrderId == externalWoId); if (workOrder == null) { skipped++; continue; } var comment = new Comments { ExternalCommentId = externalCommentId, WorkerOrderId = workOrder.Id, Commenttext = GetString(item, "text"), Commenter = GetString(item, "commenter"), RecordType = GetString(item, "record_type"), CommentType = "customer", CreatedDate = ParseDate(GetString(item, "created_at")) ?? DateTime.UtcNow, }; _db.Comments.Add(comment); created++; synced++; if (synced % 100 == 0) await _db.SaveChangesAsync(); } lastKey = response.LastEvaluatedKey; } while (lastKey != null && lastKey.Count > 0); await _db.SaveChangesAsync(); return Ok(new { synced, created, skipped }); } [HttpPost("BackfillInternalWONumbers")] public async Task BackfillInternalWONumbers() { if (SyncDisabledResult() is { } disabled) return disabled; var wosMissing = await _db.workOrders .Where(w => w.InternalWONumber == null || w.InternalWONumber == "") .OrderBy(w => w.Id) .ToListAsync(); if (wosMissing.Count == 0) return Ok(new { updated = 0 }); var lastWO = await _db.workOrders .Where(w => w.InternalWONumber != null && w.InternalWONumber != "") .OrderByDescending(w => w.InternalWONumber) .Select(w => w.InternalWONumber) .FirstOrDefaultAsync(); 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 _db.SaveChangesAsync(); return Ok(new { updated = wosMissing.Count }); } [HttpPost("BackfillDispatchNumbers")] public async Task BackfillDispatchNumbers() { if (SyncDisabledResult() is { } disabled) return disabled; var missing = await _db.Dispatches .Where(d => d.DispatchNumber == null || d.DispatchNumber == "") .OrderBy(d => d.Id) .ToListAsync(); if (missing.Count == 0) return Ok(new { updated = 0 }); int next = 1; var last = await _db.Dispatches .Where(d => d.DispatchNumber != null && d.DispatchNumber != "") .OrderByDescending(d => d.DispatchNumber) .Select(d => d.DispatchNumber) .FirstOrDefaultAsync(); 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 _db.SaveChangesAsync(); return Ok(new { updated = missing.Count }); } [HttpPost("VendorReplies")] public async Task SyncVendorReplies() { if (SyncDisabledResult() is { } disabled) return disabled; var synced = 0; var created = 0; var skipped = 0; Dictionary? lastKey = null; do { var request = new ScanRequest { TableName = VENDOR_REPLIES_TABLE, Limit = 100, ExclusiveStartKey = lastKey }; var response = await _dynamo.ScanAsync(request); foreach (var item in response.Items) { var replyId = GetString(item, "reply_id"); var woNumber = GetString(item, "internal_wo_number"); var dispatchNumber = GetString(item, "dispatch_number"); if (string.IsNullOrEmpty(replyId) || string.IsNullOrEmpty(woNumber)) continue; var alreadyExists = await _db.Comments .AnyAsync(c => c.ExternalCommentId == replyId); if (alreadyExists) { skipped++; synced++; continue; } var workOrder = await _db.workOrders .FirstOrDefaultAsync(w => w.InternalWONumber == woNumber); if (workOrder == null) { skipped++; continue; } int? dispatchId = null; if (!string.IsNullOrWhiteSpace(dispatchNumber)) { var dispatch = await _db.Dispatches .FirstOrDefaultAsync(d => d.DispatchNumber == dispatchNumber && d.WorkOrderId == workOrder.Id); dispatchId = dispatch?.Id; } var comment = new Comments { ExternalCommentId = replyId, WorkerOrderId = workOrder.Id, DispatchId = dispatchId, Commenttext = GetString(item, "reply_text"), Commenter = GetString(item, "sender_email"), CommentType = "vendor", RecordType = "vendor_reply", CreatedDate = ParseDate(GetString(item, "received_at")) ?? DateTime.UtcNow, }; _db.Comments.Add(comment); created++; synced++; if (synced % 50 == 0) await _db.SaveChangesAsync(); } lastKey = response.LastEvaluatedKey; } while (lastKey != null && lastKey.Count > 0); await _db.SaveChangesAsync(); return Ok(new { synced, created, skipped }); } [HttpPost("BackfillCommentTypes")] public async Task BackfillCommentTypes() { if (SyncDisabledResult() is { } disabled) return disabled; var updated = await _db.Database.ExecuteSqlRawAsync( "UPDATE Comments SET CommentType = 'customer' WHERE ExternalCommentId IS NOT NULL AND (CommentType IS NULL OR CommentType = '')"); return Ok(new { updated }); } [HttpPost("All")] public async Task SyncAll() { if (SyncDisabledResult() is { } disabled) return disabled; var woResult = await SyncWorkOrders() as OkObjectResult; var commentResult = await SyncComments() as OkObjectResult; return Ok(new { workOrders = woResult?.Value, comments = commentResult?.Value }); } private async Task<(int? LocationId, bool Created)> ResolveLocationId(string? siteCode, string? building, string? address) { if (string.IsNullOrWhiteSpace(siteCode) && string.IsNullOrWhiteSpace(building)) return (null, false); var matchCode = siteCode ?? building; var existing = await _db.Locations .FirstOrDefaultAsync(l => l.Name == matchCode || l.Title == matchCode); if (existing != null) return (existing.Id, false); var location = new Locations { Name = matchCode, Title = building, Address1 = address, Status = "Active" }; _db.Locations.Add(location); await _db.SaveChangesAsync(); return (location.Id, true); } private static string? GetString(Dictionary item, string key) { if (item.TryGetValue(key, out var val) && val.S != null) return val.S; return 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}"; } } }