diff --git a/tests/Vulthil.SharedKernel.Outbox.EntityFrameworkCore.Tests/EntityFrameworkOutboxStoreTests.cs b/tests/Vulthil.SharedKernel.Outbox.EntityFrameworkCore.Tests/EntityFrameworkOutboxStoreTests.cs index 6e69511b..d8fb8124 100644 --- a/tests/Vulthil.SharedKernel.Outbox.EntityFrameworkCore.Tests/EntityFrameworkOutboxStoreTests.cs +++ b/tests/Vulthil.SharedKernel.Outbox.EntityFrameworkCore.Tests/EntityFrameworkOutboxStoreTests.cs @@ -113,6 +113,68 @@ await store.ProcessBatchAsync((data, _) => dispatched.ShouldBe([first, second]); } + [Fact] + public async Task ThrottlesParallelDispatchToTheConfiguredMaxDegreeOfParallelism() + { + // Arrange + await using var seed = NewContext(); + for (var i = 0; i < 6; i++) + { + seed.OutboxMessages.Add(NewMessage(Guid.CreateVersion7(), DateTimeOffset.UtcNow)); + } + await seed.SaveChangesAsync(CancellationToken); + await using var context = NewContext(); + var store = NewStore(context, maxRetries: 3, enableParallelPublishing: true, maxDegreeOfParallelism: 2); + var dispatcher = new ConcurrencyTrackingDispatcher(saturationCount: 2); + + // Act + var processTask = store.ProcessBatchAsync((_, token) => dispatcher.DispatchAsync(token), CancellationToken); + await dispatcher.SaturationReached.WaitAsync(TimeSpan.FromSeconds(30), CancellationToken); + dispatcher.Release(); + var processed = await processTask; + + // Assert + processed.ShouldBe(6); + dispatcher.PeakConcurrency.ShouldBe(2); + } + + [Fact] + public async Task RecordsEachMessageOutcomeIndividuallyWhenAParallelBatchHasAFailure() + { + // Arrange + var failing = new Guid("00000000-0000-0000-0000-000000000001"); + await using var seed = NewContext(); + seed.OutboxMessages.Add(NewMessage(failing, DateTimeOffset.UtcNow)); + seed.OutboxMessages.Add(NewMessage(new Guid("00000000-0000-0000-0000-000000000002"), DateTimeOffset.UtcNow)); + seed.OutboxMessages.Add(NewMessage(new Guid("00000000-0000-0000-0000-000000000003"), DateTimeOffset.UtcNow)); + await seed.SaveChangesAsync(CancellationToken); + await using var context = NewContext(); + var store = NewStore(context, maxRetries: 3, enableParallelPublishing: true); + + // Act + var processed = await store.ProcessBatchAsync( + (data, _) => Task.FromResult(data.Id == failing ? "boom" : null), + CancellationToken); + + // Assert + processed.ShouldBe(2); + await using var verify = NewContext(); + var failed = await verify.OutboxMessages.SingleAsync(message => message.Id == failing, CancellationToken); + failed.ProcessedOnUtc.ShouldBeNull(); + failed.FailedOnUtc.ShouldBeNull(); + failed.RetryCount.ShouldBe(1); + failed.Error.ShouldBe("boom"); + var succeeded = await verify.OutboxMessages.Where(message => message.Id != failing).ToListAsync(CancellationToken); + succeeded.Count.ShouldBe(2); + foreach (var message in succeeded) + { + message.ProcessedOnUtc.ShouldNotBeNull(); + message.FailedOnUtc.ShouldBeNull(); + message.RetryCount.ShouldBe(0); + message.Error.ShouldBeNull(); + } + } + [Fact] public async Task DeleteProcessedRemovesOldTerminalRowsButKeepsRecentAndPending() { @@ -139,8 +201,13 @@ public async Task DeleteProcessedRemovesOldTerminalRowsButKeepsRecentAndPending( remaining.ShouldContain(message => message.ProcessedOnUtc >= now.AddDays(-7)); } - private static EntityFrameworkOutboxStore NewStore(TestDbContext context, int maxRetries) => - new(context, TimeProvider.System, Options.Create(new OutboxProcessingOptions { MaxRetries = maxRetries })); + private static EntityFrameworkOutboxStore NewStore(TestDbContext context, int maxRetries, bool enableParallelPublishing = false, int maxDegreeOfParallelism = 4) => + new(context, TimeProvider.System, Options.Create(new OutboxProcessingOptions + { + MaxRetries = maxRetries, + EnableParallelPublishing = enableParallelPublishing, + MaxDegreeOfParallelism = maxDegreeOfParallelism, + })); private static OutboxMessage NewMessage(Guid id, DateTimeOffset occurredOn, DateTimeOffset? failedOnUtc = null, DateTimeOffset? processedOnUtc = null) => new() { @@ -155,6 +222,53 @@ private static EntityFrameworkOutboxStore NewStore(TestDbContext private TestDbContext NewContext() => new(new DbContextOptionsBuilder().UseSqlite(_connection).Options); + private sealed class ConcurrencyTrackingDispatcher(int saturationCount) + { + private readonly TaskCompletionSource _saturated = new(TaskCreationOptions.RunContinuationsAsynchronously); + private readonly TaskCompletionSource _gate = new(TaskCreationOptions.RunContinuationsAsynchronously); + private int _inFlight; + private int _peak; + private int _arrivals; + + public Task SaturationReached => _saturated.Task; + + public int PeakConcurrency => Volatile.Read(ref _peak); + + public async Task DispatchAsync(CancellationToken cancellationToken) + { + RecordArrival(); + await _gate.Task.WaitAsync(cancellationToken); + Interlocked.Decrement(ref _inFlight); + return null; + } + + public void Release() => _gate.TrySetResult(); + + private void RecordArrival() + { + UpdatePeak(Interlocked.Increment(ref _inFlight)); + if (Interlocked.Increment(ref _arrivals) >= saturationCount) + { + _saturated.TrySetResult(); + } + } + + private void UpdatePeak(int current) + { + var seen = Volatile.Read(ref _peak); + while (current > seen) + { + var previous = Interlocked.CompareExchange(ref _peak, current, seen); + if (previous == seen) + { + return; + } + + seen = previous; + } + } + } + public sealed class TestDbContext(DbContextOptions options) : DbContext(options), ISaveOutboxMessages { public DbSet OutboxMessages => Set();