fix(work-orders): serialize uplift create and atomic cancel (SH-196)

This commit is contained in:
arthur.bassi 2026-08-18 10:20:05 -03:00
parent 4c15669aff
commit aeface594a
7 changed files with 241 additions and 9 deletions

View file

@ -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<int, SemaphoreSlim> 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<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);
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);
}
}
}

View file

@ -28,5 +28,9 @@ namespace SeaHaven.DataServices.Interfaces
Task<List<DispatchUpliftRequest>> GetDueForEscalationAsync(int count, CancellationToken cancellationToken);
Task StageAsync(DispatchUpliftRequest request, CancellationToken cancellationToken);
Task SaveChangesAsync(CancellationToken cancellationToken);
Task<T> ExecuteWorkOrderMutationAsync<T>(
int workOrderId,
Func<CancellationToken, Task<T>> work,
CancellationToken cancellationToken);
}
}

View file

@ -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,

View file

@ -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<WorkOrderUpliftDto?> 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<WorkOrderUpliftDto> 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);

View file

@ -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<ApplicationDbContext>()
.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<InvalidOperationException>(() =>
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<bool> VendorExistsAsync(int vendorId, CancellationToken cancellationToken)
=> _inner.VendorExistsAsync(vendorId, cancellationToken);
public Task<int> GetMaxWorkOrderIdAsync(CancellationToken cancellationToken)
=> _inner.GetMaxWorkOrderIdAsync(cancellationToken);
public Task<WorkOrder?> GetTrackedWorkOrderAsync(int workOrderId, CancellationToken cancellationToken)
=> _inner.GetTrackedWorkOrderAsync(workOrderId, cancellationToken);
public Task<Dispatch?> GetTrackedDispatchAsync(int dispatchId, int workOrderId, CancellationToken cancellationToken)
=> _inner.GetTrackedDispatchAsync(dispatchId, workOrderId, cancellationToken);
public Task<Dispatch?> 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<CancellationToken, Task> work, CancellationToken cancellationToken)
=> _inner.ExecuteTransactionalAsync(work, cancellationToken);
public Task<BoardSaveOutcome> SaveAsync(CancellationToken cancellationToken)
=> throw new InvalidOperationException("forced late failure");
}
private sealed class NoOpUpliftService : IWorkOrderUpliftService
{
public Task<IReadOnlyList<WorkOrderUpliftDto>?> ListAsync(

View file

@ -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]

View file

@ -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<ApplicationDbContext>()
.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);
}
}