From aeface594af5bf31ac80040480890ca461397a93 Mon Sep 17 00:00:00 2001 From: "arthur.bassi" Date: Tue, 18 Aug 2026 10:20:05 -0300 Subject: [PATCH] fix(work-orders): serialize uplift create and atomic cancel (SH-196) --- .../Implementation/UpliftDataService.cs | 45 +++++++ .../Interfaces/IUpliftDataService.cs | 4 + .../Implementation/WorkOrderDetailService.cs | 1 + .../Implementation/WorkOrderUpliftService.cs | 45 +++++-- .../WorkOrderBoardCancelServiceTests.cs | 113 ++++++++++++++++++ .../WorkOrderPhase6Tests.cs | 5 + .../WorkOrderUpliftServiceTests.cs | 37 ++++++ 7 files changed, 241 insertions(+), 9 deletions(-) diff --git a/SeaHaven.DataServices/Implementation/UpliftDataService.cs b/SeaHaven.DataServices/Implementation/UpliftDataService.cs index b65093e..fa8d76a 100644 --- a/SeaHaven.DataServices/Implementation/UpliftDataService.cs +++ b/SeaHaven.DataServices/Implementation/UpliftDataService.cs @@ -1,3 +1,4 @@ +using System.Collections.Concurrent; using Data.SeaHavenIndustries; using Microsoft.EntityFrameworkCore; using SeaHaven.DataServices.Interfaces; @@ -6,6 +7,7 @@ namespace SeaHaven.DataServices.Implementation { public class UpliftDataService : IUpliftDataService { + private static readonly ConcurrentDictionary WorkOrderGates = new(); private readonly ApplicationDbContext _context; public UpliftDataService(ApplicationDbContext context) @@ -338,5 +340,48 @@ namespace SeaHaven.DataServices.Implementation { await _context.SaveChangesAsync(cancellationToken); } + + public async Task ExecuteWorkOrderMutationAsync( + int workOrderId, + Func> work, + CancellationToken cancellationToken) + { + var gate = WorkOrderGates.GetOrAdd(workOrderId, _ => new SemaphoreSlim(1, 1)); + await gate.WaitAsync(cancellationToken); + 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); + } } } diff --git a/SeaHaven.DataServices/Interfaces/IUpliftDataService.cs b/SeaHaven.DataServices/Interfaces/IUpliftDataService.cs index b64bc68..1c4062b 100644 --- a/SeaHaven.DataServices/Interfaces/IUpliftDataService.cs +++ b/SeaHaven.DataServices/Interfaces/IUpliftDataService.cs @@ -28,5 +28,9 @@ namespace SeaHaven.DataServices.Interfaces Task> GetDueForEscalationAsync(int count, CancellationToken cancellationToken); Task StageAsync(DispatchUpliftRequest request, CancellationToken cancellationToken); Task SaveChangesAsync(CancellationToken cancellationToken); + Task ExecuteWorkOrderMutationAsync( + int workOrderId, + Func> work, + CancellationToken cancellationToken); } } diff --git a/SeaHaven.Services/Implementation/WorkOrderDetailService.cs b/SeaHaven.Services/Implementation/WorkOrderDetailService.cs index 0bb6eac..5fa8537 100644 --- a/SeaHaven.Services/Implementation/WorkOrderDetailService.cs +++ b/SeaHaven.Services/Implementation/WorkOrderDetailService.cs @@ -129,6 +129,7 @@ namespace SeaHaven.Services.Implementation RowVersion = row.RowVersion, DispatchRowVersion = row.DispatchRowVersion, PendingUpliftCount = row.PendingUpliftCount, + UpliftSummary = row.UpliftSummary, Description = extended?.Description, Trade = extended?.Trade, Problem = extended?.Problem, diff --git a/SeaHaven.Services/Implementation/WorkOrderUpliftService.cs b/SeaHaven.Services/Implementation/WorkOrderUpliftService.cs index 1111186..45ff8ea 100644 --- a/SeaHaven.Services/Implementation/WorkOrderUpliftService.cs +++ b/SeaHaven.Services/Implementation/WorkOrderUpliftService.cs @@ -63,10 +63,37 @@ namespace SeaHaven.Services.Implementation if (request.Amount <= 0) throw new InvalidOperationException("Uplift amount must be greater than zero"); + var userId = user.FindFirstValue(ClaimTypes.NameIdentifier); + var requesterName = await ResolveUserDisplayNameAsync(userId, cancellationToken); + var notes = request.Notes?.Trim() ?? ""; + var accountFilter = _accountResolver.ResolveAccountFilter(user); + + return await _upliftData.ExecuteWorkOrderMutationAsync( + workOrderId, + ct => CreateLockedAsync( + workOrderId, + request.Amount, + notes, + userId, + requesterName, + accountFilter, + ct), + cancellationToken); + } + + private async Task CreateLockedAsync( + int workOrderId, + decimal amount, + string notes, + string? userId, + string requesterName, + int? accountFilter, + CancellationToken cancellationToken) + { var workOrder = await _detailData.GetWorkOrderForMediaAsync( workOrderId, cancellationToken, - _accountResolver.ResolveAccountFilter(user)); + accountFilter); if (workOrder?.PrimaryDispatchId is not int dispatchId) throw new InvalidOperationException("Work order has no primary dispatch for uplift requests"); @@ -80,17 +107,14 @@ namespace SeaHaven.Services.Implementation if (await _upliftData.HasPendingForWorkOrderAsync(workOrderId, cancellationToken)) throw new InvalidOperationException("An open uplift request already exists for this work order"); - var userId = user.FindFirstValue(ClaimTypes.NameIdentifier); - var requesterName = await ResolveUserDisplayNameAsync(userId, cancellationToken); var now = _timeProvider.GetUtcNow().UtcDateTime; var current = dispatch.NTEAmount ?? 0m; - var notes = request.Notes?.Trim() ?? ""; var consumed = await _upliftData.SumAutoApprovedAmountForWorkOrderAsync(workOrderId, cancellationToken); var remaining = WorkOrderUpliftAllowance.Remaining( WorkOrderUpliftAllowance.CapFor(workOrder.WorkOrderType), consumed); - if (WorkOrderUpliftAllowance.AutoApproves(request.Amount, remaining)) + if (WorkOrderUpliftAllowance.AutoApproves(amount, remaining)) { return await PersistCreatedAsync( dispatch, @@ -98,7 +122,7 @@ namespace SeaHaven.Services.Implementation userId, requesterName, current, - request.Amount, + amount, notes, UpliftStatus.NoApprovalRequired, requiredTier: 0, @@ -115,7 +139,7 @@ namespace SeaHaven.Services.Implementation userId, requesterName, current, - request.Amount, + amount, notes, UpliftStatus.Pending, requiredTier: 1, @@ -281,8 +305,6 @@ namespace SeaHaven.Services.Implementation cancellationToken, isStatusTransition: true); } - - await _upliftData.SaveChangesAsync(cancellationToken); } private async Task PersistCreatedAsync( @@ -315,6 +337,11 @@ namespace SeaHaven.Services.Implementation ExpiresAt = expiresAt, NotificationStatus = notificationStatus, }; + if (status == UpliftStatus.NoApprovalRequired) + { + dispatch.NTEAmount = currentNte + amount; + dispatch.LastModificationTime = now; + } await _upliftData.StageAsync(created, cancellationToken); await StageAuditAsync(dispatch, workOrderId, userId, currentNte, amount, auditAction, now, cancellationToken); await _upliftData.SaveChangesAsync(cancellationToken); diff --git a/SeaHavenIndustries.Tests/WorkOrderBoardCancelServiceTests.cs b/SeaHavenIndustries.Tests/WorkOrderBoardCancelServiceTests.cs index 1f8e2b0..2aae357 100644 --- a/SeaHavenIndustries.Tests/WorkOrderBoardCancelServiceTests.cs +++ b/SeaHavenIndustries.Tests/WorkOrderBoardCancelServiceTests.cs @@ -3,6 +3,7 @@ using Data.SeaHavenIndustries.Enums; using Microsoft.EntityFrameworkCore; using Microsoft.Extensions.Options; using SeaHaven.DataServices.Implementation; +using SeaHaven.DataServices.Interfaces; using SeaHaven.Services.Configuration; using SeaHaven.Services.DTOs; using SeaHaven.Services.Exceptions; @@ -214,6 +215,118 @@ public class WorkOrderBoardCancelServiceTests Assert.Contains(context.WorkOrderAuditLogs, log => log.Action == "StatusChanged"); } + [Fact] + public async Task Cancel_LateSaveFailure_LeavesPendingUpliftAndOpenWorkOrder() + { + var databaseName = Guid.NewGuid().ToString(); + var options = new DbContextOptionsBuilder() + .UseInMemoryDatabase(databaseName) + .Options; + await using var context = new ApplicationDbContext(options); + await WorkOrderAccountTestHelpers.EnsureAccountAsync(context); + context.Vendors.Add(new Vendor { Id = 1, CompanyName = "Acme HVAC" }); + context.Dispatches.Add(new Dispatch + { + Id = 10, + VendorId = 1, + WorkOrderId = 1, + DispatchNumber = "DIS-10", + Status = "Scheduled", + }); + context.workOrders.Add(new WorkOrder + { + Id = 1, + LifecycleStatus = LifecycleStatus.Incomplete, + Status = "Incomplete", + PrimaryDispatchId = 10, + AccountId = 1, + RowVersion = new byte[] { 1, 0, 0, 0, 0, 0, 0, 1 } + }); + context.DispatchUpliftRequests.Add(new DispatchUpliftRequest + { + Id = 100, + DispatchId = 10, + RequestedNTE = 600m, + Status = "Pending", + RequiredTier = 1, + NotificationStatus = "Pending", + }); + await context.SaveChangesAsync(); + + var boardData = new WorkOrderBoardDataService(context); + var mutationData = new WorkOrderBoardMutationDataService(context); + var boardService = new WorkOrderBoardService(boardData, WorkOrderAccountTestHelpers.Resolver(context)); + var fieldLocks = new WorkOrderFieldLockService(new WorkOrderFieldLockDataService(context)); + var audit = new WorkOrderAuditService(new WorkOrderAuditDataService(context), fieldLocks); + var uplifts = new WorkOrderUpliftService( + new UpliftDataService(context), + new DispatchDataService(context), + new WorkOrderDetailDataService(context), + WorkOrderAccountTestHelpers.Resolver(context), + new UserDataService(context), + TimeProvider.System, + Options.Create(new ApprovalsOptions())); + var cancel = new WorkOrderBoardCancelService( + new ThrowingSaveMutationData(mutationData), + boardService, + audit, + uplifts); + + await Assert.ThrowsAsync(() => + cancel.CancelAsync(1, WorkOrderAccountTestHelpers.OrgWideAdmin(), "actor-1")); + + await using var verify = new ApplicationDbContext(options); + Assert.Equal(LifecycleStatus.Incomplete, Assert.Single(verify.workOrders).LifecycleStatus); + Assert.Equal("Pending", Assert.Single(verify.DispatchUpliftRequests).Status); + Assert.Empty(verify.WorkOrderAuditLogs); + } + + private sealed class ThrowingSaveMutationData : IWorkOrderBoardMutationDataService + { + private readonly IWorkOrderBoardMutationDataService _inner; + + public ThrowingSaveMutationData(IWorkOrderBoardMutationDataService inner) + { + _inner = inner; + } + + public Task VendorExistsAsync(int vendorId, CancellationToken cancellationToken) + => _inner.VendorExistsAsync(vendorId, cancellationToken); + + public Task GetMaxWorkOrderIdAsync(CancellationToken cancellationToken) + => _inner.GetMaxWorkOrderIdAsync(cancellationToken); + + public Task GetTrackedWorkOrderAsync(int workOrderId, CancellationToken cancellationToken) + => _inner.GetTrackedWorkOrderAsync(workOrderId, cancellationToken); + + public Task GetTrackedDispatchAsync(int dispatchId, int workOrderId, CancellationToken cancellationToken) + => _inner.GetTrackedDispatchAsync(dispatchId, workOrderId, cancellationToken); + + public Task GetTrackedDispatchByIdAsync(int dispatchId, CancellationToken cancellationToken) + => _inner.GetTrackedDispatchByIdAsync(dispatchId, cancellationToken); + + public void TrackNewWorkOrder(WorkOrder workOrder) => _inner.TrackNewWorkOrder(workOrder); + + public void TrackNewDispatch(Dispatch dispatch) => _inner.TrackNewDispatch(dispatch); + + public void TrackWorkOrderContact(int workOrderId, int contactId, string? notes) + => _inner.TrackWorkOrderContact(workOrderId, contactId, notes); + + public void SetExpectedWorkOrderVersion(WorkOrder workOrder, byte[] version) + => _inner.SetExpectedWorkOrderVersion(workOrder, version); + + public void SetExpectedDispatchVersion(Dispatch dispatch, byte[] version) + => _inner.SetExpectedDispatchVersion(dispatch, version); + + public bool HasPendingChanges() => _inner.HasPendingChanges(); + + public Task ExecuteTransactionalAsync(Func work, CancellationToken cancellationToken) + => _inner.ExecuteTransactionalAsync(work, cancellationToken); + + public Task SaveAsync(CancellationToken cancellationToken) + => throw new InvalidOperationException("forced late failure"); + } + private sealed class NoOpUpliftService : IWorkOrderUpliftService { public Task?> ListAsync( diff --git a/SeaHavenIndustries.Tests/WorkOrderPhase6Tests.cs b/SeaHavenIndustries.Tests/WorkOrderPhase6Tests.cs index d6bd47e..b5dae37 100644 --- a/SeaHavenIndustries.Tests/WorkOrderPhase6Tests.cs +++ b/SeaHavenIndustries.Tests/WorkOrderPhase6Tests.cs @@ -208,6 +208,11 @@ public class WorkOrderDetailServiceTests Assert.NotNull(detail); Assert.Equal(1, detail!.Info.PendingUpliftCount); + Assert.NotNull(detail.Info.UpliftSummary); + Assert.True(detail.Info.UpliftSummary!.HasUplift); + Assert.Equal(1, detail.Info.UpliftSummary.PendingCount); + Assert.Equal("pending", detail.Info.UpliftSummary.PrimaryStatus); + Assert.Equal(1500m, detail.Info.UpliftSummary.Amount); } [Fact] diff --git a/SeaHavenIndustries.Tests/WorkOrderUpliftServiceTests.cs b/SeaHavenIndustries.Tests/WorkOrderUpliftServiceTests.cs index 0ac7d7a..9ca211e 100644 --- a/SeaHavenIndustries.Tests/WorkOrderUpliftServiceTests.cs +++ b/SeaHavenIndustries.Tests/WorkOrderUpliftServiceTests.cs @@ -120,6 +120,7 @@ public sealed class WorkOrderUpliftServiceTests Assert.NotNull(created); Assert.Equal("auto_approved", created!.Status); + Assert.Equal(1400m, context.Dispatches.Single(d => d.Id == 10).NTEAmount); } [Fact] @@ -251,6 +252,7 @@ public sealed class WorkOrderUpliftServiceTests var service = NewService(context); await service.WithdrawPendingForWorkOrderAsync(workOrder.Id, "actor-1", CancellationToken.None); + await context.SaveChangesAsync(); Assert.Equal("Withdrawn", Assert.Single(context.DispatchUpliftRequests).Status); Assert.Contains(context.WorkOrderAuditLogs, log => log.Action == "uplift_cancel"); @@ -337,4 +339,39 @@ public sealed class WorkOrderUpliftServiceTests Assert.Equal("revoked", revoked!.Status); Assert.Equal(1000m, dispatch.NTEAmount); } + + [Fact] + public async Task CreateAsync_ConcurrentRequests_PreserveOnePendingAndCap() + { + var databaseName = Guid.NewGuid().ToString(); + var options = new DbContextOptionsBuilder() + .UseInMemoryDatabase(databaseName) + .Options; + + await using (var seed = new ApplicationDbContext(options)) + { + await SeedWorkOrderAsync(seed); + } + + await using var firstContext = new ApplicationDbContext(options); + await using var secondContext = new ApplicationDbContext(options); + var first = NewService(firstContext); + var second = NewService(secondContext); + var request = new CreateWorkOrderUpliftRequestDto { Amount = 400m, Notes = "Concurrent" }; + + var results = await Task.WhenAll( + first.CreateAsync(1, request, Dispatcher(), CancellationToken.None), + second.CreateAsync(1, request, Dispatcher(), CancellationToken.None)); + + await using var verify = new ApplicationDbContext(options); + var rows = verify.DispatchUpliftRequests.ToList(); + var autoApproved = rows.Where(row => row.Status == "NoApprovalRequired").ToList(); + var pending = rows.Where(row => row.Status == "Pending").ToList(); + + Assert.Equal(2, results.Length); + Assert.True(results.All(result => result != null)); + Assert.Single(autoApproved); + Assert.Single(pending); + Assert.True(autoApproved.Sum(row => row.RequestedNTE) <= 500m); + } }