diff --git a/Api.SeaHavenIndustries/Controllers/WorkOrderIngestController.cs b/Api.SeaHavenIndustries/Controllers/WorkOrderIngestController.cs new file mode 100644 index 0000000..47b91a1 --- /dev/null +++ b/Api.SeaHavenIndustries/Controllers/WorkOrderIngestController.cs @@ -0,0 +1,50 @@ +using Api.SeaHavenIndustries.Filters; +using Microsoft.AspNetCore.Mvc; +using SeaHaven.Services.DTOs; +using SeaHaven.Services.Interfaces; + +namespace Api.SeaHavenIndustries.Controllers +{ + [ApiController] + [Route("api/workorders")] + public class WorkOrderIngestController : ControllerBase + { + private const int MaxBatchSize = 100; + private readonly IWorkOrderIngestService _ingestService; + + public WorkOrderIngestController(IWorkOrderIngestService ingestService) + { + _ingestService = ingestService; + } + + /// + /// Idempotent upsert by ExternalWorkOrderId (Lambda cutover target). + /// Auth: X-Ingest-Key header. + /// + [HttpPost("ingest")] + [IngestApiKey] + [ProducesResponseType(typeof(WorkOrderIngestBatchResultDto), StatusCodes.Status200OK)] + public async Task> Ingest( + [FromBody] WorkOrderIngestBatchRequestDto? batch, + CancellationToken cancellationToken) + { + var items = NormalizeItems(batch); + if (items.Count == 0) + return BadRequest("At least one item with externalWorkOrderId is required."); + + if (items.Count > MaxBatchSize) + return BadRequest($"Max {MaxBatchSize} items per request."); + + var result = await _ingestService.UpsertBatchAsync(items, cancellationToken); + return Ok(result); + } + + private static List NormalizeItems(WorkOrderIngestBatchRequestDto? batch) + { + if (batch?.Items is { Count: > 0 }) + return batch.Items; + + return new List(); + } + } +} diff --git a/Api.SeaHavenIndustries/Filters/IngestApiKeyFilter.cs b/Api.SeaHavenIndustries/Filters/IngestApiKeyFilter.cs new file mode 100644 index 0000000..b98bfcc --- /dev/null +++ b/Api.SeaHavenIndustries/Filters/IngestApiKeyFilter.cs @@ -0,0 +1,67 @@ +using System.Security.Cryptography; +using System.Text; +using Api.SeaHavenIndustries.Options; +using Microsoft.AspNetCore.Mvc; +using Microsoft.AspNetCore.Mvc.Filters; +using Microsoft.Extensions.Options; + +namespace Api.SeaHavenIndustries.Filters +{ + /// Validates X-Ingest-Key for POST /api/workorders/ingest. + public class IngestApiKeyFilter : IAsyncActionFilter + { + private readonly WorkOrderIngestOptions _options; + + public IngestApiKeyFilter(IOptions options) + { + _options = options.Value; + } + + public async Task OnActionExecutionAsync(ActionExecutingContext context, ActionExecutionDelegate next) + { + if (!_options.Enabled) + { + context.Result = new ObjectResult(new { message = "Work order ingest is disabled." }) + { + StatusCode = StatusCodes.Status503ServiceUnavailable + }; + return; + } + + if (string.IsNullOrWhiteSpace(_options.ApiKey) + || _options.ApiKey.Contains("${", StringComparison.Ordinal)) + { + context.Result = new ObjectResult(new { message = "Ingest API key is not configured." }) + { + StatusCode = StatusCodes.Status503ServiceUnavailable + }; + return; + } + + if (!context.HttpContext.Request.Headers.TryGetValue("X-Ingest-Key", out var provided) + || !FixedTimeEquals(provided.ToString(), _options.ApiKey)) + { + context.Result = new UnauthorizedObjectResult(new { message = "Invalid or missing X-Ingest-Key." }); + return; + } + + await next(); + } + + private static bool FixedTimeEquals(string provided, string expected) + { + var providedBytes = Encoding.UTF8.GetBytes(provided); + var expectedBytes = Encoding.UTF8.GetBytes(expected); + return providedBytes.Length == expectedBytes.Length + && CryptographicOperations.FixedTimeEquals(providedBytes, expectedBytes); + } + } + + [AttributeUsage(AttributeTargets.Class | AttributeTargets.Method)] + public sealed class IngestApiKeyAttribute : ServiceFilterAttribute + { + public IngestApiKeyAttribute() : base(typeof(IngestApiKeyFilter)) + { + } + } +} diff --git a/Api.SeaHavenIndustries/Options/WorkOrderIngestOptions.cs b/Api.SeaHavenIndustries/Options/WorkOrderIngestOptions.cs new file mode 100644 index 0000000..96bed87 --- /dev/null +++ b/Api.SeaHavenIndustries/Options/WorkOrderIngestOptions.cs @@ -0,0 +1,13 @@ +namespace Api.SeaHavenIndustries.Options +{ + public class WorkOrderIngestOptions + { + public const string SectionName = "WorkOrderIngest"; + + /// When false, ingest endpoints return 503. Default false until a real ApiKey is provisioned. + public bool Enabled { get; set; } + + /// Shared secret for X-Ingest-Key header (Lambda / service accounts). + public string? ApiKey { get; set; } + } +} diff --git a/SeaHaven.Services/DTOs/WorkOrderPhase7DTOs.cs b/SeaHaven.Services/DTOs/WorkOrderPhase7DTOs.cs new file mode 100644 index 0000000..4f1087a --- /dev/null +++ b/SeaHaven.Services/DTOs/WorkOrderPhase7DTOs.cs @@ -0,0 +1,54 @@ +namespace SeaHaven.Services.DTOs +{ + public sealed class WorkOrderIngestPayloadDto + { + public required string ExternalWorkOrderId { get; set; } + public string? Description { get; set; } + public string? WoStatus { get; set; } + public string? Severity { get; set; } + public string? Customer { get; set; } + public string? SiteCode { get; set; } + public string? Building { get; set; } + public string? Address { get; set; } + public DateTime? DueDate { get; set; } + public DateTime? DateReported { get; set; } + public DateTime? ScheduledStart { get; set; } + public string? SourceEmailS3Key { get; set; } + public DateTime? CreatedAt { get; set; } + } + + public sealed class WorkOrderIngestBatchRequestDto + { + public List? Items { get; set; } + } + + public sealed class WorkOrderIngestItemResultDto + { + public required string ExternalWorkOrderId { get; set; } + public int WorkOrderId { get; set; } + public bool Created { get; set; } + } + + public sealed class WorkOrderIngestBatchResultDto + { + public int Created { get; set; } + public int Updated { get; set; } + public int Skipped { get; set; } + public List Results { get; set; } = new(); + } + + public sealed class WorkOrderOpsHealthDto + { + public DateTime? LastWeekRolledRunUtc { get; set; } + public DateTime? LastPastDueCacheRunUtc { get; set; } + public string? LastWeekRolledError { get; set; } + public string? LastPastDueCacheError { get; set; } + public int SyncRejectedLast24h { get; set; } + public int FieldLockCount { get; set; } + public bool SyncEnabled { get; set; } + public bool IngestEnabled { get; set; } + public bool LegacyDeprecationEnabled { get; set; } + public string? LegacySunsetDate { get; set; } + public DateTime CheckedAtUtc { get; set; } + } +} diff --git a/SeaHaven.Services/Helpers/WorkOrderIngestFieldMapper.cs b/SeaHaven.Services/Helpers/WorkOrderIngestFieldMapper.cs new file mode 100644 index 0000000..cc48e5e --- /dev/null +++ b/SeaHaven.Services/Helpers/WorkOrderIngestFieldMapper.cs @@ -0,0 +1,27 @@ +namespace SeaHaven.Services.Helpers +{ + /// DynamoDB / ingest field mapping shared by Sync and direct ingest. + public static class WorkOrderIngestFieldMapper + { + public 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", + _ => string.IsNullOrWhiteSpace(dynamoStatus) ? null : "Open" + }; + } + + public static string? MapSeverityToPriority(string? severity) + { + if (string.IsNullOrWhiteSpace(severity)) return null; + return $"Sev {severity}"; + } + } +} diff --git a/SeaHaven.Services/Implementation/WorkOrderIngestService.cs b/SeaHaven.Services/Implementation/WorkOrderIngestService.cs new file mode 100644 index 0000000..aa0638e --- /dev/null +++ b/SeaHaven.Services/Implementation/WorkOrderIngestService.cs @@ -0,0 +1,216 @@ +using Data.SeaHavenIndustries; +using Microsoft.EntityFrameworkCore; +using SeaHaven.Services.DTOs; +using SeaHaven.Services.Helpers; +using SeaHaven.Services.Interfaces; + +namespace SeaHaven.Services.Implementation +{ + public class WorkOrderIngestService : IWorkOrderIngestService + { + private readonly ApplicationDbContext _db; + private readonly ISyncFieldMergePolicy _mergePolicy; + private readonly IWorkOrderFieldLockService _fieldLocks; + private readonly IWorkOrderAuditService _audit; + + public WorkOrderIngestService( + ApplicationDbContext db, + ISyncFieldMergePolicy mergePolicy, + IWorkOrderFieldLockService fieldLocks, + IWorkOrderAuditService audit) + { + _db = db; + _mergePolicy = mergePolicy; + _fieldLocks = fieldLocks; + _audit = audit; + } + + public async Task UpsertBatchAsync( + IReadOnlyList items, + CancellationToken cancellationToken = default) + { + var result = new WorkOrderIngestBatchResultDto(); + if (items.Count == 0) + return result; + + var useTransaction = _db.Database.IsRelational(); + await using var transaction = useTransaction + ? await _db.Database.BeginTransactionAsync(cancellationToken) + : null; + + try + { + var nextSeed = await AllocateNextWoSeedAsync(cancellationToken); + + foreach (var item in items) + { + if (string.IsNullOrWhiteSpace(item.ExternalWorkOrderId)) + { + result.Skipped++; + continue; + } + + var (upsert, updatedSeed) = await UpsertOneAsync(item, nextSeed, cancellationToken); + nextSeed = updatedSeed; + result.Results.Add(upsert); + if (upsert.Created) + result.Created++; + else + result.Updated++; + } + + await _db.SaveChangesAsync(cancellationToken); + if (transaction != null) + await transaction.CommitAsync(cancellationToken); + } + catch + { + if (transaction != null) + await transaction.RollbackAsync(cancellationToken); + throw; + } + + return result; + } + + private async Task<(WorkOrderIngestItemResultDto Result, int NextSeed)> UpsertOneAsync( + WorkOrderIngestPayloadDto item, + int nextSeed, + CancellationToken cancellationToken) + { + var existing = await _db.workOrders + .FirstOrDefaultAsync(w => w.ExternalWorkOrderId == item.ExternalWorkOrderId, cancellationToken); + + var locationId = await ResolveLocationIdAsync( + item.SiteCode, item.Building, item.Address, cancellationToken); + + if (existing == null) + { + 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, + 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 + }; + _db.workOrders.Add(wo); + await _db.SaveChangesAsync(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 _db.workOrders.MaxAsync(w => (int?)w.Id, cancellationToken) ?? 0; + 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 _db.Locations + .FirstOrDefaultAsync(l => l.Name == matchCode || l.Title == matchCode, cancellationToken); + + if (existing != null) + return existing.Id; + + var location = new Locations + { + Name = matchCode, + Title = building, + Address1 = address, + Status = "Active" + }; + + _db.Locations.Add(location); + await _db.SaveChangesAsync(cancellationToken); + return location.Id; + } + } +} diff --git a/SeaHaven.Services/Interfaces/IWorkOrderPhase7Services.cs b/SeaHaven.Services/Interfaces/IWorkOrderPhase7Services.cs new file mode 100644 index 0000000..0242a38 --- /dev/null +++ b/SeaHaven.Services/Interfaces/IWorkOrderPhase7Services.cs @@ -0,0 +1,16 @@ +using SeaHaven.Services.DTOs; + +namespace SeaHaven.Services.Interfaces +{ + public interface IWorkOrderIngestService + { + Task UpsertBatchAsync( + IReadOnlyList items, + CancellationToken cancellationToken = default); + } + + public interface IWorkOrderOpsHealthService + { + Task GetHealthAsync(CancellationToken cancellationToken = default); + } +}