shoc-backend/SeaHaven.Services/Implementation/SyncService.cs
Arthur Bassi 1edcf479ae fix(work-orders): apply account scope across create and reads [SH-221]
Stamp WorkOrder.AccountId on all create paths and filter board/list/search/detail by server-derived account claims so scoped callers cannot cross accounts.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-08-11 10:45:37 -03:00

315 lines
13 KiB
C#

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<SyncWorkOrdersResult> 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<SyncCommentsResult> 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<BackfillResult> 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<BackfillResult> 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<SyncCommentsResult> 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<BackfillResult> BackfillCommentTypesAsync(CancellationToken cancellationToken)
{
var updated = await _data.ExecuteBackfillCommentTypesAsync(cancellationToken);
return new BackfillResult { Updated = updated };
}
public async Task<SyncAllResult> 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<string, string?> 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}";
}
}
}