Skip to content

Commit caaa6d2

Browse files
authored
Merge pull request #293 from Vulthil/fix/outbox-sync-capture-and-validation
Capture domain events and validate outbox delay options for sync SaveChanges callers
2 parents b19d9fd + 7444feb commit caaa6d2

5 files changed

Lines changed: 109 additions & 14 deletions

File tree

src/Vulthil.SharedKernel.Outbox.EntityFrameworkCore/DomainEventsToOutboxMessageSaveChangesInterceptor.cs

Lines changed: 41 additions & 14 deletions
Original file line numberDiff line numberDiff line change
@@ -15,16 +15,53 @@ public sealed class DomainEventsToOutboxMessageSaveChangesInterceptor(TimeProvid
1515
{
1616
private readonly TimeProvider _timeProvider = timeProvider;
1717

18+
/// <summary>
19+
/// Captures domain events from tracked aggregate roots and stores them as outbox messages before persisting changes.
20+
/// </summary>
21+
public override InterceptionResult<int> SavingChanges(DbContextEventData eventData, InterceptionResult<int> result)
22+
{
23+
CaptureDomainEvents(eventData.Context);
24+
return result;
25+
}
26+
1827
/// <summary>
1928
/// Captures domain events from tracked aggregate roots and stores them as outbox messages before persisting changes.
2029
/// </summary>
2130
public override ValueTask<InterceptionResult<int>> SavingChangesAsync(DbContextEventData eventData, InterceptionResult<int> result, CancellationToken cancellationToken = default)
2231
{
23-
var dbContext = eventData.Context;
32+
CaptureDomainEvents(eventData.Context);
33+
return base.SavingChangesAsync(eventData, result, cancellationToken);
34+
}
35+
36+
/// <summary>
37+
/// Wakes the outbox relay after a save that committed outside an explicit transaction, so domain events captured
38+
/// by a bare <c>SaveChanges</c> are relayed promptly instead of waiting for the next poll. When a transaction is
39+
/// open the relay is signalled on commit by the transaction-commit interceptor instead, so this skips that case to
40+
/// avoid waking the relay before the rows are committed and visible.
41+
/// </summary>
42+
public override int SavedChanges(SaveChangesCompletedEventData eventData, int result)
43+
{
44+
WakeRelayIfNeeded(eventData.Context);
45+
return result;
46+
}
47+
48+
/// <summary>
49+
/// Wakes the outbox relay after a save that committed outside an explicit transaction, so domain events captured
50+
/// by a bare <c>SaveChanges</c> are relayed promptly instead of waiting for the next poll. When a transaction is
51+
/// open the relay is signalled on commit by the transaction-commit interceptor instead, so this skips that case to
52+
/// avoid waking the relay before the rows are committed and visible.
53+
/// </summary>
54+
public override ValueTask<int> SavedChangesAsync(SaveChangesCompletedEventData eventData, int result, CancellationToken cancellationToken = default)
55+
{
56+
WakeRelayIfNeeded(eventData.Context);
57+
return base.SavedChangesAsync(eventData, result, cancellationToken);
58+
}
2459

60+
private void CaptureDomainEvents(DbContext? dbContext)
61+
{
2562
if (dbContext is not ISaveOutboxMessages dbContextWithOutboxMessages)
2663
{
27-
return base.SavingChangesAsync(eventData, result, cancellationToken);
64+
return;
2865
}
2966

3067
Activity? activity = null;
@@ -55,24 +92,14 @@ public override ValueTask<InterceptionResult<int>> SavingChangesAsync(DbContextE
5592
}).ToList();
5693

5794
dbContextWithOutboxMessages.OutboxMessages.AddRange(outboxMessages);
58-
59-
return base.SavingChangesAsync(eventData, result, cancellationToken);
6095
}
6196

62-
/// <summary>
63-
/// Wakes the outbox relay after a save that committed outside an explicit transaction, so domain events captured
64-
/// by a bare <c>SaveChanges</c> are relayed promptly instead of waiting for the next poll. When a transaction is
65-
/// open the relay is signalled on commit by the transaction-commit interceptor instead, so this skips that case to
66-
/// avoid waking the relay before the rows are committed and visible.
67-
/// </summary>
68-
public override ValueTask<int> SavedChangesAsync(SaveChangesCompletedEventData eventData, int result, CancellationToken cancellationToken = default)
97+
private void WakeRelayIfNeeded(DbContext? dbContext)
6998
{
70-
if (ShouldWakeRelay(eventData.Context))
99+
if (ShouldWakeRelay(dbContext))
71100
{
72101
signal.Notify();
73102
}
74-
75-
return base.SavedChangesAsync(eventData, result, cancellationToken);
76103
}
77104

78105
private static bool ShouldWakeRelay(DbContext? dbContext) =>
Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1 +1,3 @@
11
#nullable enable
2+
override Vulthil.SharedKernel.Outbox.EntityFrameworkCore.DomainEventsToOutboxMessageSaveChangesInterceptor.SavedChanges(Microsoft.EntityFrameworkCore.Diagnostics.SaveChangesCompletedEventData! eventData, int result) -> int
3+
override Vulthil.SharedKernel.Outbox.EntityFrameworkCore.DomainEventsToOutboxMessageSaveChangesInterceptor.SavingChanges(Microsoft.EntityFrameworkCore.Diagnostics.DbContextEventData! eventData, Microsoft.EntityFrameworkCore.Diagnostics.InterceptionResult<int> result) -> Microsoft.EntityFrameworkCore.Diagnostics.InterceptionResult<int>

src/Vulthil.SharedKernel.Outbox/OutboxEngineServiceCollectionExtensions.cs

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -39,6 +39,9 @@ public static IServiceCollection AddOutboxEngine(
3939
&& o.Retention.SweepInterval > TimeSpan.Zero
4040
&& o.Retention.BatchSize >= 1,
4141
"Outbox retention requires RetentionPeriod and SweepInterval greater than zero and BatchSize of at least 1 when enabled.")
42+
.Validate(
43+
o => o.MaxDelaySeconds >= o.OutboxProcessingDelaySeconds,
44+
"MaxDelaySeconds must be greater than or equal to OutboxProcessingDelaySeconds.")
4245
.ValidateDataAnnotations()
4346
.ValidateOnStart();
4447

tests/Vulthil.SharedKernel.Outbox.EntityFrameworkCore.Tests/DomainEventsToOutboxMessageSaveChangesInterceptorTests.cs

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,24 @@ public async Task NonTransactionalSaveCapturingDomainEventsWakesTheRelay()
3737
GetMock<IOutboxSignal>().Verify(signal => signal.Notify(), Times.Once());
3838
}
3939

