diff --git a/src/EventTests/Projections/Composite/CompositeReplayExecutorProgressionTests.cs b/src/EventTests/Projections/Composite/CompositeReplayExecutorProgressionTests.cs new file mode 100644 index 00000000..c2ee1daf --- /dev/null +++ b/src/EventTests/Projections/Composite/CompositeReplayExecutorProgressionTests.cs @@ -0,0 +1,153 @@ +using JasperFx; +using JasperFx.Events; +using JasperFx.Events.Daemon; +using JasperFx.Events.Projections; +using JasperFx.Events.Projections.Composite; +using Microsoft.Extensions.Logging.Abstractions; +using NSubstitute; +using Shouldly; + +namespace EventTests.Projections.Composite; + +// A single-pass composite replay that ends on an empty page must still persist progression, or the +// first continuous page after it updates a row that is missing or behind and the shard stops. +public class CompositeReplayExecutorProgressionTests +{ + private const long Ceiling = 10; + + private readonly IEventStore theStore = + Substitute.For>(); + + private readonly IEventDatabase theDatabase = Substitute.For(); + private readonly IEventLoader theLoader = Substitute.For(); + private readonly ISubscriptionAgent theAgent = Substitute.For(); + private readonly FakeProgressionTable theProgression = new(); + private readonly ShardName theShardName = new("Composite", ShardName.All, 1); + private readonly CompositeExecution theExecution; + + public CompositeReplayExecutorProgressionTests() + { + var batch = Substitute.For>(); + batch.ExecuteAsync(Arg.Any()).Returns(Task.CompletedTask); + batch.RecordProgress(Arg.Any()).Returns(call => + { + theProgression.Record(call.Arg()); + return ValueTask.CompletedTask; + }); + + // Like Marten, starting a batch records the composite's own progress for the range + theStore.StartProjectionBatchAsync(Arg.Any(), Arg.Any(), + Arg.Any(), Arg.Any(), Arg.Any()) + .Returns(call => + { + theProgression.Record(call.Arg()); + return new ValueTask>(batch); + }); + theStore.ErrorHandlingOptions(Arg.Any()) + .Returns(new ErrorHandlingOptions { SkipApplyErrors = false }); + + var member = Substitute.For(); + member.ShardName.Returns(new ShardName("Member", ShardName.All, 1)); + member.CompactCachesAsync().Returns(Task.CompletedTask); + + theExecution = new CompositeExecution(theShardName, new AsyncOptions(), + theStore, theDatabase, Substitute.For>(), NullLogger.Instance, + [new ExecutionStage([member])]); + + theAgent.Name.Returns(theShardName); + theAgent.Metrics.Returns(Substitute.For()); + theAgent.Status.Returns(AgentStatus.Running); + } + + [Fact] + public async Task replay_over_only_filtered_out_events_persists_progression_at_the_ceiling() + { + theLoader.LoadAsync(Arg.Any(), Arg.Any()) + .Returns(call => pageOf(call.Arg())); + + await replayAsync(batchSize: 500); + + theProgression.PositionOf("Composite:All").ShouldBe(Ceiling); + theProgression.PositionOf("Member:All").ShouldBe(Ceiling); + await continuousPageAfterTheReplayShouldNotFail(); + } + + [Fact] + public async Task replay_ending_on_a_full_page_followed_by_an_empty_page_persists_progression_at_the_ceiling() + { + theLoader.LoadAsync(Arg.Any(), Arg.Any()) + .Returns(call => + { + var request = call.Arg(); + return request.Floor == 0 ? pageOf(request, "ab") : pageOf(request); + }); + + await replayAsync(batchSize: 2); + + theProgression.PositionOf("Composite:All").ShouldBe(Ceiling); + await continuousPageAfterTheReplayShouldNotFail(); + } + + private async Task replayAsync(int batchSize) + { + var executor = new CompositeReplayExecutor(theShardName, theLoader, theExecution, theDatabase, + new AsyncOptions { BatchSize = batchSize }, NullLogger.Instance); + + var request = new SubscriptionExecutionRequest(0, ShardExecutionMode.CatchUp, new ErrorHandlingOptions(), + Substitute.For()) { StartingHighWater = Ceiling }; + + await executor.StartAsync(request, theAgent, CancellationToken.None); + + await theAgent.Received().MarkSuccessAsync(Ceiling); + } + + private async Task continuousPageAfterTheReplayShouldNotFail() + { + theExecution.Mode = ShardExecutionMode.Continuous; + await theExecution.ProcessRangeAsync(new EventRange(theAgent, Ceiling, Ceiling + 5) { Events = [] }); + + await theAgent.DidNotReceive().ReportCriticalFailureAsync(Arg.Any()); + theProgression.PositionOf("Composite:All").ShouldBe(Ceiling + 5); + } + + private static EventPage pageOf(EventRequest request, string letters = "") + { + var page = new EventPage(request.Floor); + long sequence = request.Floor; + foreach (var @event in letters.ToLetterEventsWithWrapper()) + { + @event.Sequence = ++sequence; + page.Add(@event); + } + + page.CalculateCeiling(request.BatchSize, request.HighWater); + return page; + } + + // Mirrors Marten's progression writes: a range from floor 0 inserts the row, any other range + // updates it only where it is still at that floor. + private class FakeProgressionTable + { + private readonly Dictionary _rows = new(); + + public void Record(EventRange range) + { + var name = range.ShardName.Identity; + if (range.SequenceFloor == 0) + { + _rows.Add(name, range.SequenceCeiling); + return; + } + + if (!_rows.TryGetValue(name, out var current) || current != range.SequenceFloor) + { + throw new ProgressionProgressOutOfOrderException(name, range.SequenceFloor, range.SequenceCeiling); + } + + _rows[name] = range.SequenceCeiling; + } + + public long? PositionOf(string identityPrefix) => + _rows.Where(x => x.Key.StartsWith(identityPrefix)).Select(x => (long?)x.Value).SingleOrDefault(); + } +} diff --git a/src/EventTests/Projections/Composite/CompositeReplayExecutorTests.cs b/src/EventTests/Projections/Composite/CompositeReplayExecutorTests.cs index 4af6cc98..8e4eb288 100644 --- a/src/EventTests/Projections/Composite/CompositeReplayExecutorTests.cs +++ b/src/EventTests/Projections/Composite/CompositeReplayExecutorTests.cs @@ -89,7 +89,7 @@ public async Task advances_progression_to_ceiling_when_store_is_empty() } [Fact] - public async Task advances_to_ceiling_when_no_events_match_below_high_water() + public async Task commits_an_empty_range_to_the_ceiling_when_no_events_match_below_high_water() { theAgent.Status.Returns(AgentStatus.Running); theDatabase.FetchHighestEventSequenceNumber(Arg.Any()).Returns(Task.FromResult(25L)); @@ -99,8 +99,10 @@ public async Task advances_to_ceiling_when_no_events_match_below_high_water() await theExecutor().StartAsync(rebuildRequest(), theAgent, CancellationToken.None); await theLoader.Received(1).LoadAsync(Arg.Any(), Arg.Any()); - await theExecution.DidNotReceive().ProcessRangeAsync(Arg.Any()); - await theAgent.Received(1).MarkSuccessAsync(25); + + // The empty range still goes through the execution, which persists progression and marks success + await theExecution.Received(1).ProcessRangeAsync(Arg.Is(r => + r.SequenceFloor == 0 && r.SequenceCeiling == 25 && r.Events.Count == 0)); } [Fact] diff --git a/src/JasperFx.Events/Projections/Composite/CompositeReplayExecutor.cs b/src/JasperFx.Events/Projections/Composite/CompositeReplayExecutor.cs index d22361c2..958b4801 100644 --- a/src/JasperFx.Events/Projections/Composite/CompositeReplayExecutor.cs +++ b/src/JasperFx.Events/Projections/Composite/CompositeReplayExecutor.cs @@ -86,8 +86,11 @@ public async Task StartAsync(SubscriptionExecutionRequest request, ISubscription if (page.Count == 0) { - // No more matching events below the ceiling; advance progression to the ceiling and finish. - await controller.MarkSuccessAsync(ceiling).ConfigureAwait(false); + // No more matching events below the ceiling. Commit the empty range anyway so the stored + // progression moves to the ceiling with the in-memory one; MarkSuccessAsync alone leaves + // the row missing or behind, and the next continuous page then fails to update it. + var emptyRange = new EventRange(agent, page.Floor, ceiling) { Events = page }; + await _execution.ProcessRangeAsync(emptyRange).ConfigureAwait(false); break; }