From 065a7a0ef32ddc61e7ddc33e64e0e3d8a725af55 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Thu, 24 Sep 2026 01:37:57 -0400 Subject: [PATCH 1/3] fix(grpc): preserve expected-version failures with unknown revisions Signed-off-by: Yordis Prieto --- .../AppendAcrossRestartGrpcTests.cs | 78 +++++++++ .../StreamsTests/GrpcStreamEdgeOperations.cs | 148 ++++++++++++++++++ .../HashCollisionGrpcBoundaryTests.cs | 90 +++++++++++ .../StreamRevisionAboveIntMaxTests.cs | 105 +++++++++++++ .../Services/Transport/Grpc/Status.cs | 5 +- .../Services/Transport/Grpc/Streams.Append.cs | 2 +- .../Transport/Grpc/Streams.BatchAppend.cs | 2 +- .../Transport/Grpc/WrongExpectedVersion.cs | 26 ++- 8 files changed, 443 insertions(+), 13 deletions(-) create mode 100644 src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/AppendAcrossRestartGrpcTests.cs create mode 100644 src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs create mode 100644 src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs create mode 100644 src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/AppendAcrossRestartGrpcTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/AppendAcrossRestartGrpcTests.cs new file mode 100644 index 0000000000..062bd1858c --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/AppendAcrossRestartGrpcTests.cs @@ -0,0 +1,78 @@ +using System; +using System.IO; +using System.Linq; +using System.Threading.Tasks; +using EventStore.Client.Streams; +using EventStore.Core.Tests.Helpers; +using NUnit.Framework; + +namespace EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests; + +[Category("LongRunning")] +public class AppendAcrossRestartGrpcTests : SpecificationWithDirectoryPerTestFixture +{ + private MiniNode _node; + private GrpcStreamEdgeOperations _grpc; + private string _dbPath; + + [OneTimeSetUp] + public override async Task TestFixtureSetUp() + { + await base.TestFixtureSetUp(); + _dbPath = Path.Combine(PathName, "restart-node-db"); + await StartNode(waitForAdminUserCreation: true); + } + + [OneTimeTearDown] + public override async Task TestFixtureTearDown() + { + _grpc?.Dispose(); + if (_node is not null) + await _node.Shutdown(); + await base.TestFixtureTearDown(); + } + + [Test] + public async Task detects_existing_streams_and_metadata_after_restart() + { + const string stream = "grpc-existing-stream-across-restart"; + const string metadataStream = "$$grpc-metadata-across-restart"; + AssertSuccess(await _grpc.Append(stream, count: 10, noStream: true), 9); + AssertSuccess(await _grpc.Append(metadataStream, noStream: true, + data: "{\"$maxCount\":5}", eventType: "$metadata"), 0); + AssertSuccess(await _grpc.Append("grpc-last-stream-before-restart", noStream: true), 0); + + await Task.Delay(500); + await _node.Shutdown(keepDb: true); + _grpc.Dispose(); + await StartNode(waitForAdminUserCreation: false); + + AssertSuccess(await _grpc.Append(stream, expectedRevision: 9), 10); + AssertSuccess(await _grpc.Append(metadataStream, expectedRevision: 0, + data: "{\"$maxCount\":6}", eventType: "$metadata"), 1); + + var events = await _grpc.Read(stream, 0, 20); + Assert.That(events.Count(x => x.Event is not null), Is.EqualTo(11)); + Assert.That(events.Last(x => x.Event is not null).Event.Event.StreamRevision, + Is.EqualTo(10)); + } + + private async Task StartNode(bool waitForAdminUserCreation) + { + _node = new MiniNode(PathName, + dbPath: _dbPath, + streamExistenceFilterSize: 10_000, + streamExistenceFilterCheckpointIntervalMs: 100, + streamExistenceFilterCheckpointDelayMs: 0); + await _node.Start(); + if (waitForAdminUserCreation) + await _node.AdminUserCreated; + _grpc = new GrpcStreamEdgeOperations(_node); + } + + private static void AssertSuccess(BatchAppendResp response, ulong expectedRevision) + { + Assert.That(response.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Success)); + Assert.That(response.Success.CurrentRevision, Is.EqualTo(expectedRevision)); + } +} diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs new file mode 100644 index 0000000000..452ace9dcf --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs @@ -0,0 +1,148 @@ +using System; +using System.Linq; +using System.Text; +using System.Threading.Tasks; +using EventStore.Client.Streams; +using EventStore.Core.Services.Transport.Grpc; +using EventStore.Core.Tests.Helpers; +using Google.Protobuf; +using Grpc.Core; +using Grpc.Net.Client; +using NUnit.Framework; +using GrpcMetadata = EventStore.Core.Services.Transport.Grpc.Constants.Metadata; +using Streams = EventStore.Client.Streams.Streams; + +namespace EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests; + +internal sealed class GrpcStreamEdgeOperations : IDisposable +{ + private readonly GrpcChannel _channel; + private readonly Streams.StreamsClient _client; + private readonly CallCredentials _credentials; + private CallOptions CallOptions => new(credentials: _credentials, + deadline: DateTime.UtcNow.AddSeconds(20)); + + public GrpcStreamEdgeOperations(MiniNode node) + { + _channel = GrpcChannel.ForAddress(new UriBuilder { Scheme = Uri.UriSchemeHttps }.Uri, + new GrpcChannelOptions { HttpClient = node.HttpClient, DisposeHttpClient = false }); + _client = new Streams.StreamsClient(_channel); + _credentials = CallCredentials.FromInterceptor((_, metadata) => + { + metadata.Add("authorization", "Basic " + Convert.ToBase64String( + Encoding.ASCII.GetBytes("admin:changeit"))); + return Task.CompletedTask; + }); + } + + public GrpcStreamEdgeOperations(GrpcChannel channel) + { + _client = new Streams.StreamsClient(channel); + _credentials = CallCredentials.FromInterceptor((_, metadata) => + { + metadata.Add("authorization", "Basic " + Convert.ToBase64String( + Encoding.ASCII.GetBytes("admin:changeit"))); + return Task.CompletedTask; + }); + } + + public async Task Append( + string streamName, + int count = 1, + ulong? expectedRevision = null, + bool noStream = false, + string data = "event", + string eventType = "event") + { + using var call = _client.BatchAppend(CallOptions); + var options = new BatchAppendReq.Types.Options + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(streamName) } + }; + if (expectedRevision.HasValue) + options.StreamPosition = expectedRevision.Value; + else if (noStream) + options.NoStream = new(); + else + options.Any = new(); + + var request = new BatchAppendReq + { + CorrelationId = Uuid.NewUuid().ToDto(), + IsFinal = true, + Options = options + }; + for (var index = 0; index < count; index++) + { + request.ProposedMessages.Add(new BatchAppendReq.Types.ProposedMessage + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFromUtf8(data), + Metadata = + { + [GrpcMetadata.Type] = eventType, + [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson + } + }); + } + await call.RequestStream.WriteAsync(request); + await call.RequestStream.CompleteAsync(); + Assert.True(await call.ResponseStream.MoveNext()); + return call.ResponseStream.Current; + } + + public async Task AppendSingle(string streamName) + { + using var call = _client.Append(CallOptions); + await call.RequestStream.WriteAsync(new AppendReq + { + Options = new() + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(streamName) }, + Any = new() + } + }); + await call.RequestStream.WriteAsync(new AppendReq + { + ProposedMessage = new() + { + Id = Uuid.NewUuid().ToDto(), + Data = ByteString.CopyFromUtf8("event"), + Metadata = + { + [GrpcMetadata.Type] = "event", + [GrpcMetadata.ContentType] = GrpcMetadata.ContentTypes.ApplicationJson + } + } + }); + await call.RequestStream.CompleteAsync(); + return await call.ResponseAsync; + } + + public async Task Read( + string streamName, + ulong revision, + ulong count, + ReadReq.Types.Options.Types.ReadDirection direction = + ReadReq.Types.Options.Types.ReadDirection.Forwards) + { + using var call = _client.Read(new ReadReq + { + Options = new() + { + Stream = new() + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(streamName) }, + Revision = revision + }, + Count = count, + ReadDirection = direction, + NoFilter = new(), + UuidOption = new() { Structured = new() } + } + }, CallOptions); + return await call.ResponseStream.ReadAllAsync().ToArrayAsync(); + } + + public void Dispose() => _channel?.Dispose(); +} diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs new file mode 100644 index 0000000000..4a6aab17a8 --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs @@ -0,0 +1,90 @@ +using System; +using System.IO; +using System.Linq; +using System.Threading.Tasks; +using EventStore.Client.Streams; +using EventStore.Core.Index; +using EventStore.Core.Tests.Helpers; +using NUnit.Framework; + +namespace EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests; + +[Category("LongRunning")] +public class HashCollisionGrpcBoundaryTests : SpecificationWithDirectoryPerTestFixture +{ + private const string FirstStream = "account--696193173"; + private const string SecondStream = "LPN-FC002_LPK51001"; + private MiniNode _node; + private GrpcStreamEdgeOperations _grpc; + private string _dbPath; + + [OneTimeSetUp] + public override async Task TestFixtureSetUp() + { + await base.TestFixtureSetUp(); + _dbPath = Path.Combine(PathName, "collision-node-db"); + await StartNode(waitForAdminUserCreation: true); + } + + [OneTimeTearDown] + public override async Task TestFixtureTearDown() + { + _grpc?.Dispose(); + if (_node is not null) + await _node.Shutdown(); + await base.TestFixtureTearDown(); + } + + [Test] + public async Task does_not_return_a_colliding_stream_after_the_read_limit_is_reached() + { + AssertSuccess(await _grpc.Append(FirstStream, noStream: true), 0); + for (var revision = 0; revision < 100; revision++) + AssertSuccess(await _grpc.Append(SecondStream), (ulong)revision); + + await _node.Shutdown(keepDb: true); + _grpc.Dispose(); + await StartNode(waitForAdminUserCreation: false); + + var firstRead = await _grpc.Read(FirstStream, 0, 1); + Assert.That(firstRead.Single().ContentCase, + Is.EqualTo(ReadResp.ContentOneofCase.StreamNotFound)); + + var secondRead = await _grpc.Read(SecondStream, 99, 1); + Assert.That(secondRead.Single(x => x.Event is not null).Event.Event.StreamRevision, + Is.EqualTo(99)); + + var append = await _grpc.AppendSingle(FirstStream); + Assert.That(append.ResultCase, Is.EqualTo(AppendResp.ResultOneofCase.WrongExpectedVersion)); + Assert.That(append.WrongExpectedVersion.CurrentRevisionOptionCase, + Is.EqualTo(AppendResp.Types.WrongExpectedVersion.CurrentRevisionOptionOneofCase.None)); + + var batchAppend = await _grpc.Append(FirstStream); + Assert.That(batchAppend.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Error)); + Assert.That(batchAppend.Error.Code, Is.EqualTo(Google.Rpc.Code.AlreadyExists)); + var error = batchAppend.Error.Details.Unpack(); + Assert.That(error.CurrentStreamRevisionOptionCase, + Is.EqualTo(EventStore.Client.WrongExpectedVersion.CurrentStreamRevisionOptionOneofCase.None)); + } + + private async Task StartNode(bool waitForAdminUserCreation) + { + _node = new MiniNode(PathName, + dbPath: _dbPath, + memTableSize: 20, + hashCollisionReadLimit: 1, + indexBitnessVersion: PTableVersions.IndexV4, + hash32bit: true, + streamExistenceFilterSize: 0); + await _node.Start(); + if (waitForAdminUserCreation) + await _node.AdminUserCreated; + _grpc = new GrpcStreamEdgeOperations(_node); + } + + private static void AssertSuccess(BatchAppendResp response, ulong expectedRevision) + { + Assert.That(response.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Success)); + Assert.That(response.Success.CurrentRevision, Is.EqualTo(expectedRevision)); + } +} diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs new file mode 100644 index 0000000000..d79124ee4d --- /dev/null +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs @@ -0,0 +1,105 @@ +using System; +using System.Linq; +using System.Threading; +using System.Threading.Tasks; +using EventStore.Client.Streams; +using Google.Protobuf; +using Grpc.Core; +using NUnit.Framework; +using Streams = EventStore.Client.Streams.Streams; + +namespace EventStore.Core.Tests.Services.Transport.Grpc.StreamsTests; + +[TestFixture(typeof(LogFormat.V2), typeof(string))] +[Category("LongRunning")] +public class StreamRevisionAboveIntMaxTests + : GrpcSpecificationWithExistingRecords +{ + private const long FirstRevision = (long)int.MaxValue + 1; + private const string StreamName = "grpc-stream-revision-above-int-max"; + private GrpcStreamEdgeOperations _grpc; + private readonly Guid[] _eventIds = new Guid[5]; + + public override async ValueTask WriteTestScenario(CancellationToken token) + { + for (var index = 0; index < _eventIds.Length; index++) + { + var record = await WriteSingleEvent(StreamName, FirstRevision + index, + new string('.', 3000), token: token); + _eventIds[index] = record.EventId; + } + } + + public override async Task Given() + { + _grpc = new GrpcStreamEdgeOperations(Channel); + var metadata = await _grpc.Append("$$" + StreamName, data: "{\"$tb\":2147483648}", + eventType: "$metadata"); + Assert.That(metadata.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Success)); + } + + [Test] + public async Task reads_revisions_above_int_max_in_both_directions() + { + var forwards = await _grpc.Read(StreamName, (ulong)FirstRevision, 5); + var backwards = await _grpc.Read(StreamName, (ulong)(FirstRevision + 4), 5, + ReadReq.Types.Options.Types.ReadDirection.Backwards); + + CollectionAssert.AreEqual(_eventIds, + forwards.Where(x => x.Event is not null) + .Select(x => EventStore.Core.Services.Transport.Grpc.Uuid.FromDto(x.Event.Event.Id).ToGuid())); + CollectionAssert.AreEqual(_eventIds.Reverse(), + backwards.Where(x => x.Event is not null) + .Select(x => EventStore.Core.Services.Transport.Grpc.Uuid.FromDto(x.Event.Event.Id).ToGuid())); + CollectionAssert.AreEqual(Enumerable.Range(0, 5).Select(x => (ulong)(FirstRevision + x)), + forwards.Where(x => x.Event is not null).Select(x => x.Event.Event.StreamRevision)); + } + + [Test] + public async Task appends_at_a_revision_above_int_max_and_rejects_an_incorrect_revision() + { + var success = await _grpc.Append(StreamName, expectedRevision: (ulong)(FirstRevision + 4)); + Assert.That(success.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Success)); + Assert.That(success.Success.CurrentRevision, Is.EqualTo((ulong)(FirstRevision + 5))); + + var mismatch = await _grpc.Append(StreamName, expectedRevision: (ulong)(FirstRevision + 15)); + Assert.That(mismatch.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Error)); + Assert.That(mismatch.Error.Code, Is.EqualTo(Google.Rpc.Code.AlreadyExists)); + } + + [Test] + public async Task catch_up_subscription_delivers_revisions_above_int_max() + { + var client = new Streams.StreamsClient(Channel); + using var subscription = client.Read(new ReadReq + { + Options = new() + { + Subscription = new(), + NoFilter = new(), + ReadDirection = ReadReq.Types.Options.Types.ReadDirection.Forwards, + UuidOption = new() { Structured = new() }, + Stream = new() + { + Revision = (ulong)(FirstRevision - 1), + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(StreamName) } + } + } + }, GetCallOptions(AdminCredentials).WithDeadline(DateTime.UtcNow.AddSeconds(20))); + + Assert.That(await subscription.ResponseStream.MoveNext(), Is.True); + Assert.That(subscription.ResponseStream.Current.ContentCase, + Is.EqualTo(ReadResp.ContentOneofCase.Confirmation)); + + for (var index = 0; index < _eventIds.Length; index++) + { + Assert.That(await subscription.ResponseStream.MoveNext(), Is.True); + var response = subscription.ResponseStream.Current; + Assert.That(response.ContentCase, Is.EqualTo(ReadResp.ContentOneofCase.Event)); + Assert.That(response.Event.Event.StreamRevision, + Is.EqualTo((ulong)(FirstRevision + index))); + Assert.That(EventStore.Core.Services.Transport.Grpc.Uuid.FromDto(response.Event.Event.Id).ToGuid(), + Is.EqualTo(_eventIds[index])); + } + } +} diff --git a/src/EventStore.Core/Services/Transport/Grpc/Status.cs b/src/EventStore.Core/Services/Transport/Grpc/Status.cs index 66c58954ee..531b5d246c 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/Status.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/Status.cs @@ -1,5 +1,4 @@ using EventStore.Client; -using EventStore.Core.Services.Transport.Common; using Google.Protobuf.WellKnownTypes; using Empty = Google.Protobuf.WellKnownTypes.Empty; @@ -8,11 +7,11 @@ namespace Google.Rpc { partial class Status { - public static Status WrongExpectedVersion(StreamRevision currentStreamRevision, + public static Status WrongExpectedVersion(long currentVersion, long expectedVersion) => new() { Message = nameof(WrongExpectedVersion), - Details = Any.Pack(EventStore.Client.WrongExpectedVersion.Create(currentStreamRevision, expectedVersion)), + Details = Any.Pack(EventStore.Client.WrongExpectedVersion.Create(currentVersion, expectedVersion)), Code = Code.AlreadyExists }; diff --git a/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs b/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs index 3cc40c215b..84ef90a8d8 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs @@ -201,7 +201,7 @@ void HandleWriteEventsCompleted(Message message) response.WrongExpectedVersion.CurrentNoStream = new Empty(); response.WrongExpectedVersion.NoStream2060 = new Empty(); } - else + else if (completed.CurrentVersion >= 0) { response.WrongExpectedVersion.CurrentRevision = StreamRevision.FromInt64(completed.CurrentVersion); diff --git a/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs b/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs index 5249ffb8f1..21206e7056 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs @@ -277,7 +277,7 @@ BatchAppendResp ConvertMessage(Message message) OperationResult.WrongExpectedVersion => new BatchAppendResp { Error = Status.WrongExpectedVersion( - StreamRevision.FromInt64(completed.CurrentVersion), + completed.CurrentVersion, clientWriteRequest.ExpectedVersion) }, OperationResult.AccessDenied => new BatchAppendResp { Error = Status.AccessDenied }, diff --git a/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs b/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs index 952b8c0a34..6db3e9d2a0 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs @@ -7,15 +7,11 @@ namespace EventStore.Client { partial class WrongExpectedVersion { - public static WrongExpectedVersion Create(StreamRevision currentStreamRevision, - long expectedStreamPosition) => new() + public static WrongExpectedVersion Create(long currentVersion, + long expectedStreamPosition) + { + var result = new WrongExpectedVersion { - currentStreamRevisionOption_ = currentStreamRevision == StreamRevision.End - ? new Google.Protobuf.WellKnownTypes.Empty() - : currentStreamRevision.ToUInt64(), - currentStreamRevisionOptionCase_ = currentStreamRevision == StreamRevision.End - ? CurrentStreamRevisionOptionOneofCase.CurrentNoStream - : CurrentStreamRevisionOptionOneofCase.CurrentStreamRevision, expectedStreamPositionOption_ = expectedStreamPosition switch { Any or NoStream or StreamExists => @@ -30,5 +26,19 @@ public static WrongExpectedVersion Create(StreamRevision currentStreamRevision, _ => ExpectedStreamPositionOptionOneofCase.ExpectedStreamPosition } }; + if (currentVersion == NoStream) + { + result.currentStreamRevisionOption_ = new Google.Protobuf.WellKnownTypes.Empty(); + result.currentStreamRevisionOptionCase_ = + CurrentStreamRevisionOptionOneofCase.CurrentNoStream; + } + else if (currentVersion >= 0) + { + result.currentStreamRevisionOption_ = StreamRevision.FromInt64(currentVersion).ToUInt64(); + result.currentStreamRevisionOptionCase_ = + CurrentStreamRevisionOptionOneofCase.CurrentStreamRevision; + } + return result; + } } } From 30c4568c896ad3375194bc0f155bd37adc394db7 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Fri, 25 Sep 2026 15:59:16 -0400 Subject: [PATCH 2/3] fix(grpc): keep edge parity checks reliable Signed-off-by: Yordis Prieto --- .../Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs | 3 +-- .../Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs | 4 ++++ 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs index 4a6aab17a8..859277ebb3 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/HashCollisionGrpcBoundaryTests.cs @@ -39,8 +39,7 @@ public override async Task TestFixtureTearDown() public async Task does_not_return_a_colliding_stream_after_the_read_limit_is_reached() { AssertSuccess(await _grpc.Append(FirstStream, noStream: true), 0); - for (var revision = 0; revision < 100; revision++) - AssertSuccess(await _grpc.Append(SecondStream), (ulong)revision); + AssertSuccess(await _grpc.Append(SecondStream, count: 100), 99); await _node.Shutdown(keepDb: true); _grpc.Dispose(); diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs index d79124ee4d..027fe470a2 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs @@ -65,6 +65,10 @@ public async Task appends_at_a_revision_above_int_max_and_rejects_an_incorrect_r var mismatch = await _grpc.Append(StreamName, expectedRevision: (ulong)(FirstRevision + 15)); Assert.That(mismatch.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Error)); Assert.That(mismatch.Error.Code, Is.EqualTo(Google.Rpc.Code.AlreadyExists)); + var detail = mismatch.Error.Details.Unpack(); + Assert.That(detail.CurrentStreamRevisionOptionCase, + Is.EqualTo(EventStore.Client.WrongExpectedVersion.CurrentStreamRevisionOptionOneofCase.CurrentStreamRevision)); + Assert.That(detail.CurrentStreamRevision, Is.EqualTo((ulong)(FirstRevision + 5))); } [Test] From 851727cdbb11869f797b32b5d7cd1b04e4f7b222 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Fri, 25 Sep 2026 17:40:44 -0400 Subject: [PATCH 3/3] fix(grpc): preserve current-version semantics Signed-off-by: Yordis Prieto --- .../StreamsTests/GrpcStreamEdgeOperations.cs | 16 +++--- .../StreamRevisionAboveIntMaxTests.cs | 22 +++++++++ .../Transport/Grpc/CurrentStreamVersion.cs | 49 +++++++++++++++++++ .../Services/Transport/Grpc/Status.cs | 3 +- .../Services/Transport/Grpc/Streams.Append.cs | 9 ++-- .../Transport/Grpc/Streams.BatchAppend.cs | 2 +- .../Transport/Grpc/WrongExpectedVersion.cs | 8 +-- 7 files changed, 93 insertions(+), 16 deletions(-) create mode 100644 src/EventStore.Core/Services/Transport/Grpc/CurrentStreamVersion.cs diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs index 452ace9dcf..c47ae5e12a 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/GrpcStreamEdgeOperations.cs @@ -91,16 +91,20 @@ public async Task Append( return call.ResponseStream.Current; } - public async Task AppendSingle(string streamName) + public async Task AppendSingle(string streamName, ulong? expectedRevision = null) { using var call = _client.Append(CallOptions); + var options = new AppendReq.Types.Options + { + StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(streamName) } + }; + if (expectedRevision.HasValue) + options.Revision = expectedRevision.Value; + else + options.Any = new(); await call.RequestStream.WriteAsync(new AppendReq { - Options = new() - { - StreamIdentifier = new() { StreamName = ByteString.CopyFromUtf8(streamName) }, - Any = new() - } + Options = options }); await call.RequestStream.WriteAsync(new AppendReq { diff --git a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs index 027fe470a2..93fa3ebbc0 100644 --- a/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs +++ b/src/EventStore.Core.Tests/Services/Transport/Grpc/StreamsTests/StreamRevisionAboveIntMaxTests.cs @@ -69,6 +69,28 @@ public async Task appends_at_a_revision_above_int_max_and_rejects_an_incorrect_r Assert.That(detail.CurrentStreamRevisionOptionCase, Is.EqualTo(EventStore.Client.WrongExpectedVersion.CurrentStreamRevisionOptionOneofCase.CurrentStreamRevision)); Assert.That(detail.CurrentStreamRevision, Is.EqualTo((ulong)(FirstRevision + 5))); + + var singleMismatch = await _grpc.AppendSingle(StreamName, (ulong)(FirstRevision + 15)); + Assert.That(singleMismatch.ResultCase, Is.EqualTo(AppendResp.ResultOneofCase.WrongExpectedVersion)); + Assert.That(singleMismatch.WrongExpectedVersion.CurrentRevisionOptionCase, + Is.EqualTo(AppendResp.Types.WrongExpectedVersion.CurrentRevisionOptionOneofCase.CurrentRevision)); + Assert.That(singleMismatch.WrongExpectedVersion.CurrentRevision, Is.EqualTo((ulong)(FirstRevision + 5))); + } + + [Test] + public async Task incorrect_revision_on_a_missing_stream_reports_no_stream() + { + var mismatch = await _grpc.Append(StreamName + "-missing", expectedRevision: 0); + Assert.That(mismatch.ResultCase, Is.EqualTo(BatchAppendResp.ResultOneofCase.Error)); + Assert.That(mismatch.Error.Code, Is.EqualTo(Google.Rpc.Code.AlreadyExists)); + var detail = mismatch.Error.Details.Unpack(); + Assert.That(detail.CurrentStreamRevisionOptionCase, + Is.EqualTo(EventStore.Client.WrongExpectedVersion.CurrentStreamRevisionOptionOneofCase.CurrentNoStream)); + + var singleMismatch = await _grpc.AppendSingle(StreamName + "-missing", expectedRevision: 0); + Assert.That(singleMismatch.ResultCase, Is.EqualTo(AppendResp.ResultOneofCase.WrongExpectedVersion)); + Assert.That(singleMismatch.WrongExpectedVersion.CurrentRevisionOptionCase, + Is.EqualTo(AppendResp.Types.WrongExpectedVersion.CurrentRevisionOptionOneofCase.CurrentNoStream)); } [Test] diff --git a/src/EventStore.Core/Services/Transport/Grpc/CurrentStreamVersion.cs b/src/EventStore.Core/Services/Transport/Grpc/CurrentStreamVersion.cs new file mode 100644 index 0000000000..331c397fe2 --- /dev/null +++ b/src/EventStore.Core/Services/Transport/Grpc/CurrentStreamVersion.cs @@ -0,0 +1,49 @@ +using System; +using EventStore.Core.Data; +using EventStore.Core.Services.Transport.Common; + +namespace EventStore.Core.Services.Transport.Grpc; + +internal readonly record struct CurrentStreamVersion +{ + private enum VersionKind + { + Unknown, + NoStream, + Known + } + + private readonly VersionKind _kind; + private readonly StreamRevision _revision; + + private CurrentStreamVersion(VersionKind kind, StreamRevision revision) + { + _kind = kind; + _revision = revision; + } + + public static CurrentStreamVersion Unknown { get; } = default; + public static CurrentStreamVersion NoStream { get; } = new(VersionKind.NoStream, default); + + public static CurrentStreamVersion Known(StreamRevision revision) + { + if (revision == StreamRevision.End) + throw new ArgumentOutOfRangeException(nameof(revision)); + return new CurrentStreamVersion(VersionKind.Known, revision); + } + + public static CurrentStreamVersion FromInt64(long value) => value switch + { + ExpectedVersion.NoStream => NoStream, + >= 0 => Known(StreamRevision.FromInt64(value)), + _ => Unknown + }; + + public bool IsNoStream => _kind == VersionKind.NoStream; + + public bool TryGetKnownRevision(out StreamRevision revision) + { + revision = _revision; + return _kind == VersionKind.Known; + } +} diff --git a/src/EventStore.Core/Services/Transport/Grpc/Status.cs b/src/EventStore.Core/Services/Transport/Grpc/Status.cs index 531b5d246c..19e4e32f56 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/Status.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/Status.cs @@ -1,4 +1,5 @@ using EventStore.Client; +using EventStore.Core.Services.Transport.Grpc; using Google.Protobuf.WellKnownTypes; using Empty = Google.Protobuf.WellKnownTypes.Empty; @@ -7,7 +8,7 @@ namespace Google.Rpc { partial class Status { - public static Status WrongExpectedVersion(long currentVersion, + internal static Status WrongExpectedVersion(CurrentStreamVersion currentVersion, long expectedVersion) => new() { Message = nameof(WrongExpectedVersion), diff --git a/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs b/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs index 84ef90a8d8..6776130fa0 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/Streams.Append.cs @@ -196,17 +196,18 @@ void HandleWriteEventsCompleted(Message message) break; } - if (completed.CurrentVersion == -1) + var currentVersion = CurrentStreamVersion.FromInt64(completed.CurrentVersion); + if (currentVersion.IsNoStream) { response.WrongExpectedVersion.CurrentNoStream = new Empty(); response.WrongExpectedVersion.NoStream2060 = new Empty(); } - else if (completed.CurrentVersion >= 0) + else if (currentVersion.TryGetKnownRevision(out var revision)) { response.WrongExpectedVersion.CurrentRevision = - StreamRevision.FromInt64(completed.CurrentVersion); + revision; response.WrongExpectedVersion.CurrentRevision2060 = - StreamRevision.FromInt64(completed.CurrentVersion); + revision; } appendResponseSource.TrySetResult(response); diff --git a/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs b/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs index 21206e7056..29af304447 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/Streams.BatchAppend.cs @@ -277,7 +277,7 @@ BatchAppendResp ConvertMessage(Message message) OperationResult.WrongExpectedVersion => new BatchAppendResp { Error = Status.WrongExpectedVersion( - completed.CurrentVersion, + CurrentStreamVersion.FromInt64(completed.CurrentVersion), clientWriteRequest.ExpectedVersion) }, OperationResult.AccessDenied => new BatchAppendResp { Error = Status.AccessDenied }, diff --git a/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs b/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs index 6db3e9d2a0..7a3b8b513a 100644 --- a/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs +++ b/src/EventStore.Core/Services/Transport/Grpc/WrongExpectedVersion.cs @@ -7,7 +7,7 @@ namespace EventStore.Client { partial class WrongExpectedVersion { - public static WrongExpectedVersion Create(long currentVersion, + internal static WrongExpectedVersion Create(CurrentStreamVersion currentVersion, long expectedStreamPosition) { var result = new WrongExpectedVersion @@ -26,15 +26,15 @@ public static WrongExpectedVersion Create(long currentVersion, _ => ExpectedStreamPositionOptionOneofCase.ExpectedStreamPosition } }; - if (currentVersion == NoStream) + if (currentVersion.IsNoStream) { result.currentStreamRevisionOption_ = new Google.Protobuf.WellKnownTypes.Empty(); result.currentStreamRevisionOptionCase_ = CurrentStreamRevisionOptionOneofCase.CurrentNoStream; } - else if (currentVersion >= 0) + else if (currentVersion.TryGetKnownRevision(out var revision)) { - result.currentStreamRevisionOption_ = StreamRevision.FromInt64(currentVersion).ToUInt64(); + result.currentStreamRevisionOption_ = revision.ToUInt64(); result.currentStreamRevisionOptionCase_ = CurrentStreamRevisionOptionOneofCase.CurrentStreamRevision; }