diff --git a/src/Core/src/Eventuous.Persistence/EventStore/IEventReader.cs b/src/Core/src/Eventuous.Persistence/EventStore/IEventReader.cs
index 8268287b9..6e69d92c0 100644
--- a/src/Core/src/Eventuous.Persistence/EventStore/IEventReader.cs
+++ b/src/Core/src/Eventuous.Persistence/EventStore/IEventReader.cs
@@ -7,6 +7,12 @@ public interface IEventReader {
///
/// Read a fixed number of events from an existing stream as an async enumerable.
/// Throws if the stream does not exist.
+ /// Implementations either stream events as they arrive from the store, or buffer events in an amount
+ /// proportional to before yielding, so memory usage can grow with
+ /// . To read a whole stream, use ,
+ /// which reads in pages, instead of passing as the count.
+ /// Implementations must yield exactly events unless the end of the stream is reached,
+ /// and must return an empty sequence, not throw, when reading past the end of an existing stream.
///
/// Stream name
/// Where to start reading events
@@ -18,6 +24,9 @@ public interface IEventReader {
///
/// Read a number of events from a given stream, backwards (from the stream end).
/// Throws if the stream does not exist.
+ /// Implementations either stream events as they arrive from the store, or buffer events in an amount
+ /// proportional to before yielding, so memory usage can grow with
+ /// .
///
/// Stream name
/// Where to start reading events
diff --git a/src/Core/src/Eventuous.Persistence/EventStore/StoreFunctions.cs b/src/Core/src/Eventuous.Persistence/EventStore/StoreFunctions.cs
index ca2ee79f2..18a5f430f 100644
--- a/src/Core/src/Eventuous.Persistence/EventStore/StoreFunctions.cs
+++ b/src/Core/src/Eventuous.Persistence/EventStore/StoreFunctions.cs
@@ -1,6 +1,8 @@
// Copyright (C) Eventuous HQ OÜ. All rights reserved
// Licensed under the Apache License, Version 2.0.
+using System.Runtime.CompilerServices;
+
namespace Eventuous;
public static class StoreFunctions {
@@ -148,6 +150,33 @@ CancellationToken cancellationToken
}
}
+ ///
+ /// Reads a stream from the given position to the end, as an async enumerable.
+ /// Events are read in pages of and yielded as they arrive, so the whole stream
+ /// is never buffered in memory. Use this instead of calling
+ /// with as the count.
+ ///
+ /// Name of the stream to read from
+ /// Stream position to start reading from
+ /// Number of events to read per page. It bounds the memory a buffering
+ /// implementation of uses: such implementations hold at most a small
+ /// multiple of a page in memory at a time (e.g. a tiered reader combining two stores).
+ /// Set to false to complete without yielding anything when the stream isn't found,
+ /// instead of throwing . Default is true.
+ /// Cancellation token
+ /// An async enumerable of events retrieved from the stream
+ public IAsyncEnumerable ReadStreamToEnd(
+ StreamName streamName,
+ StreamReadPosition start,
+ int pageSize = 500,
+ bool failIfNotFound = true,
+ CancellationToken cancellationToken = default
+ ) {
+ ArgumentOutOfRangeException.ThrowIfNegativeOrZero(pageSize);
+
+ return ReadToEnd(eventReader, streamName, start, pageSize, failIfNotFound, cancellationToken);
+ }
+
///
/// Reads a stream from the event store to a collection of
///
@@ -163,26 +192,58 @@ public async Task ReadStream(
bool failIfNotFound = true,
CancellationToken cancellationToken = default
) {
- const int pageSize = 500;
-
var streamEvents = new List();
- var position = start;
-
- try {
- while (true) {
- var events = await eventReader.ReadEvents(streamName, position, pageSize, failIfNotFound, cancellationToken).NoContext();
- streamEvents.AddRange(events);
+ await foreach (var evt in eventReader.ReadStreamToEnd(streamName, start, failIfNotFound: failIfNotFound, cancellationToken: cancellationToken).NoContext(cancellationToken)) {
+ streamEvents.Add(evt);
+ }
- if (events.Length < pageSize) break;
+ return [.. streamEvents];
+ }
+ }
- position = new(position.Value + events.Length);
+ // Relies on readers yielding exactly `count` events unless the stream end is reached:
+ // a page shorter than pageSize means there is nothing left to read
+ static async IAsyncEnumerable ReadToEnd(
+ IEventReader eventReader,
+ StreamName streamName,
+ StreamReadPosition start,
+ int pageSize,
+ bool failIfNotFound,
+ [EnumeratorCancellation] CancellationToken cancellationToken
+ ) {
+ var position = start;
+
+ while (true) {
+ var yielded = 0;
+ long lastRevision = 0;
+
+ await using var enumerator = eventReader.ReadEvents(streamName, position, pageSize, cancellationToken).GetAsyncEnumerator(cancellationToken);
+
+ while (true) {
+ bool moved;
+
+ try {
+ moved = await enumerator.MoveNextAsync().NoContext();
+ } catch (StreamNotFound) when (!failIfNotFound) {
+ yield break;
}
- } catch (StreamNotFound) when (!failIfNotFound) {
- return [];
+
+ if (!moved) break;
+
+ var evt = enumerator.Current;
+ yielded++;
+ lastRevision = evt.Revision;
+
+ yield return evt;
}
- return [.. streamEvents];
+ if (yielded < pageSize) yield break;
+
+ // The maximum revision is the end of the representable position space
+ if (lastRevision == long.MaxValue) yield break;
+
+ position = new(lastRevision + 1);
}
}
}
diff --git a/src/Core/src/Eventuous.Persistence/EventStore/TieredEventReader.cs b/src/Core/src/Eventuous.Persistence/EventStore/TieredEventReader.cs
index 3095489b8..08d6ffe3b 100644
--- a/src/Core/src/Eventuous.Persistence/EventStore/TieredEventReader.cs
+++ b/src/Core/src/Eventuous.Persistence/EventStore/TieredEventReader.cs
@@ -13,16 +13,29 @@ namespace Eventuous;
/// Event reader pointing to archive store
public class TieredEventReader(IEventReader hotReader, IEventReader archiveReader) : IEventReader {
public async IAsyncEnumerable ReadEvents(StreamName streamName, StreamReadPosition start, int count, [EnumeratorCancellation] CancellationToken cancellationToken) {
- var hotEvents = await LoadStreamEvents(hotReader, streamName, start, count, cancellationToken).NoContext();
+ var (hotEvents, hotNotFound) = await LoadStreamEvents(hotReader, streamName, start, count, cancellationToken).NoContext();
- var archivedEvents = hotEvents.Length switch {
- > 0 when hotEvents[0].Revision > start.Value
- => (await LoadStreamEvents(archiveReader, streamName, start, (int)hotEvents[0].Revision, cancellationToken).NoContext()).Select(x => x with { FromArchive = true }),
- 0 => (await LoadStreamEvents(archiveReader, streamName, start, count, cancellationToken).NoContext()).Select(x => x with { FromArchive = true }),
- _ => []
- };
+ IEnumerable archivedEvents;
+ var archiveNotFound = false;
+
+ switch (hotEvents.Length) {
+ case > 0 when hotEvents[0].Revision > start.Value: {
+ // Fill the gap before the first hot event from the archive, bounded by the requested count
+ var gapCount = (int)Math.Min(count, hotEvents[0].Revision - start.Value);
+
+ (var events, archiveNotFound) = await LoadStreamEvents(archiveReader, streamName, start, gapCount, cancellationToken).NoContext();
+ archivedEvents = events.Select(x => x with { FromArchive = true });
+
+ break;
+ }
+ case 0:
+ (var archived, archiveNotFound) = await LoadStreamEvents(archiveReader, streamName, start, count, cancellationToken).NoContext();
+ archivedEvents = archived.Select(x => x with { FromArchive = true }); break;
+ default:
+ archivedEvents = []; break;
+ }
- var combined = archivedEvents.Concat(hotEvents).Distinct(Comparer);
+ var combined = archivedEvents.Concat(hotEvents).Distinct(Comparer).Take(count);
var any = false;
foreach (var evt in combined) {
@@ -31,28 +44,32 @@ public async IAsyncEnumerable ReadEvents(StreamName streamName, Str
yield return evt;
}
- if (!any) throw new StreamNotFound(streamName);
+ // No events with both tiers reporting a missing stream means the stream doesn't exist;
+ // otherwise an empty result can mean the read window is past the stream end
+ if (!any && hotNotFound && archiveNotFound) throw new StreamNotFound(streamName);
}
public async IAsyncEnumerable ReadEventsBackwards(StreamName streamName, StreamReadPosition start, int count, [EnumeratorCancellation] CancellationToken cancellationToken) {
- var hotEvents = await LoadStreamEvents(hotReader, streamName, start, count, cancellationToken, backwards: true).NoContext();
+ var (hotEvents, hotNotFound) = await LoadStreamEvents(hotReader, streamName, start, count, cancellationToken, backwards: true).NoContext();
IEnumerable archivedEvents;
+ var archiveNotFound = false;
switch (hotEvents.Length) {
- case > 0 when hotEvents.Length < count: {
+ // When the hot store read reached revision 0, no events can precede it
+ case > 0 when hotEvents.Length < count && hotEvents[^1].Revision > 0: {
// Hot store returned fewer events than requested, fill the gap from archive
var lastHotRevision = hotEvents[^1].Revision;
- archivedEvents = (await LoadStreamEvents(archiveReader, streamName, new(lastHotRevision - 1), count - hotEvents.Length, cancellationToken, backwards: true).NoContext())
- .Select(x => x with { FromArchive = true });
+ (var events, archiveNotFound) = await LoadStreamEvents(archiveReader, streamName, new(lastHotRevision - 1), count - hotEvents.Length, cancellationToken, backwards: true).NoContext();
+ archivedEvents = events.Select(x => x with { FromArchive = true });
break;
}
case 0:
// Hot store has no events, try archive for the full range
- archivedEvents = (await LoadStreamEvents(archiveReader, streamName, start, count, cancellationToken, backwards: true).NoContext())
- .Select(x => x with { FromArchive = true }); break;
+ (var archived, archiveNotFound) = await LoadStreamEvents(archiveReader, streamName, start, count, cancellationToken, backwards: true).NoContext();
+ archivedEvents = archived.Select(x => x with { FromArchive = true }); break;
default:
archivedEvents = []; break;
}
@@ -66,10 +83,12 @@ public async IAsyncEnumerable ReadEventsBackwards(StreamName stream
yield return evt;
}
- if (!any) throw new StreamNotFound(streamName);
+ // No events with both tiers reporting a missing stream means the stream doesn't exist;
+ // otherwise an empty result can mean the read window is past the stream end
+ if (!any && hotNotFound && archiveNotFound) throw new StreamNotFound(streamName);
}
- static async Task LoadStreamEvents(
+ static async Task<(StreamEvent[] Events, bool NotFound)> LoadStreamEvents(
IEventReader reader,
StreamName streamName,
StreamReadPosition startPosition,
@@ -78,11 +97,13 @@ static async Task LoadStreamEvents(
bool backwards = false
) {
try {
- return backwards
+ var events = backwards
? await reader.ReadEventsBackwards(streamName, startPosition, localCount, true, cancellationToken).NoContext()
: await reader.ReadEvents(streamName, startPosition, localCount, true, cancellationToken).NoContext();
+
+ return (events, false);
} catch (StreamNotFound) {
- return [];
+ return ([], true);
}
}
diff --git a/src/Core/test/Eventuous.Tests.Persistence.Base/Store/Read.cs b/src/Core/test/Eventuous.Tests.Persistence.Base/Store/Read.cs
index 25adbbfba..fbec3133c 100644
--- a/src/Core/test/Eventuous.Tests.Persistence.Base/Store/Read.cs
+++ b/src/Core/test/Eventuous.Tests.Persistence.Base/Store/Read.cs
@@ -146,6 +146,117 @@ public async Task ShouldReturnWhenReadingBackwards(CancellationToken cancellatio
await Assert.That(result.Length).IsEqualTo(5);
}
+ [Test]
+ [Category("Store")]
+ public async Task ShouldThrowWhenReadingMissingStream(CancellationToken cancellationToken) {
+ var streamName = Helpers.GetStreamName();
+
+ await Assert.ThrowsAsync(() => _fixture.EventStore.ReadEvents(streamName, StreamReadPosition.Start, 10, true, cancellationToken));
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldThrowWhenReadingMissingStreamBackwards(CancellationToken cancellationToken) {
+ var streamName = Helpers.GetStreamName();
+
+ await Assert.ThrowsAsync(() => _fixture.EventStore.ReadEventsBackwards(streamName, StreamReadPosition.End, 10, true, cancellationToken));
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldReadStreamToEnd(CancellationToken cancellationToken) {
+ object[] events = [.. _fixture.CreateEvents(25)];
+ var streamName = Helpers.GetStreamName();
+ await _fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream);
+
+ var result = new List();
+
+ await foreach (var evt in _fixture.EventStore.ReadStreamToEnd(streamName, StreamReadPosition.Start, pageSize: 10, cancellationToken: cancellationToken)) {
+ result.Add(evt);
+ }
+
+ IEnumerable actual = result.Select(x => x.Payload)!;
+ await Assert.That(actual).IsEquivalentTo(events);
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldReadStreamToEndWithExactPageMultiple(CancellationToken cancellationToken) {
+ object[] events = [.. _fixture.CreateEvents(20)];
+ var streamName = Helpers.GetStreamName();
+ await _fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream);
+
+ var result = new List();
+
+ await foreach (var evt in _fixture.EventStore.ReadStreamToEnd(streamName, StreamReadPosition.Start, pageSize: 10, cancellationToken: cancellationToken)) {
+ result.Add(evt);
+ }
+
+ IEnumerable actual = result.Select(x => x.Payload)!;
+ await Assert.That(actual).IsEquivalentTo(events);
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldReadStreamToEndFromPosition(CancellationToken cancellationToken) {
+ object[] events = [.. _fixture.CreateEvents(25)];
+ var streamName = Helpers.GetStreamName();
+ await _fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream);
+
+ var result = new List();
+
+ await foreach (var evt in _fixture.EventStore.ReadStreamToEnd(streamName, new(10), pageSize: 10, cancellationToken: cancellationToken)) {
+ result.Add(evt);
+ }
+
+ var expected = events.Skip(10);
+ var actual = result.Select(x => x.Payload!);
+ await Assert.That(actual).IsEquivalentTo(expected);
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldRejectInvalidPageSizeReadingToEnd(CancellationToken cancellationToken) {
+ var streamName = Helpers.GetStreamName();
+
+ await Assert.ThrowsAsync(() => Read(0));
+ await Assert.ThrowsAsync(() => Read(-1));
+
+ return;
+
+ async Task Read(int pageSize) {
+ await foreach (var _ in _fixture.EventStore.ReadStreamToEnd(streamName, StreamReadPosition.Start, pageSize: pageSize, cancellationToken: cancellationToken)) { }
+ }
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldThrowWhenReadingMissingStreamToEnd(CancellationToken cancellationToken) {
+ var streamName = Helpers.GetStreamName();
+
+ await Assert.ThrowsAsync(ReadFunc);
+
+ return;
+
+ async Task ReadFunc() {
+ await foreach (var _ in _fixture.EventStore.ReadStreamToEnd(streamName, StreamReadPosition.Start, cancellationToken: cancellationToken)) { }
+ }
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldReturnNothingWhenReadingMissingStreamToEnd(CancellationToken cancellationToken) {
+ var streamName = Helpers.GetStreamName();
+
+ var result = new List();
+
+ await foreach (var evt in _fixture.EventStore.ReadStreamToEnd(streamName, StreamReadPosition.Start, failIfNotFound: false, cancellationToken: cancellationToken)) {
+ result.Add(evt);
+ }
+
+ await Assert.That(result).IsEmpty();
+ }
+
[Test]
[Category("Store")]
public async Task ShouldThrowWhenReadingBackwardsFromNegativePosition(CancellationToken cancellationToken) {
diff --git a/src/Core/test/Eventuous.Tests.Persistence.Base/Store/TieredStoreTests.cs b/src/Core/test/Eventuous.Tests.Persistence.Base/Store/TieredStoreTests.cs
index 9f7ca38f1..8377a49cb 100644
--- a/src/Core/test/Eventuous.Tests.Persistence.Base/Store/TieredStoreTests.cs
+++ b/src/Core/test/Eventuous.Tests.Persistence.Base/Store/TieredStoreTests.cs
@@ -9,6 +9,75 @@ public abstract class TieredStoreTestsBase where TContainer : Docker
protected async Task Should_load_hot_and_archive() {
const int count = 100;
+ var (combined, stream, testEvents) = await SeedTieredStream(count, truncateHotAt: 50);
+
+ var loaded = (await combined.ReadStream(stream, StreamReadPosition.Start)).ToArray();
+
+ var actual = loaded.Select(x => (TestEventForTiers)x.Payload!);
+ await Assert.That(actual).IsEquivalentTo(testEvents);
+
+ await Assert.That(loaded.Take(50).Select(x => x.FromArchive)).DoesNotContain(false);
+ await Assert.That(loaded.Skip(50).Select(x => x.FromArchive)).DoesNotContain(true);
+ }
+
+ protected async Task Should_read_bounded_count_across_tier_boundary() {
+ const int count = 100;
+
+ var (combined, stream, testEvents) = await SeedTieredStream(count, truncateHotAt: 50);
+
+ // The first 50 events only exist in the archive, the hot store starts at revision 50
+ var firstPage = await combined.ReadEvents(stream, StreamReadPosition.Start, 50, true, CancellationToken.None);
+
+ await Assert.That(firstPage.Length).IsEqualTo(50);
+ await Assert.That(firstPage.Select(x => (TestEventForTiers)x.Payload!)).IsEquivalentTo(testEvents.Take(50));
+
+ var loaded = new List();
+
+ await foreach (var evt in combined.ReadStreamToEnd(stream, StreamReadPosition.Start, pageSize: 50)) {
+ loaded.Add(evt);
+ }
+
+ await Assert.That(loaded.Select(x => (TestEventForTiers)x.Payload!)).IsEquivalentTo(testEvents);
+ }
+
+ protected async Task Should_return_empty_reading_past_end() {
+ const int count = 10;
+
+ var (tieredReader, stream, _) = await SeedTieredStream(count);
+
+ var loaded = await tieredReader.ReadEvents(stream, new(count), 5, true, CancellationToken.None);
+
+ await Assert.That(loaded).IsEmpty();
+ }
+
+ protected async Task Should_read_stream_to_end_with_exact_page_multiple() {
+ const int count = 100;
+
+ var (tieredReader, stream, testEvents) = await SeedTieredStream(count);
+
+ var loaded = new List();
+
+ // 100 events with page size 50 forces a final read past the stream end
+ await foreach (var evt in tieredReader.ReadStreamToEnd(stream, StreamReadPosition.Start, pageSize: 50)) {
+ loaded.Add(evt);
+ }
+
+ var actual = loaded.Select(x => (TestEventForTiers)x.Payload!);
+ await Assert.That(actual).IsEquivalentTo(testEvents);
+ }
+
+ protected async Task Should_read_backwards_more_than_available() {
+ const int count = 10;
+
+ var (combined, stream, testEvents) = await SeedTieredStream(count);
+
+ // Requesting more events than the stream holds reads the hot store down to revision 0
+ var loaded = await combined.ReadEventsBackwards(stream, StreamReadPosition.End, count * 2, true, CancellationToken.None);
+
+ await Assert.That(loaded.Select(x => (TestEventForTiers)x.Payload!).Reverse()).IsEquivalentTo(testEvents);
+ }
+
+ async Task<(TieredEventReader Reader, StreamName Stream, TestEventForTiers[] Events)> SeedTieredStream(int count, long? truncateHotAt = null) {
var store = _storeFixture.EventStore;
var archive = new ArchiveStore(_storeFixture.EventStore);
var testEvents = TestEventForTiers.CreateMany(count).ToArray();
@@ -17,15 +86,11 @@ protected async Task Should_load_hot_and_archive() {
await store.Store(stream, ExpectedStreamVersion.NoStream, testEvents);
await archive.Store(stream, ExpectedStreamVersion.NoStream, testEvents);
- await store.TruncateStream(stream, new(50), ExpectedStreamVersion.Any);
- var combined = new TieredEventReader(store, archive);
- var loaded = (await combined.ReadStream(stream, StreamReadPosition.Start)).ToArray();
+ if (truncateHotAt != null) {
+ await store.TruncateStream(stream, new(truncateHotAt.Value), ExpectedStreamVersion.Any);
+ }
- var actual = loaded.Select(x => (TestEventForTiers)x.Payload!);
- await Assert.That(actual).IsEquivalentTo(testEvents);
-
- await Assert.That(loaded.Take(50).Select(x => x.FromArchive)).DoesNotContain(false);
- await Assert.That(loaded.Skip(50).Select(x => x.FromArchive)).DoesNotContain(true);
+ return (new(store, archive), stream, testEvents);
}
readonly StoreFixtureBase _storeFixture;
diff --git a/src/KurrentDB/src/Eventuous.KurrentDB/KurrentDBEventStore.cs b/src/KurrentDB/src/Eventuous.KurrentDB/KurrentDBEventStore.cs
index 165bac57d..b62837830 100644
--- a/src/KurrentDB/src/Eventuous.KurrentDB/KurrentDBEventStore.cs
+++ b/src/KurrentDB/src/Eventuous.KurrentDB/KurrentDBEventStore.cs
@@ -216,48 +216,100 @@ EventData ToEventData(NewStreamEvent streamEvent) {
}
///
- public async IAsyncEnumerable ReadEvents(StreamName stream, StreamReadPosition start, int count, [EnumeratorCancellation] CancellationToken cancellationToken = default) {
- var read = _client.ReadStreamAsync(Direction.Forwards, stream, start.AsStreamPosition(), count, cancellationToken: cancellationToken);
-
- var events = await TryExecute(
- async () => {
- var resolvedEvents = await read.ToArrayAsync(cancellationToken).NoContext();
-
- return ToStreamEvents(resolvedEvents);
- },
+ public IAsyncEnumerable ReadEvents(StreamName stream, StreamReadPosition start, int count, CancellationToken cancellationToken = default)
+ => EnumerateStream(
+ (from, remaining) => _client.ReadStreamAsync(Direction.Forwards, stream, from ?? start.AsStreamPosition(), remaining, cancellationToken: cancellationToken),
+ forwards: true,
stream,
- true,
+ count,
() => new("Unable to read {Count} starting at {Start} events from {Stream}", count, start, stream),
- (s, ex) => new ReadFromStreamException(s, ex)
+ cancellationToken
);
- foreach (var evt in events) yield return evt;
- }
-
///
- public async IAsyncEnumerable ReadEventsBackwards(StreamName stream, StreamReadPosition start, int count, [EnumeratorCancellation] CancellationToken cancellationToken = default) {
- var read = _client.ReadStreamAsync(
- Direction.Backwards,
+ public IAsyncEnumerable ReadEventsBackwards(StreamName stream, StreamReadPosition start, int count, CancellationToken cancellationToken = default)
+ => EnumerateStream(
+ (from, remaining) => _client.ReadStreamAsync(Direction.Backwards, stream, from ?? start.AsStreamPosition(), remaining, resolveLinkTos: true, cancellationToken: cancellationToken),
+ forwards: false,
stream,
- start.AsStreamPosition(),
count,
- resolveLinkTos: true,
- cancellationToken: cancellationToken
+ () => new("Unable to read {Count} events backwards from {Stream}", count, stream),
+ cancellationToken
);
- var events = await TryExecute(
- async () => {
- var resolvedEvents = await read.ToArrayAsync(cancellationToken).NoContext();
+ // Events are yielded as they arrive from the server, so a read holds at most one
+ // deserialized event at a time, regardless of the requested count.
+ // Non-deserializable system events are skipped and compensated for with follow-up
+ // reads, so the enumeration delivers `count` events unless the stream end is reached —
+ // paged readers rely on a short read meaning the end of the stream.
+ // The exception mapping wraps each advance of the source enumerator instead of the whole
+ // loop because iterators can't yield from inside a try block with a catch clause.
+ async IAsyncEnumerable EnumerateStream(
+ Func> read,
+ bool forwards,
+ string stream,
+ int count,
+ Func getError,
+ [EnumeratorCancellation] CancellationToken cancellationToken
+ ) {
+ var remaining = count;
+ StreamPosition? from = null;
+
+ while (remaining > 0) {
+ var requested = remaining;
+ var received = 0;
+ long lastRaw = 0;
+
+ await using var enumerator = read(from, requested).GetAsyncEnumerator(cancellationToken);
+
+ while (true) {
+ var moved = false;
+ StreamEvent? streamEvent = null;
+
+ try {
+ moved = await enumerator.MoveNextAsync().NoContext();
+
+ if (moved) {
+ received++;
+ lastRaw = enumerator.Current.OriginalEventNumber.ToInt64();
+ streamEvent = ToStreamEvent(enumerator.Current);
+ }
+ } catch (StreamNotFoundException) {
+ LogStreamStreamNotFound(stream);
+
+ throw new StreamNotFound(stream);
+ } catch (OperationCanceledException) {
+ throw;
+ } catch (Exception ex) {
+ var (message, args) = getError();
+ // ReSharper disable once TemplateIsNotCompileTimeConstantProblem
+#pragma warning disable CA2254
+ _logger.LogWarning(ex, message, args);
+#pragma warning restore CA2254
- return ToStreamEvents(resolvedEvents);
- },
- stream,
- true,
- () => new("Unable to read {Count} events backwards from {Stream}", count, stream),
- (s, ex) => new ReadFromStreamException(s, ex)
- );
+ throw new ReadFromStreamException(stream, ex);
+ }
+
+ if (!moved) break;
+
+ if (streamEvent != null) {
+ remaining--;
+
+ yield return streamEvent.Value;
+ }
+ }
+
+ // Fewer events received than requested means the stream end was reached
+ if (received < requested) yield break;
+
+ // Nothing was skipped and the requested count is delivered
+ if (remaining == 0) yield break;
- foreach (var evt in events) yield return evt;
+ // Reading backwards can't continue past the first stream event
+ if (!forwards && lastRaw == 0) yield break;
+
+ from = StreamPosition.FromInt64(forwards ? lastRaw + 1 : lastRaw - 1);
+ }
}
///
@@ -362,14 +414,6 @@ StreamEvent AsStreamEvent(object payload)
);
}
- StreamEvent[] ToStreamEvents(ResolvedEvent[] resolvedEvents)
- => [
- .. resolvedEvents
- .Select(ToStreamEvent)
- .Where(x => x != null)
- .Select(x => x!.Value)
- ];
-
record ErrorInfo(string Message, params object[] Args);
[LoggerMessage(LogLevel.Warning, "Stream {stream} not found")]
diff --git a/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Store/StreamingReadTests.cs b/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Store/StreamingReadTests.cs
new file mode 100644
index 000000000..e3abbede2
--- /dev/null
+++ b/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Store/StreamingReadTests.cs
@@ -0,0 +1,139 @@
+using Eventuous.KurrentDB;
+using Eventuous.Sut.Domain;
+using Eventuous.Tests.Persistence.Base.Fixtures;
+using KurrentDB.Client;
+
+namespace Eventuous.Tests.KurrentDB.Store;
+
+[ClassDataSource]
+public class StreamingReadTests {
+ readonly StoreFixture _fixture;
+
+ public StreamingReadTests(StoreFixture fixture) {
+ fixture.TypeMapper.RegisterKnownEventTypes(typeof(BookingEvents.BookingImported).Assembly);
+ _fixture = fixture;
+ }
+
+ const int EventCount = 100;
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldStreamEventsForwardsWithoutBufferingWholeRead(CancellationToken cancellationToken) {
+ var serializer = new CountingSerializer(_fixture.Serializer);
+ var store = new KurrentDBEventStore(_fixture.Client, serializer);
+
+ object[] events = [.. _fixture.CreateEvents(EventCount)];
+ var streamName = Helpers.GetStreamName();
+ await _fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream);
+
+ var deserializedAtFirstYield = 0;
+
+ await foreach (var _ in store.ReadEvents(streamName, StreamReadPosition.Start, EventCount, cancellationToken)) {
+ if (deserializedAtFirstYield == 0) deserializedAtFirstYield = serializer.DeserializedCount;
+ }
+
+ await Assert.That(deserializedAtFirstYield).IsEqualTo(1);
+ await Assert.That(serializer.DeserializedCount).IsEqualTo(EventCount);
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldStreamEventsBackwardsWithoutBufferingWholeRead(CancellationToken cancellationToken) {
+ var serializer = new CountingSerializer(_fixture.Serializer);
+ var store = new KurrentDBEventStore(_fixture.Client, serializer);
+
+ object[] events = [.. _fixture.CreateEvents(EventCount)];
+ var streamName = Helpers.GetStreamName();
+ await _fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream);
+
+ var deserializedAtFirstYield = 0;
+
+ await foreach (var _ in store.ReadEventsBackwards(streamName, new(EventCount - 1), EventCount, cancellationToken)) {
+ if (deserializedAtFirstYield == 0) deserializedAtFirstYield = serializer.DeserializedCount;
+ }
+
+ await Assert.That(deserializedAtFirstYield).IsEqualTo(1);
+ await Assert.That(serializer.DeserializedCount).IsEqualTo(EventCount);
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldReadRequestedCountWhenSystemEventsAreSkipped(CancellationToken cancellationToken) {
+ var (streamName, events) = await SeedStreamWithSystemEvent(cancellationToken);
+
+ var result = new List();
+
+ await foreach (var evt in _fixture.EventStore.ReadEvents(streamName, StreamReadPosition.Start, events.Length, cancellationToken)) {
+ result.Add(evt);
+ }
+
+ IEnumerable actual = result.Select(x => x.Payload)!;
+ await Assert.That(actual).IsEquivalentTo(events);
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldReadRequestedCountBackwardsWhenSystemEventsAreSkipped(CancellationToken cancellationToken) {
+ var (streamName, events) = await SeedStreamWithSystemEvent(cancellationToken);
+
+ var result = new List();
+
+ await foreach (var evt in _fixture.EventStore.ReadEventsBackwards(streamName, StreamReadPosition.End, events.Length, cancellationToken)) {
+ result.Add(evt);
+ }
+
+ IEnumerable actual = result.Select(x => x.Payload).Reverse()!;
+ await Assert.That(actual).IsEquivalentTo(events);
+ }
+
+ [Test]
+ [Category("Store")]
+ public async Task ShouldReadStreamToEndWhenSystemEventsAreSkipped(CancellationToken cancellationToken) {
+ var (streamName, events) = await SeedStreamWithSystemEvent(cancellationToken);
+
+ var result = new List();
+
+ // The system event lands inside the first page, which then yields fewer events than the page size
+ await foreach (var evt in _fixture.EventStore.ReadStreamToEnd(streamName, StreamReadPosition.Start, pageSize: 6, cancellationToken: cancellationToken)) {
+ result.Add(evt);
+ }
+
+ IEnumerable actual = result.Select(x => x.Payload)!;
+ await Assert.That(actual).IsEquivalentTo(events);
+ }
+
+ // Seeds a stream of 12 events where revision 5 is a non-deserializable $-typed event,
+ // which the store skips when reading. Returns the 11 deserializable events.
+ async Task<(StreamName Stream, object[] Events)> SeedStreamWithSystemEvent(CancellationToken cancellationToken) {
+ var streamName = Helpers.GetStreamName();
+ object[] first = [.. _fixture.CreateEvents(5)];
+ object[] rest = [.. _fixture.CreateEvents(6)];
+
+ await _fixture.AppendEvents(streamName, first, ExpectedStreamVersion.NoStream);
+
+ await _fixture.Client.AppendToStreamAsync(
+ streamName.ToString(),
+ StreamState.Any,
+ [new EventData(Uuid.NewUuid(), "$test-skipped", "{}"u8.ToArray())],
+ cancellationToken: cancellationToken
+ );
+
+ await _fixture.AppendEvents(streamName, rest, ExpectedStreamVersion.Any);
+
+ return (streamName, [.. first, .. rest]);
+ }
+
+ class CountingSerializer(IEventSerializer inner) : IEventSerializer {
+ int _deserializedCount;
+
+ public int DeserializedCount => _deserializedCount;
+
+ public DeserializationResult DeserializeEvent(ReadOnlySpan data, string eventType, string contentType) {
+ Interlocked.Increment(ref _deserializedCount);
+
+ return inner.DeserializeEvent(data, eventType, contentType);
+ }
+
+ public SerializationResult SerializeEvent(object evt) => inner.SerializeEvent(evt);
+ }
+}
diff --git a/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Store/TieredStoreTests.cs b/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Store/TieredStoreTests.cs
index e1efd488b..08d8c4ec9 100644
--- a/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Store/TieredStoreTests.cs
+++ b/src/KurrentDB/test/Eventuous.Tests.KurrentDB/Store/TieredStoreTests.cs
@@ -9,4 +9,24 @@ public class TieredStoreTests(StoreFixture storeFixture) : TieredStoreTestsBase<
public async Task Esdb_should_load_hot_and_archive() {
await Should_load_hot_and_archive();
}
+
+ [Test]
+ public async Task Esdb_should_return_empty_reading_past_end() {
+ await Should_return_empty_reading_past_end();
+ }
+
+ [Test]
+ public async Task Esdb_should_read_stream_to_end_with_exact_page_multiple() {
+ await Should_read_stream_to_end_with_exact_page_multiple();
+ }
+
+ [Test]
+ public async Task Esdb_should_read_bounded_count_across_tier_boundary() {
+ await Should_read_bounded_count_across_tier_boundary();
+ }
+
+ [Test]
+ public async Task Esdb_should_read_backwards_more_than_available() {
+ await Should_read_backwards_more_than_available();
+ }
}
\ No newline at end of file
diff --git a/src/Postgres/src/Eventuous.Postgresql/PostgresStore.cs b/src/Postgres/src/Eventuous.Postgresql/PostgresStore.cs
index 4017811c7..25bf8b1db 100644
--- a/src/Postgres/src/Eventuous.Postgresql/PostgresStore.cs
+++ b/src/Postgres/src/Eventuous.Postgresql/PostgresStore.cs
@@ -56,7 +56,8 @@ protected override DbCommand GetReadCommand(NpgsqlConnection connection, StreamN
protected override DbCommand GetReadBackwardsCommand(NpgsqlConnection connection, StreamName stream, StreamReadPosition start, int count)
=> connection.GetCommand(Schema.ReadStreamBackwards)
.Add("_stream_name", NpgsqlDbType.Varchar, stream.ToString())
- .Add("_from_position", NpgsqlDbType.Integer, start.Value)
+ // Stream positions are 32-bit, so StreamReadPosition.End gets clamped, and the function trims it to the stream head
+ .Add("_from_position", NpgsqlDbType.Integer, (int)Math.Min(start.Value, int.MaxValue))
.Add("_count", NpgsqlDbType.Integer, count);
protected override bool IsStreamNotFound(Exception exception)
diff --git a/src/Redis/src/Eventuous.Redis/RedisStore.cs b/src/Redis/src/Eventuous.Redis/RedisStore.cs
index ab7edb3ab..4a667af67 100644
--- a/src/Redis/src/Eventuous.Redis/RedisStore.cs
+++ b/src/Redis/src/Eventuous.Redis/RedisStore.cs
@@ -36,17 +36,52 @@ public RedisStore(
const string ContentType = "application/json";
+ ///
+ /// Reads events from a stream. Positions are inclusive of the start position.
+ /// Streams containing entries written by pre-0.16 versions with auto-generated IDs whose
+ /// sequence number exceeds 9 are not readable from a non-zero position: positions for such
+ /// entries don't round-trip through the position encoding, so resumed reads are rejected with
+ /// instead of risking silently skipped events. Read such
+ /// streams from the start, which fails loudly on the first unrepresentable entry, and migrate them.
+ /// To support that rejection, every read from a non-zero position validates the stream prefix
+ /// below the position in bounded batches, so resumed reads cost extra roundtrips proportional
+ /// to the prefix length.
+ /// Resumed reads require exclusive write ownership of the stream by this store version: any
+ /// writer that doesn't use its explicit entry ID scheme — a pre-0.16 store version, or any
+ /// external XADD with auto-generated IDs — must be quiesced first. Such a writer racing the
+ /// gap between validation and the data read can append an unrepresentable entry below the
+ /// requested position, which that read won't see. The stream is rejected by the next resumed
+ /// read, but the racing read itself can't detect it. Concurrent writers going through this
+ /// store version are safe.
+ ///
public async IAsyncEnumerable ReadEvents(StreamName stream, StreamReadPosition start, int count, [EnumeratorCancellation] CancellationToken cancellationToken) {
StreamEvent[] events;
+ var database = _getDatabase();
try {
- var result = await _getDatabase().StreamReadAsync(stream.ToString(), start.Value.ToRedisValue(), count).NoContext();
+ // A resumed position is only unambiguous when every entry ID in the stream round-trips
+ // through the position encoding (sequence numbers 0-9). Entries the encoding can't
+ // represent can hide below the decoded start position while falling inside the
+ // requested range, so reads from a non-zero position are conservatively rejected for
+ // streams holding any such entry.
+ if (start.Value >= 10) {
+ await EnsureStreamPositionsRoundTrip(database, stream, start, cancellationToken).NoContext();
+ }
+
+ // Range read is inclusive of the start position, matching the IEventReader contract
+ // and the paged read extensions, which advance pages from the last revision + 1
+ var result = await database.StreamRangeAsync(stream.ToString(), start.Value.ToRedisValue(), count: count).NoContext();
if (result == null! || result.Length == 0) {
- throw new StreamNotFound(stream);
+ // An empty result can also mean the read window is past the stream end
+ if (!await database.KeyExistsAsync(stream.ToString()).NoContext()) {
+ throw new StreamNotFound(stream);
+ }
+
+ events = [];
+ } else {
+ events = [.. result.Select(x => ToStreamEvent(x, _serializer, _metaSerializer))];
}
-
- events = [.. result.Select(x => ToStreamEvent(x, _serializer, _metaSerializer))];
} catch (InvalidOperationException e) when (e.Message.Contains("Reading is not allowed after reader was completed") ||
cancellationToken.IsCancellationRequested) {
throw new OperationCanceledException("Redis read operation terminated", e, cancellationToken);
@@ -58,6 +93,53 @@ public async IAsyncEnumerable ReadEvents(StreamName stream, StreamR
public IAsyncEnumerable ReadEventsBackwards(StreamName stream, StreamReadPosition start, int count, CancellationToken cancellationToken)
=> throw new NotImplementedException();
+ const int ValidationPageSize = 1000;
+
+ // Validates that every entry ID below the decoded start position round-trips through the
+ // position encoding, scanning the current stream contents in bounded pages on every call.
+ // No verdict is cached: Redis has no immutable per-key generation identity, so a cached
+ // verdict can go stale when a key is deleted, recreated, or restored under the same name.
+ // Entries at or above the decoded position don't need validation here — the read
+ // materializes them, and converting an unrepresentable ID to a revision fails loudly.
+ // A key replaced concurrently with an in-flight read can still change underneath the scan,
+ // which no non-atomic paged read can detect; that also holds for the data reads themselves.
+ // Likewise, a writer not using the explicit entry ID scheme (a pre-fix store version or any
+ // external XADD with auto-generated IDs) appending an unrepresentable entry between this scan
+ // and the data read escapes the racing read (the next resumed read rejects the stream) — such
+ // writers must be quiesced before resumed reads are used, as documented on ReadEvents.
+ async ValueTask EnsureStreamPositionsRoundTrip(IDatabase database, string stream, StreamReadPosition start, CancellationToken cancellationToken) {
+ RedisValue from = "-";
+ var end = $"({start.Value.ToRedisValue()}";
+
+ while (true) {
+ cancellationToken.ThrowIfCancellationRequested();
+
+ var batch = await database.StreamRangeAsync(stream, from, end, count: ValidationPageSize).NoContext();
+
+ if (batch.Length == 0) break;
+
+ foreach (var entry in batch) {
+ if (EntrySequence(entry.Id) > 9) {
+ throw new NotSupportedException(
+ $"Stream {stream} can't be read from a non-zero position: it contains entry ID {entry.Id}, which the position encoding can't represent (only ID sequence numbers 0-9 are supported). " +
+ "Entries with higher sequence numbers were written with auto-generated IDs by an older version of the store. Read the stream from the start and migrate it."
+ );
+ }
+ }
+
+ if (batch.Length < ValidationPageSize) break;
+
+ from = $"({batch[^1].Id}";
+ }
+ }
+
+ // Redis stream ID sequence components are unsigned 64-bit values
+ static ulong EntrySequence(RedisValue id) {
+ var value = Ensure.NotNull(id);
+
+ return ulong.Parse(value.AsSpan(value.IndexOf('-') + 1));
+ }
+
public async Task AppendEvents(
StreamName stream,
ExpectedStreamVersion expectedVersion,
@@ -132,6 +214,6 @@ static StreamEvent ToStreamEvent(StreamEntry evt, IEventSerializer serializer, I
};
StreamEvent AsStreamEvent(object payload)
- => new(Guid.Parse(evt[MessageId].ToString()), payload, meta ?? new Metadata(), ContentType, evt.Id.ToLong(), DateTime.Parse(evt[Created]!, CultureInfo.InvariantCulture));
+ => new(Guid.Parse(evt[MessageId].ToString()), payload, meta ?? new Metadata(), ContentType, evt.Id.ToRevision(), DateTime.Parse(evt[Created]!, CultureInfo.InvariantCulture));
}
}
diff --git a/src/Redis/src/Eventuous.Redis/Scripts/AppendEvents.lua b/src/Redis/src/Eventuous.Redis/Scripts/AppendEvents.lua
index 5a8831b47..80a2ab475 100644
--- a/src/Redis/src/Eventuous.Redis/Scripts/AppendEvents.lua
+++ b/src/Redis/src/Eventuous.Redis/Scripts/AppendEvents.lua
@@ -1,5 +1,17 @@
#!lua name=append_events
+-- Entry IDs are assigned explicitly as '-0' with the millisecond part bumped past
+-- the last entry when needed. Auto-generated IDs ('*') bump the sequence part instead, and the
+-- client-side position encoding can only represent sequence numbers 0-9.
+local function last_id_ms(key)
+ local entries = redis.call('XREVRANGE', key, '+', '-', 'COUNT', 1)
+ if #entries == 0 then
+ return 0
+ end
+ local id = entries[1][1]
+ return tonumber(string.sub(id, 1, string.find(id, '-', 1, true) - 1))
+end
+
local function append_events(keys, args)
local stream_name = keys[1]
local expected_version = tonumber(keys[2])
@@ -22,20 +34,29 @@ local function append_events(keys, args)
local global_position
local items_inserted = 0
+ local time = redis.call('TIME')
+ local now_ms = tonumber(time[1]) * 1000 + math.floor(tonumber(time[2]) / 1000)
+ local stream_ms = last_id_ms(stream_name)
+ local all_ms = last_id_ms('_all')
+
for i=1, table.getn(args), 4 do
+ stream_ms = math.max(now_ms, stream_ms + 1)
+
local stream_position = redis.call(
- 'XADD', stream_name, '*',
+ 'XADD', stream_name, string.format('%.0f', stream_ms) .. '-0',
'message_id', args[i],
- 'message_type', args[i+1],
- 'json_data', args[i+2],
+ 'message_type', args[i+1],
+ 'json_data', args[i+2],
'json_metadata', args[i+3],
'created', created
)
+ all_ms = math.max(now_ms, all_ms + 1)
+
global_position = redis.call(
- 'XADD', '_all', '*',
- 'stream', stream_name,
+ 'XADD', '_all', string.format('%.0f', all_ms) .. '-0',
+ 'stream', stream_name,
'position', stream_position
)
diff --git a/src/Redis/src/Eventuous.Redis/Tools/Conversions.cs b/src/Redis/src/Eventuous.Redis/Tools/Conversions.cs
index cb0a882ed..f1c20ad9a 100644
--- a/src/Redis/src/Eventuous.Redis/Tools/Conversions.cs
+++ b/src/Redis/src/Eventuous.Redis/Tools/Conversions.cs
@@ -9,6 +9,29 @@ public static long ToLong(this RedisValue value) {
return long.Parse(first) * 10 + long.Parse(second);
}
+ // Redis stream ID components are unsigned 64-bit values, so both parts are parsed as ulong
+ // and range-checked before conversion to the signed position
+ public static long ToRevision(this RedisValue value) {
+ var (first, second) = new Split(Ensure.NotNull(value).AsSpan());
+ var sequence = ulong.Parse(second);
+
+ if (sequence > 9) {
+ throw new NotSupportedException(
+ $"Redis stream entry ID {value} can't be represented as a stream position: the position encoding only supports ID sequence numbers 0-9. " +
+ "Entries with higher sequence numbers were written with auto-generated IDs by an older version of the store."
+ );
+ }
+
+ var milliseconds = ulong.Parse(first);
+
+ const ulong maxMilliseconds = long.MaxValue / 10;
+
+ // At the quotient boundary only sequences up to long.MaxValue % 10 still fit
+ return milliseconds < maxMilliseconds || (milliseconds == maxMilliseconds && sequence <= long.MaxValue % 10)
+ ? (long)milliseconds * 10 + (long)sequence
+ : throw new NotSupportedException($"Redis stream entry ID {value} can't be represented as a stream position: the encoded value exceeds the position range.");
+ }
+
public static ulong ToULong(this ReadOnlySpan valueString) {
var (first, second) = new Split(valueString);
return ulong.Parse(first) * 10 + ulong.Parse(second);
diff --git a/src/Redis/test/Eventuous.Tests.Redis/Store/Read.cs b/src/Redis/test/Eventuous.Tests.Redis/Store/Read.cs
index 8a7649224..4269017a0 100644
--- a/src/Redis/test/Eventuous.Tests.Redis/Store/Read.cs
+++ b/src/Redis/test/Eventuous.Tests.Redis/Store/Read.cs
@@ -1,5 +1,7 @@
+using System.Globalization;
using Eventuous.Tests.Redis.Fixtures;
using Shouldly;
+using StackExchange.Redis;
using static Eventuous.Tests.Redis.Store.Helpers;
namespace Eventuous.Tests.Redis.Store;
@@ -43,12 +45,237 @@ public async Task ShouldReadTail(CancellationToken cancellationToken) {
var events2 = CreateEvents(10).ToArray();
await fixture.AppendEvents(streamName, events2, ExpectedStreamVersion.Any, cancellationToken);
- var result = await fixture.EventReader.ReadEvents(streamName, new((long)position), 100, true, cancellationToken);
+ // The read position is inclusive, so start from the position right after the first batch
+ var result = await fixture.EventReader.ReadEvents(streamName, new((long)position + 1), 100, true, cancellationToken);
IEnumerable actual = result.Select(x => x.Payload)!;
await Assert.That(actual).IsEquivalentTo(events2);
}
+ [Test]
+ public async Task ShouldReadStreamToEndAcrossPages(CancellationToken cancellationToken) {
+ // A single batch this large lands in one millisecond, so all positions must still round-trip
+ var events = CreateEvents(25).ToArray();
+ var streamName = GetStreamName();
+ await fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream, cancellationToken);
+
+ var result = new List();
+
+ await foreach (var evt in fixture.EventReader.ReadStreamToEnd(streamName, StreamReadPosition.Start, pageSize: 4, cancellationToken: cancellationToken)) {
+ result.Add(evt);
+ }
+
+ IEnumerable actual = result.Select(x => x.Payload)!;
+ await Assert.That(actual).IsEquivalentTo(events);
+ }
+
+ [Test]
+ public async Task ShouldRejectLegacyUnrepresentableEntryId(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+
+ // Entries written by older versions can carry auto-generated IDs with sequence numbers
+ // the position encoding can't represent; reading them must fail loudly, not garble positions
+ await AddLegacyEntry(fixture.GetDatabase(), streamName, "12345-10");
+
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, StreamReadPosition.Start, 10, true, cancellationToken));
+ }
+
+ [Test]
+ public async Task ShouldRejectLegacyEntryIdHiddenBehindPageBoundary(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+
+ // Legacy auto-generated IDs from a same-millisecond burst
+ var database = fixture.GetDatabase();
+
+ for (var sequence = 0; sequence <= 10; sequence++) {
+ await AddLegacyEntry(database, streamName, $"12345-{sequence}");
+ }
+
+ // The first page ends at 12345-9 and the advanced position decodes past 12345-10,
+ // which must fail loudly instead of being silently skipped
+ await Assert.ThrowsAsync(ReadFunc);
+
+ return;
+
+ async Task ReadFunc() {
+ await foreach (var _ in fixture.EventReader.ReadStreamToEnd(streamName, StreamReadPosition.Start, pageSize: 10, cancellationToken: cancellationToken)) { }
+ }
+ }
+
+ [Test]
+ public async Task ShouldRejectEntryIdWithSequenceAboveLongRange(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+
+ // Redis ID sequence components are unsigned 64-bit; values beyond long range must still
+ // surface as the documented NotSupportedException, both when materialized and when validated
+ await AddLegacyEntry(fixture.GetDatabase(), streamName, "12345-9223372036854775808");
+
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, StreamReadPosition.Start, 10, true, cancellationToken));
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, new(123470), 10, true, cancellationToken));
+ }
+
+ [Test]
+ public async Task ShouldHandleRevisionBoundaryAtLongMax(CancellationToken cancellationToken) {
+ // long.MaxValue / 10 = 922337203685477580, long.MaxValue % 10 = 7: sequence 7 encodes to
+ // exactly long.MaxValue, sequence 8 no longer fits and must be rejected, not wrap negative
+ var fitting = GetStreamName();
+ await AddLegacyEntry(fixture.GetDatabase(), fitting, "922337203685477580-7");
+
+ var result = await fixture.EventReader.ReadEvents(fitting, StreamReadPosition.Start, 10, true, cancellationToken);
+ await Assert.That(result[0].Revision).IsEqualTo(long.MaxValue);
+
+ var overflowing = GetStreamName();
+ await AddLegacyEntry(fixture.GetDatabase(), overflowing, "922337203685477580-8");
+
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(overflowing, StreamReadPosition.Start, 10, true, cancellationToken));
+ }
+
+ [Test]
+ public async Task ShouldReadStreamToEndAtMaxRevision(CancellationToken cancellationToken) {
+ // An event at the maximum representable revision filling an exact page must complete
+ // the paged read instead of advancing past the end of the position space
+ var streamName = GetStreamName();
+ await AddLegacyEntry(fixture.GetDatabase(), streamName, "922337203685477580-7");
+
+ var result = new List();
+
+ await foreach (var evt in fixture.EventReader.ReadStreamToEnd(streamName, StreamReadPosition.Start, pageSize: 1, cancellationToken: cancellationToken)) {
+ result.Add(evt);
+ }
+
+ await Assert.That(result).HasCount().EqualTo(1);
+ await Assert.That(result[0].Revision).IsEqualTo(long.MaxValue);
+ }
+
+ [Test]
+ public async Task ShouldRejectLegacyBurstStreamReadFromStart(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+
+ // A legacy burst with sequence numbers beyond a single decimal carry: positions minted for
+ // such entries by older versions are ambiguous, but any read from the start of the stream
+ // must reject the first unrepresentable entry it materializes
+ var database = fixture.GetDatabase();
+
+ for (var sequence = 0; sequence <= 20; sequence += 5) {
+ await AddLegacyEntry(database, streamName, $"12345-{sequence}");
+ }
+
+ await Assert.ThrowsAsync(ReadFunc);
+
+ return;
+
+ async Task ReadFunc() {
+ await foreach (var _ in fixture.EventReader.ReadStreamToEnd(streamName, StreamReadPosition.Start, cancellationToken: cancellationToken)) { }
+ }
+ }
+
+ [Test]
+ public async Task ShouldRejectResumedCursorOnLegacyBurstStream(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+
+ var database = fixture.GetDatabase();
+
+ for (var sequence = 0; sequence <= 20; sequence++) {
+ await AddLegacyEntry(database, streamName, $"12345-{sequence}");
+ }
+
+ // A cursor minted by a pre-fix reader after consuming 12345-19 (revision 123469 + 1):
+ // resuming from it must be rejected, not silently skip the remaining entries
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, new(123470), 10, true, cancellationToken));
+ }
+
+ [Test]
+ public async Task ShouldRejectResumedReadAfterStreamRecreatedWithLegacyEntries(CancellationToken cancellationToken) {
+ var events = CreateEvents(3).ToArray();
+ var streamName = GetStreamName();
+ await fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream, cancellationToken);
+
+ // A resumed read on the clean stream passes validation
+ var appended = await fixture.EventReader.ReadEvents(streamName, new(10), 10, true, cancellationToken);
+ await Assert.That(appended.Length).IsGreaterThan(0);
+
+ // Recreate the stream under the same name with legacy entries: the earlier verdict must not stick
+ var database = fixture.GetDatabase();
+ await database.KeyDeleteAsync(streamName.ToString());
+
+ for (var sequence = 0; sequence <= 20; sequence++) {
+ await AddLegacyEntry(database, streamName, $"12345-{sequence}");
+ }
+
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, new(123470), 10, true, cancellationToken));
+ }
+
+ [Test]
+ public async Task ShouldRejectResumedReadAfterStreamRestoredWithSameFirstEntry(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+ var database = fixture.GetDatabase();
+
+ // A clean stream with explicit IDs, validated by a resumed read
+ await AddLegacyEntry(database, streamName, "12345-0");
+ await AddLegacyEntry(database, streamName, "12346-0");
+ await AddLegacyEntry(database, streamName, "12347-0");
+
+ var appended = await fixture.EventReader.ReadEvents(streamName, new(123460), 10, true, cancellationToken);
+ await Assert.That(appended.Length).IsGreaterThan(0);
+
+ // Restore the stream with the same first entry but an unrepresentable entry
+ // below the previously validated range: the earlier verdict must not stick
+ await database.KeyDeleteAsync(streamName.ToString());
+ await AddLegacyEntry(database, streamName, "12345-0");
+ await AddLegacyEntry(database, streamName, "12346-10");
+
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, new(123470), 10, true, cancellationToken));
+ }
+
+ [Test]
+ public async Task ShouldRejectResumedReadAfterMissingStreamGetsLegacyEntries(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+
+ // A resumed read of a missing stream must not establish a verdict for the name
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, new(100), 10, true, cancellationToken));
+
+ var database = fixture.GetDatabase();
+
+ for (var sequence = 0; sequence <= 20; sequence++) {
+ await AddLegacyEntry(database, streamName, $"12345-{sequence}");
+ }
+
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, new(123470), 10, true, cancellationToken));
+ }
+
+ static async Task AddLegacyEntry(IDatabase database, StreamName streamName, string id) {
+ var serialized = EventSerializer.Default.SerializeEvent(CreateEvent());
+
+ await database.StreamAddAsync(
+ streamName.ToString(),
+ [
+ new("message_id", Guid.NewGuid().ToString()),
+ new("message_type", serialized.EventType),
+ new("json_data", serialized.Payload),
+ new("created", DateTime.UtcNow.ToString(CultureInfo.InvariantCulture))
+ ],
+ id
+ );
+ }
+
+ [Test]
+ public async Task ShouldReturnEmptyReadingPastEnd(CancellationToken cancellationToken) {
+ var events = CreateEvents(10).ToArray();
+ var streamName = GetStreamName();
+ var appended = await fixture.AppendEvents(streamName, events, ExpectedStreamVersion.NoStream, cancellationToken);
+
+ var result = await fixture.EventReader.ReadEvents(streamName, new((long)appended.GlobalPosition + 1000), 10, true, cancellationToken);
+
+ await Assert.That(result).IsEmpty();
+ }
+
+ [Test]
+ public async Task ShouldThrowWhenReadingMissingStream(CancellationToken cancellationToken) {
+ var streamName = GetStreamName();
+
+ await Assert.ThrowsAsync(() => fixture.EventReader.ReadEvents(streamName, StreamReadPosition.Start, 10, true, cancellationToken));
+ }
+
[Test]
public async Task ShouldReadHead(CancellationToken cancellationToken) {
// ReSharper disable once CoVariantArrayConversion
diff --git a/src/Relational/src/Eventuous.Sql.Base/SqlEventStoreBase.cs b/src/Relational/src/Eventuous.Sql.Base/SqlEventStoreBase.cs
index 8f2a41d21..8729106d7 100644
--- a/src/Relational/src/Eventuous.Sql.Base/SqlEventStoreBase.cs
+++ b/src/Relational/src/Eventuous.Sql.Base/SqlEventStoreBase.cs
@@ -103,6 +103,9 @@ public async IAsyncEnumerable ReadEvents(StreamName stream, StreamR
var events = await ReadInternal(stream, start, count, cancellationToken).NoContext();
+ // A plain query can't tell a missing stream from a read past the stream end
+ if (events.Length == 0 && !await StreamExists(stream, cancellationToken).NoContext()) throw new StreamNotFound(stream);
+
foreach (var evt in events) yield return evt;
}
@@ -112,6 +115,9 @@ public async IAsyncEnumerable ReadEventsBackwards(StreamName stream
var events = await ReadInternalBackwards(stream, start, count, cancellationToken).NoContext();
+ // A plain query can't tell a missing stream from a read past the stream end
+ if (events.Length == 0 && !await StreamExists(stream, cancellationToken).NoContext()) throw new StreamNotFound(stream);
+
foreach (var evt in events) yield return evt;
}
diff --git a/src/SqlServer/src/Eventuous.SqlServer/SqlServerStore.cs b/src/SqlServer/src/Eventuous.SqlServer/SqlServerStore.cs
index 118953802..fd5ee7484 100644
--- a/src/SqlServer/src/Eventuous.SqlServer/SqlServerStore.cs
+++ b/src/SqlServer/src/Eventuous.SqlServer/SqlServerStore.cs
@@ -41,7 +41,8 @@ protected override DbCommand GetReadBackwardsCommand(SqlConnection connection, S
=> connection
.GetStoredProcCommand(Schema.ReadStreamBackwards)
.Add("@stream_name", SqlDbType.NVarChar, stream.ToString())
- .Add("@from_position", SqlDbType.Int, start.Value)
+ // Stream positions are 32-bit, so StreamReadPosition.End gets clamped, and the procedure trims it to the stream head
+ .Add("@from_position", SqlDbType.Int, (int)Math.Min(start.Value, int.MaxValue))
.Add("@count", SqlDbType.Int, count);
protected override bool IsStreamNotFound(Exception exception) => exception is SqlException e && e.Message.StartsWith("StreamNotFound");