Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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<TLogFormat, TStreamId> : TestFixtureWithExistingEvents<TLogFormat, TStreamId>
{
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<TLogFormat, TStreamId> : TestFixtureWithExistingEvents<TLogFormat, TStreamId>
{
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -78,10 +78,13 @@ private Event CreateStreamMetadataEvent(long? tb)
private void WriteStateCompleted(Action<ResolvedEvent, OperationResult> 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);
Expand Down
Loading