From 77dc3d85c336676347c993f0d15d57baf836d6af Mon Sep 17 00:00:00 2001 From: Alexandre Brandizzi Date: Fri, 25 Sep 2026 19:15:38 -0300 Subject: [PATCH] Serialize a CRM cancel with uplift creation on the work order lock The per-work-order gate moves into a shared data-layer helper. A CRM mutation that cancels an existing work order now runs under it, so a create in flight either commits first and is cancelled, or sees the cancelled work order. --- .../Helpers/WorkOrderMutationLock.cs | 79 +++++++++++++++++++ .../Implementation/UpliftDataService.cs | 47 +---------- .../WorkOrderWebhookDataService.cs | 26 ++++++ .../WorkOrderCrmCancelUpliftTests.cs | 37 +++++++++ 4 files changed, 144 insertions(+), 45 deletions(-) create mode 100644 SeaHaven.DataServices/Helpers/WorkOrderMutationLock.cs diff --git a/SeaHaven.DataServices/Helpers/WorkOrderMutationLock.cs b/SeaHaven.DataServices/Helpers/WorkOrderMutationLock.cs new file mode 100644 index 0000000..e15df5d --- /dev/null +++ b/SeaHaven.DataServices/Helpers/WorkOrderMutationLock.cs @@ -0,0 +1,79 @@ +using System.Collections.Concurrent; +using Data.SeaHavenIndustries; +using Microsoft.EntityFrameworkCore; + +namespace SeaHaven.DataServices.Helpers +{ + /// + /// The per-work-order gate that serializes work-order-scoped invariants (uplift create, + /// decide, revoke, and cancelling a work order's pending uplifts): an in-process gate per + /// work order, a transaction, and an update lock on the work order row on SQL Server. + /// Every caller shares the same gates, so a cancel and a create on one work order never + /// interleave. + /// + public static class WorkOrderMutationLock + { + private static readonly ConcurrentDictionary Gates = new(); + + /// + /// Runs under the work order's gate and row lock. The + /// transaction commits when the work returns and (when + /// given) accepts its result; otherwise, or when the work throws, it rolls back. + /// + public static async Task RunAsync( + ApplicationDbContext context, + int workOrderId, + Func> work, + CancellationToken cancellationToken, + Func? commitWhen = null) + { + var gate = Gates.GetOrAdd(workOrderId, _ => new SemaphoreSlim(1, 1)); + await gate.WaitAsync(cancellationToken); + try + { + await using var transaction = context.Database.IsRelational() + ? await context.Database.BeginTransactionAsync(cancellationToken) + : null; + try + { + await LockRowAsync(context, workOrderId, cancellationToken); + var result = await work(cancellationToken); + if (transaction is not null) + { + if (commitWhen is null || commitWhen(result)) + await transaction.CommitAsync(cancellationToken); + else + await transaction.RollbackAsync(cancellationToken); + } + return result; + } + catch + { + if (transaction is not null) + await transaction.RollbackAsync(cancellationToken); + throw; + } + } + finally + { + gate.Release(); + } + } + + private static async Task LockRowAsync( + ApplicationDbContext context, + int workOrderId, + CancellationToken cancellationToken) + { + if (context.Database.ProviderName?.Contains("SqlServer", StringComparison.OrdinalIgnoreCase) != true) + return; + + await context.workOrders + .FromSqlRaw( + "SELECT * FROM [workOrders] WITH (UPDLOCK, ROWLOCK, HOLDLOCK) WHERE [Id] = {0}", + workOrderId) + .Select(workOrder => workOrder.Id) + .FirstOrDefaultAsync(cancellationToken); + } + } +} diff --git a/SeaHaven.DataServices/Implementation/UpliftDataService.cs b/SeaHaven.DataServices/Implementation/UpliftDataService.cs index 3f9cf21..9f4e0c5 100644 --- a/SeaHaven.DataServices/Implementation/UpliftDataService.cs +++ b/SeaHaven.DataServices/Implementation/UpliftDataService.cs @@ -1,4 +1,3 @@ -using System.Collections.Concurrent; using Data.SeaHavenIndustries; using Data.SeaHavenIndustries.Enums; using Microsoft.EntityFrameworkCore; @@ -10,8 +9,6 @@ namespace SeaHaven.DataServices.Implementation { public class UpliftDataService : IUpliftDataService { - private static readonly ConcurrentDictionary WorkOrderGates = new(); - // The Rejected queue also surfaces the legacy "Denied" spelling, which reads as Rejected. private static readonly string[] RejectedStatuses = { "Rejected", "Denied" }; private readonly ApplicationDbContext _context; @@ -590,50 +587,10 @@ namespace SeaHaven.DataServices.Implementation return message.Contains(activeDispatchIndex, StringComparison.OrdinalIgnoreCase); } - public async Task ExecuteWorkOrderMutationAsync( + public Task ExecuteWorkOrderMutationAsync( int workOrderId, Func> work, CancellationToken cancellationToken) - { - var gate = WorkOrderGates.GetOrAdd(workOrderId, _ => new SemaphoreSlim(1, 1)); - await gate.WaitAsync(cancellationToken); - try - { - await using var transaction = _context.Database.IsRelational() - ? await _context.Database.BeginTransactionAsync(cancellationToken) - : null; - try - { - await LockWorkOrderRowAsync(workOrderId, cancellationToken); - var result = await work(cancellationToken); - if (transaction is not null) - await transaction.CommitAsync(cancellationToken); - return result; - } - catch - { - if (transaction is not null) - await transaction.RollbackAsync(cancellationToken); - throw; - } - } - finally - { - gate.Release(); - } - } - - private async Task LockWorkOrderRowAsync(int workOrderId, CancellationToken cancellationToken) - { - if (_context.Database.ProviderName?.Contains("SqlServer", StringComparison.OrdinalIgnoreCase) != true) - return; - - await _context.workOrders - .FromSqlRaw( - "SELECT * FROM [workOrders] WITH (UPDLOCK, ROWLOCK, HOLDLOCK) WHERE [Id] = {0}", - workOrderId) - .Select(workOrder => workOrder.Id) - .FirstOrDefaultAsync(cancellationToken); - } + => WorkOrderMutationLock.RunAsync(_context, workOrderId, work, cancellationToken); } } diff --git a/SeaHaven.DataServices/Implementation/WorkOrderWebhookDataService.cs b/SeaHaven.DataServices/Implementation/WorkOrderWebhookDataService.cs index 49cdc7b..bc0915a 100644 --- a/SeaHaven.DataServices/Implementation/WorkOrderWebhookDataService.cs +++ b/SeaHaven.DataServices/Implementation/WorkOrderWebhookDataService.cs @@ -20,6 +20,32 @@ namespace SeaHaven.DataServices.Implementation public async Task ApplyAsync( WorkOrderWebhookMutation mutation, CancellationToken cancellationToken) + { + // A mutation that cancels an existing work order also cancels its pending uplifts, + // so it runs under the same per-work-order lock as uplift creation: a create in + // flight either commits first (and is cancelled here) or sees the cancelled order. + var current = mutation.IsStateEvent && mutation.LifecycleStatus.HasValue + ? await _context.workOrders + .AsNoTracking() + .Where(w => w.ExternalWorkOrderId == mutation.ExternalWorkOrderId) + .Select(w => new { w.Id, w.LifecycleStatus }) + .SingleOrDefaultAsync(cancellationToken) + : null; + + if (current is null || !PendingUpliftCancellation.Applies(current.LifecycleStatus, mutation.LifecycleStatus)) + return await ApplyUnlockedAsync(mutation, cancellationToken); + + return await WorkOrderMutationLock.RunAsync( + _context, + current.Id, + ct => ApplyUnlockedAsync(mutation, ct), + cancellationToken, + commitWhen: result => result.Status == WorkOrderWebhookPersistenceStatus.Applied); + } + + private async Task ApplyUnlockedAsync( + WorkOrderWebhookMutation mutation, + CancellationToken cancellationToken) { var prior = await _context.WorkOrderWebhookDeliveries .AsNoTracking() diff --git a/SeaHavenIndustries.Tests/WorkOrderCrmCancelUpliftTests.cs b/SeaHavenIndustries.Tests/WorkOrderCrmCancelUpliftTests.cs index 8b7c42a..e576095 100644 --- a/SeaHavenIndustries.Tests/WorkOrderCrmCancelUpliftTests.cs +++ b/SeaHavenIndustries.Tests/WorkOrderCrmCancelUpliftTests.cs @@ -93,6 +93,43 @@ public sealed class WorkOrderCrmCancelUpliftTests Assert.Equal(UpliftStatus.Withdrawn, upliftAudit.NewValue); } + [Fact] + public async Task CrmCancel_WaitsForAnUpliftCreateInFlightAndCancelsWhatItCommitted() + { + var options = await SeedInMemoryAsync(); + var createEntered = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + var releaseCreate = new TaskCompletionSource(TaskCreationOptions.RunContinuationsAsynchronously); + + // An uplift create holding the work order's lock, committing its Pending request + // only after the CRM cancel has started. + await using var createContext = new ApplicationDbContext(options); + var create = new UpliftDataService(createContext).ExecuteWorkOrderMutationAsync(1, async ct => + { + createEntered.SetResult(); + await releaseCreate.Task; + createContext.DispatchUpliftRequests.Add(WorkOrderBoardPatchCancelUpliftTests.Uplift(200, "Pending", 800m)); + await createContext.SaveChangesAsync(ct); + return true; + }, CancellationToken.None); + await createEntered.Task; + + await using var crmContext = new ApplicationDbContext(options); + var cancel = new WorkOrderWebhookDataService(crmContext) + .ApplyAsync(CrmMutation("work_order.cancelled", cancelled: true), CancellationToken.None); + + var finishedFirst = await Task.WhenAny(cancel, Task.Delay(TimeSpan.FromSeconds(2))); + Assert.NotSame(cancel, finishedFirst); + + releaseCreate.SetResult(); + await create; + await cancel; + + await using var verify = new ApplicationDbContext(options); + Assert.Equal(LifecycleStatus.Canceled, verify.workOrders.Single().LifecycleStatus); + Assert.Equal(UpliftStatus.Withdrawn, verify.DispatchUpliftRequests.Single(u => u.Id == 200).Status); + Assert.Single(verify.WorkOrderAuditLogs, log => log.Action == "uplift_cancel"); + } + [Fact] public async Task CrmUpdateThatDoesNotCancel_LeavesPendingUpliftPending() {