mirror of
https://github.com/Sea-Haven-Industries/shoc-backend.git
synced 2026-09-30 11:53:12 +00:00
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>
341 lines
14 KiB
C#
341 lines
14 KiB
C#
using Data.SeaHavenIndustries;
|
|
using Microsoft.EntityFrameworkCore;
|
|
using Microsoft.EntityFrameworkCore.Storage;
|
|
using SeaHaven.DataServices.Interfaces;
|
|
|
|
namespace SeaHaven.DataServices.Implementation
|
|
{
|
|
public sealed class WorkOrderWebhookDataService : IWorkOrderWebhookDataService
|
|
{
|
|
private static readonly SemaphoreSlim InMemoryNumberLock = new(1, 1);
|
|
private readonly ApplicationDbContext _context;
|
|
|
|
public WorkOrderWebhookDataService(ApplicationDbContext context)
|
|
{
|
|
_context = context;
|
|
}
|
|
|
|
public async Task<WorkOrderWebhookPersistenceResult> ApplyAsync(
|
|
WorkOrderWebhookMutation mutation,
|
|
CancellationToken cancellationToken)
|
|
{
|
|
var prior = await _context.WorkOrderWebhookDeliveries
|
|
.AsNoTracking()
|
|
.SingleOrDefaultAsync(d => d.DeliveryId == mutation.DeliveryId, cancellationToken);
|
|
|
|
if (prior != null)
|
|
return ExistingDeliveryResult(prior, mutation.BodySha256);
|
|
|
|
var workOrder = await _context.workOrders
|
|
.SingleOrDefaultAsync(
|
|
w => w.ExternalWorkOrderId == mutation.ExternalWorkOrderId,
|
|
cancellationToken);
|
|
|
|
if (workOrder == null)
|
|
{
|
|
var accountId = await TryResolveAccountIdAsync(mutation.Customer, cancellationToken);
|
|
if (accountId is not int resolvedAccountId)
|
|
{
|
|
return new WorkOrderWebhookPersistenceResult(
|
|
WorkOrderWebhookPersistenceStatus.AccountUnresolved);
|
|
}
|
|
|
|
var internalNumber = await AllocateInternalNumberAsync(cancellationToken);
|
|
workOrder = new WorkOrder
|
|
{
|
|
ExternalWorkOrderId = mutation.ExternalWorkOrderId,
|
|
ExternalSource = mutation.Source,
|
|
WorkerOrderNumber = mutation.WorkerOrderNumber,
|
|
InternalWONumber = internalNumber.ToString("D11"),
|
|
WorkerOrderTitle = mutation.Title ?? mutation.Description
|
|
?? $"Imported work order {mutation.ExternalWorkOrderId}",
|
|
Description = mutation.Description,
|
|
Customer = mutation.Customer,
|
|
AccountId = resolvedAccountId,
|
|
Source = mutation.Source,
|
|
istemplate = false,
|
|
CreatedDate = mutation.CreatedAt ?? mutation.ProcessedAt.UtcDateTime
|
|
};
|
|
_context.workOrders.Add(workOrder);
|
|
}
|
|
|
|
var staleState = mutation.IsStateEvent
|
|
&& !IsNewer(
|
|
mutation.UpdatedAt,
|
|
mutation.VersionHash,
|
|
workOrder.ExternalUpdatedAt,
|
|
workOrder.ExternalVersionHash);
|
|
|
|
if (mutation.IsStateEvent && !staleState)
|
|
{
|
|
var locationKey = mutation.SiteCode ?? mutation.Building;
|
|
if (!string.IsNullOrWhiteSpace(locationKey))
|
|
{
|
|
var locationIsNewer = await UpsertReceiptAsync(
|
|
mutation.Source,
|
|
"location",
|
|
locationKey,
|
|
mutation.UpdatedAt,
|
|
mutation.VersionHash,
|
|
mutation.ProcessedAt,
|
|
cancellationToken);
|
|
var location = await _context.Locations
|
|
.FirstOrDefaultAsync(
|
|
l => l.ExternalSource == mutation.Source
|
|
&& l.ExternalLocationId == locationKey,
|
|
cancellationToken);
|
|
if (location == null)
|
|
{
|
|
location = new Locations
|
|
{
|
|
ExternalSource = mutation.Source,
|
|
ExternalLocationId = locationKey,
|
|
Name = locationKey,
|
|
Title = mutation.Building,
|
|
Address1 = mutation.Address,
|
|
Status = "Active"
|
|
};
|
|
_context.Locations.Add(location);
|
|
}
|
|
else if (locationIsNewer)
|
|
{
|
|
location.Name = locationKey;
|
|
location.Title = mutation.Building;
|
|
location.Address1 = mutation.Address;
|
|
location.Status = "Active";
|
|
}
|
|
|
|
workOrder.Locations = location;
|
|
}
|
|
|
|
workOrder.WorkerOrderNumber = mutation.WorkerOrderNumber;
|
|
workOrder.WorkerOrderTitle = mutation.Title ?? mutation.Description;
|
|
workOrder.Description = mutation.Description;
|
|
workOrder.Status = mutation.IsCancelled ? "Cancelled" : mutation.Status;
|
|
if (mutation.LifecycleStatus.HasValue)
|
|
workOrder.LifecycleStatus = mutation.LifecycleStatus.Value;
|
|
workOrder.Severity = mutation.Severity;
|
|
workOrder.Priority = mutation.Priority;
|
|
workOrder.ExternalAssignedTo = mutation.AssignedTo;
|
|
workOrder.ExternalRecordType = mutation.RecordType;
|
|
workOrder.Customer = mutation.Customer;
|
|
workOrder.SiteCode = mutation.SiteCode;
|
|
workOrder.Building = mutation.Building;
|
|
workOrder.DueDate = mutation.DueDate;
|
|
workOrder.DateReported = mutation.DateReported;
|
|
workOrder.ScheduledStart = mutation.ScheduledStart;
|
|
workOrder.Source = mutation.Source;
|
|
workOrder.SourceEmailS3Key = mutation.SourceEmailS3Key;
|
|
workOrder.ExternalSource = mutation.Source;
|
|
workOrder.istemplate = false;
|
|
workOrder.ExternalLastOccurredAt = mutation.OccurredAt;
|
|
workOrder.ExternalUpdatedAt = mutation.UpdatedAt;
|
|
workOrder.ExternalVersionHash = mutation.VersionHash;
|
|
}
|
|
|
|
if (mutation.CommentId != null)
|
|
{
|
|
var receiptIsNewer = await UpsertReceiptAsync(
|
|
mutation.Source,
|
|
"comment",
|
|
mutation.CommentId,
|
|
mutation.UpdatedAt,
|
|
mutation.VersionHash,
|
|
mutation.ProcessedAt,
|
|
cancellationToken);
|
|
var comment = await _context.Comments
|
|
.SingleOrDefaultAsync(
|
|
c => c.ExternalSource == mutation.Source
|
|
&& c.ExternalCommentId == mutation.CommentId,
|
|
cancellationToken);
|
|
|
|
if (comment == null)
|
|
{
|
|
comment = new Comments
|
|
{
|
|
ExternalSource = mutation.Source,
|
|
ExternalCommentId = mutation.CommentId,
|
|
RecordType = "WorkOrder",
|
|
WorkOrder = workOrder,
|
|
CreatedDate = mutation.OccurredAt.UtcDateTime,
|
|
};
|
|
_context.Comments.Add(comment);
|
|
}
|
|
|
|
if (receiptIsNewer)
|
|
{
|
|
comment.Commenttext = mutation.CommentText;
|
|
comment.Commenter = mutation.Commenter;
|
|
comment.CommentType = mutation.CommentType;
|
|
comment.ExternalUpdatedAt = mutation.UpdatedAt;
|
|
comment.ExternalVersionHash = mutation.VersionHash;
|
|
comment.ExternalSourceEmailS3Key = mutation.SourceEmailS3Key;
|
|
}
|
|
}
|
|
|
|
_context.WorkOrderWebhookDeliveries.Add(new WorkOrderWebhookDelivery
|
|
{
|
|
DeliveryId = mutation.DeliveryId,
|
|
EventType = mutation.EventType,
|
|
OccurredAt = mutation.OccurredAt,
|
|
ProcessedAt = mutation.ProcessedAt,
|
|
BodySha256 = mutation.BodySha256
|
|
});
|
|
|
|
try
|
|
{
|
|
await _context.SaveChangesAsync(cancellationToken);
|
|
return new WorkOrderWebhookPersistenceResult(
|
|
WorkOrderWebhookPersistenceStatus.Applied,
|
|
staleState);
|
|
}
|
|
catch (DbUpdateException)
|
|
{
|
|
_context.ChangeTracker.Clear();
|
|
var concurrent = await _context.WorkOrderWebhookDeliveries
|
|
.AsNoTracking()
|
|
.SingleOrDefaultAsync(d => d.DeliveryId == mutation.DeliveryId, cancellationToken);
|
|
|
|
if (concurrent == null)
|
|
throw;
|
|
|
|
return ExistingDeliveryResult(concurrent, mutation.BodySha256);
|
|
}
|
|
}
|
|
|
|
private async Task<long> AllocateInternalNumberAsync(CancellationToken cancellationToken)
|
|
{
|
|
if (_context.Database.IsSqlServer())
|
|
{
|
|
return await AllocateSqlServerInternalNumberAsync(cancellationToken);
|
|
}
|
|
|
|
await InMemoryNumberLock.WaitAsync(cancellationToken);
|
|
try
|
|
{
|
|
var values = await _context.workOrders
|
|
.AsNoTracking()
|
|
.Where(w => w.InternalWONumber != null)
|
|
.Select(w => w.InternalWONumber!)
|
|
.ToListAsync(cancellationToken);
|
|
return values
|
|
.Select(v => long.TryParse(v, out var parsed) ? parsed : 0L)
|
|
.DefaultIfEmpty()
|
|
.Max() + 1;
|
|
}
|
|
finally
|
|
{
|
|
InMemoryNumberLock.Release();
|
|
}
|
|
}
|
|
|
|
private async Task<long> AllocateSqlServerInternalNumberAsync(CancellationToken cancellationToken)
|
|
{
|
|
var connection = _context.Database.GetDbConnection();
|
|
var openedHere = false;
|
|
if (connection.State != System.Data.ConnectionState.Open)
|
|
{
|
|
await _context.Database.OpenConnectionAsync(cancellationToken);
|
|
openedHere = true;
|
|
}
|
|
|
|
try
|
|
{
|
|
using var command = connection.CreateCommand();
|
|
command.CommandText = "SELECT NEXT VALUE FOR dbo.WorkOrderInternalNumberSequence";
|
|
var currentTransaction = _context.Database.CurrentTransaction;
|
|
if (currentTransaction != null)
|
|
{
|
|
command.Transaction = currentTransaction.GetDbTransaction();
|
|
}
|
|
var result = await command.ExecuteScalarAsync(cancellationToken);
|
|
return Convert.ToInt64(result, System.Globalization.CultureInfo.InvariantCulture);
|
|
}
|
|
finally
|
|
{
|
|
if (openedHere)
|
|
await _context.Database.CloseConnectionAsync();
|
|
}
|
|
}
|
|
|
|
private async Task<bool> UpsertReceiptAsync(
|
|
string source,
|
|
string kind,
|
|
string externalId,
|
|
DateTimeOffset updatedAt,
|
|
string versionHash,
|
|
DateTimeOffset processedAt,
|
|
CancellationToken cancellationToken)
|
|
{
|
|
var receipt = await _context.WorkOrderExternalReceipts.SingleOrDefaultAsync(
|
|
r => r.Source == source && r.Kind == kind && r.ExternalId == externalId,
|
|
cancellationToken);
|
|
if (receipt == null)
|
|
{
|
|
_context.WorkOrderExternalReceipts.Add(new WorkOrderExternalReceipt
|
|
{
|
|
Source = source,
|
|
Kind = kind,
|
|
ExternalId = externalId,
|
|
UpdatedAt = updatedAt,
|
|
VersionHash = versionHash,
|
|
ProcessedAt = processedAt
|
|
});
|
|
return true;
|
|
}
|
|
|
|
if (!IsNewer(updatedAt, versionHash, receipt.UpdatedAt, receipt.VersionHash))
|
|
return false;
|
|
|
|
receipt.UpdatedAt = updatedAt;
|
|
receipt.VersionHash = versionHash;
|
|
receipt.ProcessedAt = processedAt;
|
|
return true;
|
|
}
|
|
|
|
private async Task<int?> TryResolveAccountIdAsync(
|
|
string? customer,
|
|
CancellationToken cancellationToken)
|
|
{
|
|
if (string.IsNullOrWhiteSpace(customer))
|
|
return null;
|
|
|
|
var trimmed = customer.Trim();
|
|
var matches = await _context.Accounts
|
|
.AsNoTracking()
|
|
.Where(a =>
|
|
a.Name == trimmed
|
|
&& (a.IsDeleted == false || a.IsDeleted == null))
|
|
.Select(a => a.Id)
|
|
.Take(2)
|
|
.ToListAsync(cancellationToken);
|
|
|
|
return matches.Count == 1 ? matches[0] : null;
|
|
}
|
|
|
|
private static bool IsNewer(
|
|
DateTimeOffset candidateUpdatedAt,
|
|
string candidateHash,
|
|
DateTimeOffset? currentUpdatedAt,
|
|
string? currentHash)
|
|
{
|
|
if (!currentUpdatedAt.HasValue)
|
|
return true;
|
|
|
|
var timestampComparison = candidateUpdatedAt.CompareTo(currentUpdatedAt.Value);
|
|
return timestampComparison > 0
|
|
|| (timestampComparison == 0
|
|
&& string.CompareOrdinal(candidateHash, currentHash ?? string.Empty) > 0);
|
|
}
|
|
|
|
private static WorkOrderWebhookPersistenceResult ExistingDeliveryResult(
|
|
WorkOrderWebhookDelivery delivery,
|
|
string bodySha256)
|
|
{
|
|
return new WorkOrderWebhookPersistenceResult(
|
|
string.Equals(delivery.BodySha256, bodySha256, StringComparison.OrdinalIgnoreCase)
|
|
? WorkOrderWebhookPersistenceStatus.Duplicate
|
|
: WorkOrderWebhookPersistenceStatus.HashConflict);
|
|
}
|
|
}
|
|
}
|