40+
[Fact]
41+
public async Task SyncSaveCapturingDomainEventsWakesTheRelay()
42+
{
43+
// Arrange
44+
await using var context = NewContext();
45+
var aggregate = new TestAggregate(Guid.NewGuid());
46+
aggregate.RaiseSomething();
47+
context.Aggregates.Add(aggregate);
48+
49+
// Act
50+
SaveChangesSynchronously(context);
51+
52+
// Assert
53+
GetMock<IOutboxSignal>().Verify(signal => signal.Notify(), Times.Once());
54+
var captured = await context.OutboxMessages.SingleAsync(CancellationToken);
55+
captured.Type.ShouldBe(typeof(TestDomainEvent).FullName);
56+
}
57+
4058
[Fact]
4159
public async Task SaveInsideAnExplicitTransactionDoesNotWakeTheRelay()
4260
{
@@ -113,6 +131,8 @@ private TestDbContext NewContext(bool withInterceptor = true)
113131
return new TestDbContext(builder.Options);
114132
}
115133

134+
private static void SaveChangesSynchronously(TestDbContext context) => context.SaveChanges();
135+
116136
public sealed class TestDbContext(DbContextOptions<TestDbContext> options) : DbContext(options), ISaveOutboxMessages
117137
{
118138
public DbSet<OutboxMessage> OutboxMessages => Set<OutboxMessage>();
Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,43 @@
1+
using Microsoft.Extensions.DependencyInjection;
2+
using Microsoft.Extensions.Options;
3+
using Vulthil.xUnit;
4+
5+
namespace Vulthil.SharedKernel.Outbox.Tests;
6+
7+
public sealed class OutboxEngineServiceCollectionExtensionsTests : BaseUnitTestCase
8+
{
9+
[Fact]
10+
public void MaxDelaySecondsLessThanOutboxProcessingDelaySecondsFailsValidationAtStartup()
11+
{
12+
// Arrange
13+
var services = new ServiceCollection();
14+
services.AddOutboxEngine(o =>
15+
{
16+
o.OutboxProcessingDelaySeconds = 100;
17+
o.MaxDelaySeconds = 1;
18+
});
19+
using var provider = services.BuildServiceProvider();
20+
21+
// Act & Assert
22+
Should.Throw<OptionsValidationException>(() => provider.GetRequiredService<IOptions<OutboxProcessingOptions>>().Value);
23+
}
24+
25+
[Fact]
26+
public void MaxDelaySecondsEqualToOutboxProcessingDelaySecondsPassesValidation()
27+
{
28+
// Arrange
29+
var services = new ServiceCollection();
30+
services.AddOutboxEngine(o =>
31+
{
32+
o.OutboxProcessingDelaySeconds = 5;
33+
o.MaxDelaySeconds = 5;
34+
});
35+
using var provider = services.BuildServiceProvider();
36+
37+
// Act
38+
var options = provider.GetRequiredService<IOptions<OutboxProcessingOptions>>().Value;
39+
40+
// Assert
41+
options.MaxDelaySeconds.ShouldBe(5);
42+
}
43+
}

0 commit comments

Comments
 (0)