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