Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<string?>(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()
{
Expand All @@ -139,8 +201,13 @@ public async Task DeleteProcessedRemovesOldTerminalRowsButKeepsRecentAndPending(
remaining.ShouldContain(message => message.ProcessedOnUtc >= now.AddDays(-7));
}

private static EntityFrameworkOutboxStore<TestDbContext> NewStore(TestDbContext context, int maxRetries) =>
new(context, TimeProvider.System, Options.Create(new OutboxProcessingOptions { MaxRetries = maxRetries }));
private static EntityFrameworkOutboxStore<TestDbContext> 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()
{
Expand All @@ -155,6 +222,53 @@ private static EntityFrameworkOutboxStore<TestDbContext> NewStore(TestDbContext

private TestDbContext NewContext() => new(new DbContextOptionsBuilder<TestDbContext>().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<string?> 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<TestDbContext> options) : DbContext(options), ISaveOutboxMessages
{
public DbSet<OutboxMessage> OutboxMessages => Set<OutboxMessage>();
Expand Down