Skip to content
Open
Show file tree
Hide file tree
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 @@ -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();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,8 @@ namespace Eventuous.Tests.Subscriptions;
/// event in flight when a run's own teardown cancelled a parked handler.
/// </summary>
public class CancelledMessageTests {
static readonly TimeSpan RecoveryTimeout = TimeSpan.FromSeconds(30);

/// <summary>
/// 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.
Expand All @@ -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");
Expand Down Expand Up @@ -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);

Expand All @@ -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;

Expand All @@ -123,6 +130,10 @@ sealed class SingleEventSubscription(
null
) {
SubscriptionRun? _run;
int _deliveriesQueued;
readonly TaskCompletionSource _redeliveryQueued = new(TaskCreationOptions.RunContinuationsAsynchronously);

public Task RedeliveryQueued => _redeliveryQueued.Task;

/// <summary>
/// Fails the current run, standing in for a transport drop or any other reason the supervisor tears a run down.
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -199,7 +212,7 @@ public override async ValueTask<EventHandlingStatus> HandleEvent(IMessageConsume
break;
case 2:
_redeliveryStarted.TrySetResult();
await _proceedWithSuccess.Task.NoContext();
await _proceedWithSuccess.Task.WaitAsync(context.CancellationToken).NoContext();
break;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -14,17 +14,17 @@ namespace Eventuous.Tests.Subscriptions;
/// </summary>
public class ResubscribeConcurrencyTests {
/// <summary>
/// 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.
/// </summary>
[Test]
public async Task Burst_of_nacks_produces_a_single_resubscribe(CancellationToken ct) {
var logs = new CapturingLoggerFactory(LogLevel.Trace);

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,
Expand All @@ -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) {
Expand All @@ -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();
}
}

/// <summary>
Expand Down Expand Up @@ -805,6 +812,31 @@ void TrackPumpStarted() {
}
}

/// <summary>
/// Holds all deliveries until the test releases the burst of failures.
/// </summary>
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<EventHandlingStatus> 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}");
}
}

/// <summary>
/// Defers by throwing when a precondition isn't met, so the message is redelivered after the resubscribe.
/// </summary>
Expand Down
3 changes: 2 additions & 1 deletion src/Core/test/Eventuous.Tests.Subscriptions/SequenceTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -62,7 +62,8 @@ public void ShouldWorkForNormalCase() {
}

public static IEnumerable<Func<(CommitPositionSequence, CommitPosition)>> 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));
Expand Down
19 changes: 14 additions & 5 deletions src/Diagnostics/test/Eventuous.Tests.OpenTelemetry/MetricsTests.cs
Original file line number Diff line number Diff line change
Expand Up @@ -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();

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -14,56 +15,61 @@ public class ProjectWithBuilder(IntegrationFixture fixture) {
[Test]
[Retry(3)]
[MethodDataSource(typeof(CollectionSource), nameof(CollectionSource.TestOptions))]
public async Task ShouldProjectImported(MongoProjectionOptions<BookingDocument>? options) {
public async Task ShouldProjectImported(MongoProjectionOptions<BookingDocument>? options, CancellationToken cancellationToken) {
var evt = DomainFixture.CreateImportBookingEvent();
var projectionFixture = new ProjectionTestBase<SutProjection>(nameof(ProjectWithBuilder), fixture);
var id = new BookingId(projectionFixture.CreateId());
var stream = StreamNameFactory.For<Booking, BookingState, BookingId>(id);

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<T>(ProjectionTestBase<SutProjection> f, StreamName stream, T evt) where T : class {
static async Task<(AppendEventsResult Append, BookingDocument? Doc)> Act<T>(ProjectionTestBase<SutProjection> 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<BookingDocument>(stream.GetId());
await f.WaitForPosition(append.GlobalPosition, cancellationToken);
var actual = await f.Fixture.Mongo.LoadDocument<BookingDocument>(stream.GetId(), cancellationToken: cancellationToken);

return (append, actual);
}
Expand Down
Loading
Loading