using System.Data.Common; using Data.SeaHavenIndustries; using Microsoft.Data.Sqlite; using Microsoft.EntityFrameworkCore; using Microsoft.EntityFrameworkCore.Diagnostics; using SeaHaven.DataServices.Implementation; namespace SeaHavenIndustries.Tests; public sealed class UpliftDataServiceTransactionTests { [Theory] [InlineData("Cannot insert duplicate key row with unique index 'IX_DispatchUpliftRequests_DispatchId'.", true)] [InlineData("Cannot insert duplicate key row with unique index 'IX_DispatchUpliftRequests_DispatchId_RequestKey'.", false)] public void IsActiveDispatchIndexViolationMessage_MatchesOnlyTheActiveDispatchIndex( string message, bool expected) { Assert.Equal(expected, UpliftDataService.IsActiveDispatchIndexViolationMessage(message)); } [Fact] public async Task SaveChangesAsync_ActiveUpliftOnSameDispatch_MapsConcurrentInsertConflict() { await using var connection = new SqliteConnection("DataSource=:memory:"); await connection.OpenAsync(); var options = new DbContextOptionsBuilder() .UseSqlite(connection) .Options; await using var context = new ApplicationDbContext(options); await context.Database.EnsureCreatedAsync(); context.Vendors.Add(new Vendor { Id = 1, CompanyName = "Acme HVAC" }); context.Dispatches.Add(new Dispatch { Id = 10, VendorId = 1, DispatchNumber = "DIS-10", Status = "Scheduled", }); await context.SaveChangesAsync(); var data = new UpliftDataService(context); await data.StageAsync(new DispatchUpliftRequest { DispatchId = 10, RequestedNTE = 1500m, Status = "Pending", RequiredTier = 1, NotificationStatus = "Pending", }, CancellationToken.None); await data.SaveChangesAsync(CancellationToken.None); // The database index is the final guard when competing requests pass // their dispatch-scoped active-request checks. var exception = await Assert.ThrowsAsync( () => data.ExecuteWorkOrderMutationAsync( 2, async cancellationToken => { await data.StageAsync(new DispatchUpliftRequest { DispatchId = 10, RequestedNTE = 1700m, Status = "Pending", RequiredTier = 1, NotificationStatus = "Pending", }, cancellationToken); await data.SaveChangesAsync(cancellationToken); return true; }, CancellationToken.None)); Assert.Equal("An active uplift request already exists for this dispatch.", exception.Message); Assert.Equal(1, await context.DispatchUpliftRequests.CountAsync()); } [Fact] public async Task SaveChangesAsync_RequestKeyConflict_IsNotMappedToActiveDispatchConflict() { await using var connection = new SqliteConnection("DataSource=:memory:"); await connection.OpenAsync(); var options = new DbContextOptionsBuilder() .UseSqlite(connection) .Options; await using var context = new ApplicationDbContext(options); await context.Database.EnsureCreatedAsync(); context.Vendors.Add(new Vendor { Id = 1, CompanyName = "Acme HVAC" }); context.Dispatches.Add(new Dispatch { Id = 10, VendorId = 1, DispatchNumber = "DIS-10", Status = "Scheduled", }); await context.SaveChangesAsync(); var data = new UpliftDataService(context); await data.StageAsync(new DispatchUpliftRequest { DispatchId = 10, RequestKey = "same-key", RequestedNTE = 1500m, Status = "Approved", RequiredTier = 1, NotificationStatus = "Sent", }, CancellationToken.None); await data.SaveChangesAsync(CancellationToken.None); await data.StageAsync(new DispatchUpliftRequest { DispatchId = 10, RequestKey = "same-key", RequestedNTE = 1700m, Status = "Approved", RequiredTier = 1, NotificationStatus = "Sent", }, CancellationToken.None); await Assert.ThrowsAsync( () => data.SaveChangesAsync(CancellationToken.None)); } [Fact] public async Task ExecuteWorkOrderMutationAsync_ReleasesGate_WhenTransactionInitializationFails() { await using var connection = new SqliteConnection("DataSource=:memory:"); await connection.OpenAsync(); var failingOptions = new DbContextOptionsBuilder() .UseSqlite(connection) .AddInterceptors(new ThrowOnceTransactionInterceptor()) .Options; await using (var failedContext = new ApplicationDbContext(failingOptions)) { var failedService = new UpliftDataService(failedContext); await Assert.ThrowsAsync(() => failedService.ExecuteWorkOrderMutationAsync( 387, _ => Task.FromResult(true), CancellationToken.None)); } var retryOptions = new DbContextOptionsBuilder() .UseSqlite(connection) .Options; await using var retryContext = new ApplicationDbContext(retryOptions); var retryService = new UpliftDataService(retryContext); var retry = retryService.ExecuteWorkOrderMutationAsync( 387, _ => Task.FromResult("entered"), CancellationToken.None); Assert.Equal("entered", await retry.WaitAsync(TimeSpan.FromSeconds(1))); } private sealed class ThrowOnceTransactionInterceptor : DbTransactionInterceptor { private int _remaining = 1; public override InterceptionResult TransactionStarting( DbConnection connection, TransactionStartingEventData eventData, InterceptionResult result) { ThrowOnce(); return base.TransactionStarting(connection, eventData, result); } public override ValueTask> TransactionStartingAsync( DbConnection connection, TransactionStartingEventData eventData, InterceptionResult result, CancellationToken cancellationToken = default) { ThrowOnce(); return base.TransactionStartingAsync(connection, eventData, result, cancellationToken); } private void ThrowOnce() { if (Interlocked.Exchange(ref _remaining, 0) == 1) throw new InvalidOperationException("Forced transaction initialization failure"); } } }