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");