diff --git a/src/Azure/test/Eventuous.Tests.Azure.ServiceBus/IsSerialisableByServiceBus.cs b/src/Azure/test/Eventuous.Tests.Azure.ServiceBus/IsSerialisableByServiceBus.cs index fdf6eab4b..7f8d56729 100644 --- a/src/Azure/test/Eventuous.Tests.Azure.ServiceBus/IsSerialisableByServiceBus.cs +++ b/src/Azure/test/Eventuous.Tests.Azure.ServiceBus/IsSerialisableByServiceBus.cs @@ -19,9 +19,10 @@ public class IsSerialisableByServiceBus { yield return () => 12.34m; yield return () => true; yield return () => 'c'; - yield return () => Guid.NewGuid(); - yield return () => DateTime.UtcNow; - yield return () => DateTimeOffset.UtcNow; + // Parameter values appear in test names, so keep them stable between runs. + yield return () => new Guid("9a8c03f6-9bdf-4072-a43a-772bdae7bb21"); + yield return () => new DateTime(2026, 1, 1, 12, 0, 0, DateTimeKind.Utc); + yield return () => new DateTimeOffset(2026, 1, 1, 12, 0, 0, TimeSpan.Zero); yield return () => TimeSpan.FromMinutes(5); yield return () => new Uri("https://example.com"); yield return () => new MemoryStream(); diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/CancelledMessageTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/CancelledMessageTests.cs index 97698b75b..41d8df628 100644 --- a/src/Core/test/Eventuous.Tests.Subscriptions/CancelledMessageTests.cs +++ b/src/Core/test/Eventuous.Tests.Subscriptions/CancelledMessageTests.cs @@ -16,6 +16,8 @@ namespace Eventuous.Tests.Subscriptions; /// event in flight when a run's own teardown cancelled a parked handler. /// public class CancelledMessageTests { + static readonly TimeSpan RecoveryTimeout = TimeSpan.FromSeconds(30); + /// /// A handler cancelled because the run is ending was never given a verdict, so it must not be /// acknowledged — only redelivered once the successor run comes up. @@ -40,20 +42,20 @@ public async Task Handler_cancelled_by_shutdown_is_not_acknowledged_and_is_redel // The default retry delay is 2s; a short one keeps the test deterministic and fast. options.RetryDelay = TimeSpan.FromMilliseconds(20); - var subscription = new SingleEventSubscription(options, checkpointStore, pipe, loggerFactory); + await using var subscription = new SingleEventSubscription(options, checkpointStore, pipe, loggerFactory); await subscription.Subscribe(_ => { }, (_, _, _) => { }, ct); - (await Wait.Until(() => handler.Parked.IsCompleted, TimeSpan.FromSeconds(5))) - .ShouldBeTrue("the handler should have received the first delivery and parked on it"); + await WaitFor(handler.Parked, "the handler to receive the first delivery", ct).NoContext(); committed.ShouldBeEmpty("nothing should commit while the only delivery so far is still parked, undecided"); subscription.FailCurrentRun(); + await WaitFor(subscription.RedeliveryQueued, "the successor run to queue redelivery", ct).NoContext(); + // The handler blocks again on redelivery, so the test can inspect the checkpoint before it's allowed to succeed. - (await Wait.Until(() => handler.RedeliveryStarted.IsCompleted, TimeSpan.FromSeconds(5))) - .ShouldBeTrue("the successor run should redeliver the event the cancelled handler never finished"); + await WaitFor(handler.RedeliveryStarted, "the handler to start redelivery", ct).NoContext(); // If the fix regresses, DelayedConsume acknowledges the cancelled delivery and this fires. committed.ShouldBeEmpty("the checkpoint must never move past an event whose only delivery was cancelled by shutdown, not decided"); @@ -89,7 +91,7 @@ public async Task Handler_self_cancellation_is_an_ordinary_failure_and_is_skippe CheckpointCommitDelayMs = 10 }; - var subscription = new SingleEventSubscription(options, checkpointStore, pipe, loggerFactory); + await using var subscription = new SingleEventSubscription(options, checkpointStore, pipe, loggerFactory); await subscription.Subscribe(_ => { }, (_, _, _) => { }, ct); @@ -98,8 +100,13 @@ public async Task Handler_self_cancellation_is_an_ordinary_failure_and_is_skippe await subscription.Unsubscribe(_ => { }, ct); } - - + static async Task WaitFor(Task signal, string phase, CancellationToken ct) { + try { + await signal.WaitAsync(RecoveryTimeout, ct).NoContext(); + } catch (TimeoutException exception) { + throw new TimeoutException($"Timed out after {RecoveryTimeout.TotalSeconds} seconds waiting for {phase}.", exception); + } + } record TestOptions : SubscriptionWithCheckpointOptions; @@ -123,6 +130,10 @@ sealed class SingleEventSubscription( null ) { SubscriptionRun? _run; + int _deliveriesQueued; + readonly TaskCompletionSource _redeliveryQueued = new(TaskCreationOptions.RunContinuationsAsynchronously); + + public Task RedeliveryQueued => _redeliveryQueued.Task; /// /// Fails the current run, standing in for a transport drop or any other reason the supervisor tears a run down. @@ -167,6 +178,8 @@ async Task DeliverOnce(SubscriptionRun run) { await HandleInternal(run, context).NoContext(); + if (Interlocked.Increment(ref _deliveriesQueued) == 2) _redeliveryQueued.TrySetResult(); + // Parked rather than returned: a pump ending while its connection is up is read as a drop. await run.Ended.NoContext(); } @@ -199,7 +212,7 @@ public override async ValueTask HandleEvent(IMessageConsume break; case 2: _redeliveryStarted.TrySetResult(); - await _proceedWithSuccess.Task.NoContext(); + await _proceedWithSuccess.Task.WaitAsync(context.CancellationToken).NoContext(); break; } diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/ResubscribeConcurrencyTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/ResubscribeConcurrencyTests.cs index 1c4fe41bf..67f830286 100644 --- a/src/Core/test/Eventuous.Tests.Subscriptions/ResubscribeConcurrencyTests.cs +++ b/src/Core/test/Eventuous.Tests.Subscriptions/ResubscribeConcurrencyTests.cs @@ -14,7 +14,7 @@ namespace Eventuous.Tests.Subscriptions; /// public class ResubscribeConcurrencyTests { /// - /// A whole page of messages fails at once, so Dropped is called once per message, all in one drop window. + /// A whole page of messages is in flight before any handler fails, exercising one concurrent drop window. /// [Test] public async Task Burst_of_nacks_produces_a_single_resubscribe(CancellationToken ct) { @@ -22,9 +22,9 @@ public async Task Burst_of_nacks_produces_a_single_resubscribe(CancellationToken const int messageCount = 8; - var handler = new DeferringHandler(_ => true); + var handler = new BurstFailingHandler(messageCount, ct); - var subscription = new PumpingSubscription( + await using var subscription = new PumpingSubscription( new() { SubscriptionId = "burst-of-nacks", ThrowOnError = true, @@ -34,7 +34,7 @@ public async Task Burst_of_nacks_produces_a_single_resubscribe(CancellationToken new NoOpCheckpointStore(), new ConsumePipe().AddDefaultConsumer(handler), logs, - concurrencyLimit: 4, + concurrencyLimit: messageCount, // Only the first run delivers — a replacement redelivering the same messages would open a second drop window. pump: async (sub, transport, start, run) => { if (transport.Index == 0) { @@ -45,22 +45,29 @@ public async Task Burst_of_nacks_produces_a_single_resubscribe(CancellationToken } ) { ResubscribeDelay = TimeSpan.FromMilliseconds(500) }; - await subscription.Subscribe(_ => { }, (_, _, _) => { }, ct); + try { + await subscription.Subscribe(_ => { }, (_, _, _) => { }, ct); - // All of them have to fail before the resubscribe fires, or this proves nothing about concurrent drops. - (await Wait.Until(() => handler.HandledCount >= messageCount, TimeSpan.FromSeconds(5))) - .ShouldBeTrue($"all {messageCount} messages should have been handled and nacked, got {handler.HandledCount}"); + // A first nack can cancel the pump. Get every delivery into its handler before allowing any nack. + await handler.AllStarted.WaitAsync(TimeSpan.FromSeconds(5), ct).NoContext(); + handler.Release(); - (await Wait.Until(() => subscription.SubscribeCalls > 1, TimeSpan.FromSeconds(5))).ShouldBeTrue("the subscription should have resubscribed"); + (await Wait.Until(() => handler.FailedCount == messageCount, TimeSpan.FromSeconds(5))) + .ShouldBeTrue($"all {messageCount} in-flight handlers should have failed, got {handler.FailedCount}"); - // Give any extra resubscribes scheduled by the other nacks time to show up before asserting. - await Task.Delay(TimeSpan.FromSeconds(1), ct); + (await Wait.Until(() => subscription.SubscribeCalls > 1, TimeSpan.FromSeconds(5))).ShouldBeTrue("the subscription should have resubscribed"); - await subscription.Unsubscribe(_ => { }, ct); + // Give any extra resubscribes scheduled by the other nacks time to show up before asserting. + await Task.Delay(TimeSpan.FromSeconds(1), ct); - subscription.SubscribeCalls.ShouldBe(2, "one initial subscribe plus exactly one resubscribe for the whole drop cycle"); - logs.Count("Resubscribing").ShouldBe(1, "a burst of nacks is one drop cycle, so it gets one 'Resubscribing' line"); - logs.Count("Dropped:").ShouldBe(1, "the drop is reported once per cycle, not once per failing message"); + await subscription.Unsubscribe(_ => { }, ct); + + subscription.SubscribeCalls.ShouldBe(2, "one initial subscribe plus exactly one resubscribe for the whole drop cycle"); + logs.Count("Resubscribing").ShouldBe(1, "a burst of nacks is one drop cycle, so it gets one 'Resubscribing' line"); + logs.Count("Dropped:").ShouldBe(1, "the drop is reported once per cycle, not once per failing message"); + } finally { + handler.Release(); + } } /// @@ -805,6 +812,31 @@ void TrackPumpStarted() { } } + /// + /// Holds all deliveries until the test releases the burst of failures. + /// + sealed class BurstFailingHandler(int messageCount, CancellationToken testCancellation) : BaseEventHandler { + readonly TaskCompletionSource _allStarted = new(TaskCreationOptions.RunContinuationsAsynchronously); + readonly TaskCompletionSource _release = new(TaskCreationOptions.RunContinuationsAsynchronously); + int _started; + int _failed; + + public Task AllStarted => _allStarted.Task; + public int FailedCount => Volatile.Read(ref _failed); + + public void Release() => _release.TrySetResult(); + + public override async ValueTask HandleEvent(IMessageConsumeContext context) { + if (Interlocked.Increment(ref _started) == messageCount) _allStarted.TrySetResult(); + + // All in-flight operations must fail even if the first nack cancels their subscription run. + await _release.Task.WaitAsync(testCancellation).NoContext(); + Interlocked.Increment(ref _failed); + + throw new InvalidOperationException($"Precondition not met for {context.Stream}:{context.GlobalPosition}"); + } + } + /// /// Defers by throwing when a precondition isn't met, so the message is redelivered after the resubscribe. /// diff --git a/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs b/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs index cc43beba1..61532b807 100644 --- a/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs +++ b/src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs @@ -62,7 +62,8 @@ public void ShouldWorkForNormalCase() { } public static IEnumerable> TestData() { - var timestamp = DateTime.Now; + // Parameter values appear in test names, so keep them stable between runs and time zones. + var timestamp = new DateTime(2026, 1, 1, 12, 0, 0, DateTimeKind.Utc); yield return () => ([new(0, 1, timestamp), new(0, 2, timestamp), new(0, 4, timestamp), new(0, 6, timestamp)], new(0, 2, timestamp)); yield return () => ([new(0, 1, timestamp), new(0, 2, timestamp), new(0, 8, timestamp), new(0, 6, timestamp)], new(0, 2, timestamp)); diff --git a/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs b/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs index 94237e556..c1c0bd5ef 100644 --- a/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs +++ b/src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs @@ -36,14 +36,23 @@ protected async Task ShouldMeasureSubscriptionGapCountBase() { static MetricValue? GetValue(MetricValue[] values, string metric) => values.FirstOrDefault(x => x.Name == metric); [Before(Test)] - public async Task InitializeAsync() { + public async Task InitializeAsync(CancellationToken cancellationToken) { var testEvents = TestEvent.CreateMany(fixture.Count); - await fixture.Producer.Produce(fixture.Stream, testEvents, new()); - - while (fixture.Counter.Count < fixture.Count / 2) { - await Task.Delay(100); + await fixture.Producer.Produce(fixture.Stream, testEvents, new(), cancellationToken: cancellationToken); + + var expectedCount = fixture.Count / 2; + using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + cts.CancelAfter(TimeSpan.FromSeconds(30)); + + try { + while (fixture.Counter.Count < expectedCount) { + await Task.Delay(100, cts.Token); + } + } catch (OperationCanceledException ex) when (cts.IsCancellationRequested && !cancellationToken.IsCancellationRequested) { + throw new TimeoutException($"Expected at least {expectedCount} consumed events within 30 seconds, but observed {fixture.Counter.Count}.", ex); } + cancellationToken.ThrowIfCancellationRequested(); fixture.Exporter.Collect(Timeout.Infinite); _values = fixture.Exporter.CollectValues(); diff --git a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBuilder.cs b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBuilder.cs index 6e70c7a03..251fe6fec 100644 --- a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBuilder.cs +++ b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBuilder.cs @@ -3,6 +3,7 @@ using Eventuous.Subscriptions.Consumers; using Eventuous.Sut.Domain; using Eventuous.Tests.Projections.MongoDB.Fixtures; +using Eventuous.Tools; using JetBrains.Annotations; using MongoDB.Driver; using static Eventuous.Sut.Domain.BookingEvents; @@ -14,7 +15,7 @@ public class ProjectWithBuilder(IntegrationFixture fixture) { [Test] [Retry(3)] [MethodDataSource(typeof(CollectionSource), nameof(CollectionSource.TestOptions))] - public async Task ShouldProjectImported(MongoProjectionOptions? options) { + public async Task ShouldProjectImported(MongoProjectionOptions? options, CancellationToken cancellationToken) { var evt = DomainFixture.CreateImportBookingEvent(); var projectionFixture = new ProjectionTestBase(nameof(ProjectWithBuilder), fixture); var id = new BookingId(projectionFixture.CreateId()); @@ -22,48 +23,53 @@ public async Task ShouldProjectImported(MongoProjectionOptions? await projectionFixture.InitializeAsync(); - var first = await Act(projectionFixture, stream, evt); + BookingDocument? deletedDocument; - var expected = new BookingDocument(id.ToString()) { - RoomId = evt.RoomId, - CheckInDate = evt.CheckIn, - CheckOutDate = evt.CheckOut, - BookingPrice = evt.Price, - Outstanding = evt.Price, - Position = first.Append.GlobalPosition, - StreamPosition = (ulong)first.Append.NextExpectedVersion - }; + try { + var first = await Act(projectionFixture, stream, evt, cancellationToken); - await Assert.That(first.Doc).IsEquivalentTo(expected); + var expected = new BookingDocument(id.ToString()) { + RoomId = evt.RoomId, + CheckInDate = evt.CheckIn, + CheckOutDate = evt.CheckOut, + BookingPrice = evt.Price, + Outstanding = evt.Price, + Position = first.Append.GlobalPosition, + StreamPosition = (ulong)first.Append.NextExpectedVersion + }; - var payment = new BookingPaymentRegistered(Guid.NewGuid().ToString(), evt.Price); + await Assert.That(first.Doc).IsEquivalentTo(expected); - var second = await Act(projectionFixture, stream, payment); + var payment = new BookingPaymentRegistered(Guid.NewGuid().ToString(), evt.Price); - expected = expected with { - PaidAmount = payment.AmountPaid, - Position = second.Append.GlobalPosition, - StreamPosition = (ulong)second.Append.NextExpectedVersion - }; + var second = await Act(projectionFixture, stream, payment, cancellationToken); - await Assert.That(second.Doc).IsEquivalentTo(expected); + expected = expected with { + PaidAmount = payment.AmountPaid, + Position = second.Append.GlobalPosition, + StreamPosition = (ulong)second.Append.NextExpectedVersion + }; - var cancellation = new BookingCancelled(); + await Assert.That(second.Doc).IsEquivalentTo(expected); - var third = await Act(projectionFixture, stream, cancellation); + var cancellation = new BookingCancelled(); - await projectionFixture.DisposeAsync(); + var third = await Act(projectionFixture, stream, cancellation, cancellationToken); + deletedDocument = third.Doc; + } finally { + await projectionFixture.DisposeAsync().NoContext(); + } - await Assert.That(third.Doc).IsNull(); + await Assert.That(deletedDocument).IsNull(); // Extra test to make sure that generated context conversions were used await Assert.That(MessageConsumeContextConverter.ConversionCache).IsEmpty(); } - static async Task<(AppendEventsResult Append, BookingDocument? Doc)> Act(ProjectionTestBase f, StreamName stream, T evt) where T : class { + static async Task<(AppendEventsResult Append, BookingDocument? Doc)> Act(ProjectionTestBase f, StreamName stream, T evt, CancellationToken cancellationToken) where T : class { var append = await f.Fixture.AppendEvent(stream, evt); - await f.WaitForPosition(append.GlobalPosition); - var actual = await f.Fixture.Mongo.LoadDocument(stream.GetId()); + await f.WaitForPosition(append.GlobalPosition, cancellationToken); + var actual = await f.Fixture.Mongo.LoadDocument(stream.GetId(), cancellationToken: cancellationToken); return (append, actual); } diff --git a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBulkBuilder.cs b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBulkBuilder.cs index 9636f590c..eb822bd2e 100644 --- a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBulkBuilder.cs +++ b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectWithBulkBuilder.cs @@ -2,6 +2,7 @@ using Eventuous.Projections.MongoDB.Tools; using Eventuous.Sut.Domain; using Eventuous.Tests.Projections.MongoDB.Fixtures; +using Eventuous.Tools; using MongoDB.Driver; using static Eventuous.Sut.Domain.BookingEvents; @@ -10,45 +11,52 @@ namespace Eventuous.Tests.Projections.MongoDB; [ClassDataSource] public class ProjectWithBulkBuilder(IntegrationFixture fixture) : ProjectionTestBase(nameof(ProjectWithBulkBuilder), fixture) { [Test] - public async Task ShouldProjectImported() { + public async Task ShouldProjectImported(CancellationToken cancellationToken) { await InitializeAsync(); - var evt = DomainFixture.CreateImportBookingEvent(); - var id = new BookingId(CreateId()); - var stream = StreamNameFactory.For(id); - var first = await Act(stream, evt); + BookingDocument expected; + (AppendEventsResult Append, BookingDocument? Doc) second; - var expected = new BookingDocument(id.ToString()) { - RoomId = evt.RoomId, - CheckInDate = evt.CheckIn, - CheckOutDate = evt.CheckOut, - BookingPrice = evt.Price, - Outstanding = evt.Price, - Position = first.Append.GlobalPosition, - StreamPosition = (ulong)first.Append.NextExpectedVersion - }; + try { + var evt = DomainFixture.CreateImportBookingEvent(); + var id = new BookingId(CreateId()); + var stream = StreamNameFactory.For(id); - await Assert.That(first.Doc).IsEquivalentTo(expected); + var first = await Act(stream, evt, cancellationToken); - var payment = new BookingPaymentRegistered(Guid.NewGuid().ToString(), evt.Price); + expected = new BookingDocument(id.ToString()) { + RoomId = evt.RoomId, + CheckInDate = evt.CheckIn, + CheckOutDate = evt.CheckOut, + BookingPrice = evt.Price, + Outstanding = evt.Price, + Position = first.Append.GlobalPosition, + StreamPosition = (ulong)first.Append.NextExpectedVersion + }; - var second = await Act(stream, payment); - await DisposeAsync(); + await Assert.That(first.Doc).IsEquivalentTo(expected); - expected = expected with { - PaidAmount = payment.AmountPaid, - Position = second.Append.GlobalPosition, - StreamPosition = (ulong)second.Append.NextExpectedVersion - }; + var payment = new BookingPaymentRegistered(Guid.NewGuid().ToString(), evt.Price); + + second = await Act(stream, payment, cancellationToken); + + expected = expected with { + PaidAmount = payment.AmountPaid, + Position = second.Append.GlobalPosition, + StreamPosition = (ulong)second.Append.NextExpectedVersion + }; + } finally { + await DisposeAsync().NoContext(); + } await Assert.That(second.Doc).IsEquivalentTo(expected); } - async Task<(AppendEventsResult Append, BookingDocument? Doc)> Act(StreamName stream, T evt) + async Task<(AppendEventsResult Append, BookingDocument? Doc)> Act(StreamName stream, T evt, CancellationToken cancellationToken) where T : class { var append = await Fixture.AppendEvent(stream, evt); - await WaitForPosition(append.GlobalPosition); - var actual = await Fixture.Mongo.LoadDocument(stream.GetId()); + await WaitForPosition(append.GlobalPosition, cancellationToken); + var actual = await Fixture.Mongo.LoadDocument(stream.GetId(), cancellationToken: cancellationToken); return (append, actual); } diff --git a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectingWithTypedHandlers.cs b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectingWithTypedHandlers.cs index ed5f8c2b1..1624a91de 100644 --- a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectingWithTypedHandlers.cs +++ b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectingWithTypedHandlers.cs @@ -2,6 +2,7 @@ using Eventuous.Projections.MongoDB.Tools; using Eventuous.Sut.Domain; using Eventuous.Tests.Projections.MongoDB.Fixtures; +using Eventuous.Tools; using MongoDB.Driver; using static Eventuous.Sut.Domain.BookingEvents; @@ -13,28 +14,31 @@ public sealed class ProjectingWithTypedHandlers(IntegrationFixture fixture) [Test] public async Task ShouldProjectImported(CancellationToken cancellationToken) { await InitializeAsync(); - var evt = DomainFixture.CreateImportBookingEvent(); - var id = new BookingId(CreateId()); - var stream = StreamNameFactory.For(id); - var append = await Fixture.AppendEvent(stream, evt); - - await WaitForPosition(append.GlobalPosition); - - var expected = new BookingDocument(id.ToString()) { - RoomId = evt.RoomId, - CheckInDate = evt.CheckIn, - CheckOutDate = evt.CheckOut, - BookingPrice = evt.Price, - Outstanding = evt.Price, - Position = append.GlobalPosition, - StreamPosition = (ulong)append.NextExpectedVersion - }; - - var actual = await Fixture.Mongo.LoadDocument(id.ToString(), cancellationToken: cancellationToken); - await Assert.That(actual).IsEquivalentTo(expected); - - await DisposeAsync(); + try { + var evt = DomainFixture.CreateImportBookingEvent(); + var id = new BookingId(CreateId()); + var stream = StreamNameFactory.For(id); + + var append = await Fixture.AppendEvent(stream, evt); + + await WaitForPosition(append.GlobalPosition, cancellationToken); + + var expected = new BookingDocument(id.ToString()) { + RoomId = evt.RoomId, + CheckInDate = evt.CheckIn, + CheckOutDate = evt.CheckOut, + BookingPrice = evt.Price, + Outstanding = evt.Price, + Position = append.GlobalPosition, + StreamPosition = (ulong)append.NextExpectedVersion + }; + + var actual = await Fixture.Mongo.LoadDocument(id.ToString(), cancellationToken: cancellationToken); + await Assert.That(actual).IsEquivalentTo(expected); + } finally { + await DisposeAsync().NoContext(); + } } public class SutProjection : MongoProjector { diff --git a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectionTestBase.cs b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectionTestBase.cs index 2324d4251..bacb9058e 100644 --- a/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectionTestBase.cs +++ b/src/Mongo/test/Eventuous.Tests.Projections.MongoDB/ProjectionTestBase.cs @@ -1,30 +1,33 @@ using Eventuous.KurrentDB.Subscriptions; using Eventuous.Projections.MongoDB; using Eventuous.Subscriptions; -using Eventuous.Subscriptions.Checkpoints; using Eventuous.TestHelpers.TUnit.Logging; using Eventuous.Tests.Projections.MongoDB.Fixtures; +using Eventuous.Tools; using Microsoft.Extensions.DependencyInjection; using Microsoft.Extensions.Hosting; +using Microsoft.Extensions.Options; +using MongoDB.Driver; +using MongoDB.Driver.Linq; using static Microsoft.Extensions.Hosting.Host; namespace Eventuous.Tests.Projections.MongoDB; public abstract class ProjectionTestBase { - readonly string _subscriptionId; + protected string SubscriptionId { get; } readonly IHostBuilder _builder; protected IHost Host = null!; protected ProjectionTestBase(string subscriptionId) { - _subscriptionId = subscriptionId; + SubscriptionId = subscriptionId; _builder = CreateDefaultBuilder().ConfigureLogging(cfg => cfg.ForTests()); } protected abstract void ConfigureServices(IServiceCollection services, string subscriptionId); public async Task InitializeAsync() { - _builder.ConfigureServices(collection => ConfigureServices(collection, _subscriptionId)); + _builder.ConfigureServices(collection => ConfigureServices(collection, SubscriptionId)); Host = _builder.Build(); Host.Services.AddEventuousLogs(); await Host.StartAsync(); @@ -49,16 +52,34 @@ protected override void ConfigureServices(IServiceCollection services, string su public string CreateId() => new(Guid.NewGuid().ToString("N")); - public async Task WaitForPosition(ulong position) { - var checkpointStore = Host.Services.GetRequiredService(); - var count = 100; + public async Task WaitForPosition(ulong position, CancellationToken cancellationToken) { + var options = Host.Services.GetRequiredService>().Value; + var checkpoints = Fixture.Mongo.GetCollection(options.CollectionName); + using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + cts.CancelAfter(TimeSpan.FromSeconds(30)); - while (count-- > 0) { - var checkpoint = await checkpointStore.GetLastCheckpoint(nameof(ProjectWithBuilder), default); + ulong? observedPosition = null; - if (checkpoint.Position.HasValue && checkpoint.Position.Value >= position) break; + try { + while (true) { + // GetLastCheckpoint initializes the store's write subject. Poll storage without replacing the active writer. + var checkpoint = await checkpoints.AsQueryable() + .Where(x => x.Id == SubscriptionId) + .SingleOrDefaultAsync(cts.Token) + .NoContext(); + observedPosition = checkpoint?.Position; - await Task.Delay(100); + if (observedPosition.HasValue && observedPosition.Value >= position) return; + + await Task.Delay(100, cts.Token).NoContext(); + } + } catch (OperationCanceledException ex) when (cts.IsCancellationRequested && !cancellationToken.IsCancellationRequested) { + throw new TimeoutException( + $"Expected subscription '{SubscriptionId}' to reach checkpoint {position} within 30 seconds, but observed {observedPosition?.ToString() ?? "no checkpoint"}.", + ex + ); } } + + record StoredCheckpoint(string Id, ulong? Position); } diff --git a/src/Postgres/test/Eventuous.Tests.Postgres/Eventuous.Tests.Postgres.csproj b/src/Postgres/test/Eventuous.Tests.Postgres/Eventuous.Tests.Postgres.csproj index 63fcb946b..0763cbe9c 100644 --- a/src/Postgres/test/Eventuous.Tests.Postgres/Eventuous.Tests.Postgres.csproj +++ b/src/Postgres/test/Eventuous.Tests.Postgres/Eventuous.Tests.Postgres.csproj @@ -16,7 +16,7 @@ - + diff --git a/src/Postgres/test/Eventuous.Tests.Postgres/Projections/ProjectorTests.cs b/src/Postgres/test/Eventuous.Tests.Postgres/Projections/ProjectorTests.cs index 860eeb013..345d64672 100644 --- a/src/Postgres/test/Eventuous.Tests.Postgres/Projections/ProjectorTests.cs +++ b/src/Postgres/test/Eventuous.Tests.Postgres/Projections/ProjectorTests.cs @@ -1,3 +1,5 @@ +extern alias Postgresql; + using Eventuous.Postgresql; using Eventuous.Postgresql.Projections; using Eventuous.Postgresql.Subscriptions; @@ -6,6 +8,7 @@ using Eventuous.Tests.Persistence.Base.Fixtures; using Eventuous.Tests.Postgres.Subscriptions; using Npgsql; +using static Postgresql::Eventuous.Tools.TaskExtensions; using Assert = TUnit.Assertions.Assert; namespace Eventuous.Tests.Postgres.Projections; @@ -26,7 +29,7 @@ public async Task ProjectImportedBookingsToTable(CancellationToken cancellationT await CreateSchema(); var commands = await GenerateAndProduceEvents(100); - await Task.Delay(1000, cancellationToken); + await WaitForBookings(commands.Count, cancellationToken).NoContext(); await using var connection = await _fixture.DataSource.OpenConnectionAsync(cancellationToken); @@ -36,12 +39,33 @@ public async Task ProjectImportedBookingsToTable(CancellationToken cancellationT await using var cmd = new NpgsqlCommand(select, connection); cmd.Parameters.AddWithValue("@bookingId", command.BookingId); await using var reader = await cmd.ExecuteReaderAsync(cancellationToken); - await reader.ReadAsync(cancellationToken); + await Assert.That(await reader.ReadAsync(cancellationToken).NoContext()).IsTrue(); await Assert.That(reader["checkin_date"]).IsEqualTo(command.CheckIn.ToDateTimeUnspecified()); await Assert.That(reader["price"]).IsEqualTo((decimal)command.Price); } } + async Task WaitForBookings(int expectedCount, CancellationToken cancellationToken) { + using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + cts.CancelAfter(TimeSpan.FromSeconds(30)); + + long projectedCount = 0; + + try { + await using var connection = await _fixture.DataSource.OpenConnectionAsync(cts.Token).NoContext(); + await using var cmd = new NpgsqlCommand($"select count(*) from {_fixture.SchemaName}.bookings", connection); + + while (true) { + projectedCount = (long)(await cmd.ExecuteScalarAsync(cts.Token).NoContext())!; + if (projectedCount == expectedCount) return; + + await Task.Delay(100, cts.Token).NoContext(); + } + } catch (OperationCanceledException ex) when (cts.IsCancellationRequested && !cancellationToken.IsCancellationRequested) { + throw new TimeoutException($"Expected {expectedCount} projected bookings within 30 seconds, but observed {projectedCount}.", ex); + } + } + async Task CreateSchema() { await using var connection = await _fixture.DataSource.OpenConnectionAsync(); diff --git a/src/RabbitMq/test/Eventuous.Tests.RabbitMq/HandlerFailureSpec.cs b/src/RabbitMq/test/Eventuous.Tests.RabbitMq/HandlerFailureSpec.cs index d9f0aa3cf..56c1508a9 100644 --- a/src/RabbitMq/test/Eventuous.Tests.RabbitMq/HandlerFailureSpec.cs +++ b/src/RabbitMq/test/Eventuous.Tests.RabbitMq/HandlerFailureSpec.cs @@ -71,15 +71,19 @@ async Task AssertRecovery(CancellationToken cancellationToken) { cts.CancelAfter(TimeSpan.FromSeconds(30)); try { - while (_handler.Handled.Count < EventCount) await Task.Delay(100, cts.Token); - } catch (OperationCanceledException) when (cts.Token.IsCancellationRequested) { + // Rejection can redeliver the failed message before the supervisor reconnects. + // Wait for both message recovery and the lifecycle callbacks asserted below. + while (_handler.Handled.Count < EventCount || Volatile.Read(ref _dropped) < 1 || Volatile.Read(ref _subscribed) < 2) { + await Task.Delay(100, cts.Token).ConfigureAwait(false); + } + } catch (OperationCanceledException) when (cts.IsCancellationRequested && !cancellationToken.IsCancellationRequested) { // Fall through to the assertions, which say more about what went wrong than a cancellation would. } + cancellationToken.ThrowIfCancellationRequested(); await Assert.That(_handler.HasFailed).IsTrue(); - // The event the handler threw on was never acknowledged, so only a resubscribe can bring it back. - // Its absence is the regression: the subscription keeps its connection and quietly loses the message. + // Verify the failed event was redelivered, not just that subsequent events were handled. await Assert.That(_handler.Handled.Order()).IsEquivalentTo(testEvents.Select(x => x.Number).Order()); await Assert.That(Volatile.Read(ref _dropped)).IsGreaterThanOrEqualTo(1); diff --git a/src/SqlServer/test/Eventuous.Tests.SqlServer/Projections/ProjectorTests.cs b/src/SqlServer/test/Eventuous.Tests.SqlServer/Projections/ProjectorTests.cs index c38d8daa2..d79d00949 100644 --- a/src/SqlServer/test/Eventuous.Tests.SqlServer/Projections/ProjectorTests.cs +++ b/src/SqlServer/test/Eventuous.Tests.SqlServer/Projections/ProjectorTests.cs @@ -5,6 +5,7 @@ using Eventuous.Sut.Domain; using Eventuous.Tests.Persistence.Base.Fixtures; using Eventuous.Tests.SqlServer.Subscriptions; +using Eventuous.Tools; using Microsoft.Data.SqlClient; namespace Eventuous.Tests.SqlServer.Projections; @@ -29,7 +30,7 @@ public async Task ProjectImportedBookingsToTable(CancellationToken cancellationT await CreateSchema(); var commands = await GenerateAndProduceEvents(100); - await Task.Delay(1000, cancellationToken); + await WaitForBookings(commands.Count, cancellationToken).NoContext(); await using var connection = await ConnectionFactory.GetConnection(_fixture.ConnectionString, cancellationToken); @@ -45,12 +46,33 @@ async Task ValidateProjectedObject(SqlConnection conn, Commands.ImportBooking co await using var cmd = new SqlCommand(select, conn); cmd.Parameters.AddWithValue("@BookingId", command.BookingId); await using var reader = await cmd.ExecuteReaderAsync(cancellationToken); - await reader.ReadAsync(cancellationToken); + await Assert.That(await reader.ReadAsync(cancellationToken).NoContext()).IsTrue(); await Assert.That(reader["CheckinDate"]).IsEqualTo(command.CheckIn.ToDateTimeUnspecified()); await Assert.That(reader.GetDecimal(1)).IsEqualTo((decimal)command.Price); } } + async Task WaitForBookings(int expectedCount, CancellationToken cancellationToken) { + using var cts = CancellationTokenSource.CreateLinkedTokenSource(cancellationToken); + cts.CancelAfter(TimeSpan.FromSeconds(30)); + + var projectedCount = 0; + + try { + await using var connection = await ConnectionFactory.GetConnection(_fixture.ConnectionString, cts.Token).NoContext(); + await using var cmd = new SqlCommand($"SELECT COUNT(*) FROM {_fixture.SchemaName}.Bookings", connection); + + while (true) { + projectedCount = (int)(await cmd.ExecuteScalarAsync(cts.Token).NoContext())!; + if (projectedCount == expectedCount) return; + + await Task.Delay(100, cts.Token).NoContext(); + } + } catch (OperationCanceledException ex) when (cts.IsCancellationRequested && !cancellationToken.IsCancellationRequested) { + throw new TimeoutException($"Expected {expectedCount} projected bookings within 30 seconds, but observed {projectedCount}.", ex); + } + } + async Task CreateSchema() { await using var connection = await ConnectionFactory.GetConnection(_fixture.ConnectionString, default);