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 Task RunAsync( ApplicationDbContext context, int workOrderId, Func> work, CancellationToken cancellationToken, Func? commitWhen = null) => RunManyAsync(context, new[] { workOrderId }, work, cancellationToken, commitWhen); /// /// for a unit of work that spans several work orders (a sync /// batch). Gates and row locks are taken in ascending id order, so two batches never /// wait on each other in a cycle; an empty set runs the work in a plain transaction. /// public static async Task RunManyAsync( ApplicationDbContext context, IEnumerable workOrderIds, Func> work, CancellationToken cancellationToken, Func? commitWhen = null) { var ids = workOrderIds.Distinct().OrderBy(id => id).ToList(); var held = new List(ids.Count); try { foreach (var id in ids) { var gate = Gates.GetOrAdd(id, _ => new SemaphoreSlim(1, 1)); await gate.WaitAsync(cancellationToken); held.Add(gate); } await using var transaction = context.Database.IsRelational() ? await context.Database.BeginTransactionAsync(cancellationToken) : null; try { foreach (var id in ids) await LockRowAsync(context, id, 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 { foreach (var gate in held) 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); } } }