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.
This commit is contained in:
Alexandre Brandizzi 2026-09-25 19:15:38 -03:00
parent 306ab159cc
commit 77dc3d85c3
4 changed files with 144 additions and 45 deletions

View file

@ -0,0 +1,79 @@
using System.Collections.Concurrent;
using Data.SeaHavenIndustries;
using Microsoft.EntityFrameworkCore;
namespace SeaHaven.DataServices.Helpers
{
/// <summary>
/// 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.
/// </summary>
public static class WorkOrderMutationLock
{
private static readonly ConcurrentDictionary<int, SemaphoreSlim> Gates = new();
/// <summary>
/// Runs <paramref name="work"/> under the work order's gate and row lock. The
/// transaction commits when the work returns and <paramref name="commitWhen"/> (when
/// given) accepts its result; otherwise, or when the work throws, it rolls back.
/// </summary>
public static async Task<T> RunAsync<T>(
ApplicationDbContext context,
int workOrderId,
Func<CancellationToken, Task<T>> work,
CancellationToken cancellationToken,
Func<T, bool>? 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);
}
}
}

View file

@ -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<int, SemaphoreSlim> 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<T> ExecuteWorkOrderMutationAsync<T>(
public Task<T> ExecuteWorkOrderMutationAsync<T>(
int workOrderId,
Func<CancellationToken, Task<T>> 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);
}
}

View file

@ -20,6 +20,32 @@ namespace SeaHaven.DataServices.Implementation
public async Task<WorkOrderWebhookPersistenceResult> 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<WorkOrderWebhookPersistenceResult> ApplyUnlockedAsync(
WorkOrderWebhookMutation mutation,
CancellationToken cancellationToken)
{
var prior = await _context.WorkOrderWebhookDeliveries
.AsNoTracking()

View file

@ -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()
{