Skip to content
Merged
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
@@ -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<FakeOperations, FakeSession> theStore =
Substitute.For<IEventStore<FakeOperations, FakeSession>>();

private readonly IEventDatabase theDatabase = Substitute.For<IEventDatabase>();
private readonly IEventLoader theLoader = Substitute.For<IEventLoader>();
private readonly ISubscriptionAgent theAgent = Substitute.For<ISubscriptionAgent>();
private readonly FakeProgressionTable theProgression = new();
private readonly ShardName theShardName = new("Composite", ShardName.All, 1);
private readonly CompositeExecution<FakeOperations, FakeSession> theExecution;

public CompositeReplayExecutorProgressionTests()
{
var batch = Substitute.For<IProjectionBatch<FakeOperations, FakeSession>>();
batch.ExecuteAsync(Arg.Any<CancellationToken>()).Returns(Task.CompletedTask);
batch.RecordProgress(Arg.Any<EventRange>()).Returns(call =>
{
theProgression.Record(call.Arg<EventRange>());
return ValueTask.CompletedTask;
});

// Like Marten, starting a batch records the composite's own progress for the range
theStore.StartProjectionBatchAsync(Arg.Any<EventRange>(), Arg.Any<IEventDatabase>(),
Arg.Any<ShardExecutionMode>(), Arg.Any<AsyncOptions>(), Arg.Any<CancellationToken>())
.Returns(call =>
{
theProgression.Record(call.Arg<EventRange>());
return new ValueTask<IProjectionBatch<FakeOperations, FakeSession>>(batch);
});
theStore.ErrorHandlingOptions(Arg.Any<ShardExecutionMode>())
.Returns(new ErrorHandlingOptions { SkipApplyErrors = false });

var member = Substitute.For<ISubscriptionExecution>();
member.ShardName.Returns(new ShardName("Member", ShardName.All, 1));
member.CompactCachesAsync().Returns(Task.CompletedTask);

theExecution = new CompositeExecution<FakeOperations, FakeSession>(theShardName, new AsyncOptions(),
theStore, theDatabase, Substitute.For<IJasperFxProjection<FakeOperations>>(), NullLogger.Instance,
[new ExecutionStage([member])]);

theAgent.Name.Returns(theShardName);
theAgent.Metrics.Returns(Substitute.For<ISubscriptionMetrics>());
theAgent.Status.Returns(AgentStatus.Running);
}

[Fact]
public async Task replay_over_only_filtered_out_events_persists_progression_at_the_ceiling()
{
theLoader.LoadAsync(Arg.Any<EventRequest>(), Arg.Any<CancellationToken>())
.Returns(call => pageOf(call.Arg<EventRequest>()));

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<EventRequest>(), Arg.Any<CancellationToken>())
.Returns(call =>
{
var request = call.Arg<EventRequest>();
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<IDaemonRuntime>()) { 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<Exception>());
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<string, long> _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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<CancellationToken>()).Returns(Task.FromResult(25L));
Expand All @@ -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<EventRequest>(), Arg.Any<CancellationToken>());
await theExecution.DidNotReceive().ProcessRangeAsync(Arg.Any<EventRange>());
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<EventRange>(r =>
r.SequenceFloor == 0 && r.SequenceCeiling == 25 && r.Events.Count == 0));
}

[Fact]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

Expand Down
Loading