From 93b85ff09971947cca3c425be8b9d5b177d97d20 Mon Sep 17 00:00:00 2001 From: Yordis Prieto Date: Mon, 10 Aug 2026 13:50:35 -0400 Subject: [PATCH] fix(persistent-subscriptions): keep parker statistics trustworthy Signed-off-by: Yordis Prieto --- ...ersistentSubscriptionMessageParkerTests.cs | 32 +++++++++++++++++++ .../PersistentSubscriptionMessageParker.cs | 9 ++++-- 2 files changed, 38 insertions(+), 3 deletions(-) diff --git a/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionMessageParkerTests.cs b/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionMessageParkerTests.cs index 2cb4674c6..2e5ace622 100644 --- a/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionMessageParkerTests.cs +++ b/src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionMessageParkerTests.cs @@ -276,6 +276,38 @@ public async Task should_count_max_retry_park_requests() } } + [TestFixture(typeof(LogFormat.V2), typeof(string))] + public class given_a_park_write_fails_after_a_message_was_parked : TestFixtureWithExistingEvents + { + private PersistentSubscriptionMessageParker _messageParker; + private OperationResult _result; + private DateTime? _oldestParkedMessage; + + protected override void Given() + { + base.Given(); + _messageParker = new PersistentSubscriptionMessageParker(Guid.NewGuid().ToString(), _ioDispatcher); + NoOtherStreams(); + AllWritesQueueUp(); + + _messageParker.BeginParkMessage(CreateResolvedEvent(0, 0), "testing", (_, __) => { }); + OneWriteCompletes(); + _oldestParkedMessage = _messageParker.GetOldestParkedMessage; + } + + [Test] + public void should_preserve_the_existing_parked_message_state() + { + _messageParker.BeginParkMessage(CreateResolvedEvent(1, 100), "testing", (_, result) => _result = result); + + CompleteWriteWithResult(OperationResult.CommitTimeout); + + Assert.AreEqual(OperationResult.CommitTimeout, _result); + Assert.AreEqual(1, _messageParker.ParkedMessageCount); + Assert.AreEqual(_oldestParkedMessage, _messageParker.GetOldestParkedMessage); + } + } + [TestFixture(typeof(LogFormat.V2), typeof(string))] public class given_messages_are_parked_and_then_replayed : TestFixtureWithExistingEvents { diff --git a/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionMessageParker.cs b/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionMessageParker.cs index 61cf3dd07..b84cc23a0 100644 --- a/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionMessageParker.cs +++ b/src/EventStore.Core/Services/PersistentSubscription/PersistentSubscriptionMessageParker.cs @@ -78,10 +78,13 @@ private Event CreateStreamMetadataEvent(long? tb) private void WriteStateCompleted(Action completed, ResolvedEvent ev, ClientMessage.WriteEventsCompleted msg, DateTime parkedMessageAdded) { - _lastParkedEventNumber = msg.LastEventNumber; - if (_oldestParkedMessage == null) + if (msg.Result == OperationResult.Success) { - _oldestParkedMessage = parkedMessageAdded.ToUniversalTime(); + _lastParkedEventNumber = msg.LastEventNumber; + if (_oldestParkedMessage == null) + { + _oldestParkedMessage = parkedMessageAdded.ToUniversalTime(); + } } completed?.Invoke(ev, msg.Result);