Skip to content

Commit 5636dd8

Browse files
authored
feat(persistent-subscriptions): allow parked message truncation (#437)
- give operators a safe way to discard parked messages without redelivering them - keep persistent subscription parked-message metrics aligned with operator cleanup actions - expose the cleanup path through gRPC so automation does not need to mutate parked streams directly Signed-off-by: Yordis Prieto <yordis.prieto@gmail.com>
1 parent 6067c86 commit 5636dd8

19 files changed

Lines changed: 766 additions & 21 deletions

proto.lock

Lines changed: 51 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4326,6 +4326,51 @@
43264326
{
43274327
"name": "ReplayParkedResp"
43284328
},
4329+
{
4330+
"name": "TruncateParkedReq",
4331+
"fields": [
4332+
{
4333+
"id": 1,
4334+
"name": "options",
4335+
"type": "Options"
4336+
}
4337+
],
4338+
"messages": [
4339+
{
4340+
"name": "Options",
4341+
"fields": [
4342+
{
4343+
"id": 1,
4344+
"name": "group_name",
4345+
"type": "string"
4346+
},
4347+
{
4348+
"id": 2,
4349+
"name": "stream_identifier",
4350+
"type": "event_store.client.StreamIdentifier"
4351+
},
4352+
{
4353+
"id": 3,
4354+
"name": "all",
4355+
"type": "event_store.client.Empty"
4356+
},
4357+
{
4358+
"id": 4,
4359+
"name": "stop_at",
4360+
"type": "int64"
4361+
},
4362+
{
4363+
"id": 5,
4364+
"name": "no_limit",
4365+
"type": "event_store.client.Empty"
4366+
}
4367+
]
4368+
}
4369+
]
4370+
},
4371+
{
4372+
"name": "TruncateParkedResp"
4373+
},
43294374
{
43304375
"name": "ListReq",
43314376
"fields": [
@@ -4441,6 +4486,11 @@
44414486
"in_type": "ReplayParkedReq",
44424487
"out_type": "ReplayParkedResp"
44434488
},
4489+
{
4490+
"name": "TruncateParked",
4491+
"in_type": "TruncateParkedReq",
4492+
"out_type": "TruncateParkedResp"
4493+
},
44444494
{
44454495
"name": "List",
44464496
"in_type": "ListReq",
@@ -7280,4 +7330,4 @@
72807330
}
72817331
}
72827332
]
7283-
}
7333+
}

src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionMessageParkerTests.cs

Lines changed: 57 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -360,6 +360,63 @@ public async Task should_have_one_parked_message()
360360
}
361361

362362

363+
[TestFixture(typeof(LogFormat.V2), typeof(string))]
364+
public class given_messages_are_parked_and_then_truncated<TLogFormat, TStreamId> : TestFixtureWithExistingEvents<TLogFormat, TStreamId>
365+
{
366+
private PersistentSubscriptionMessageParker _messageParker;
367+
private string _streamId = Guid.NewGuid().ToString();
368+
private TaskCompletionSource<bool> _parked;
369+
private TaskCompletionSource<bool> _done = new TaskCompletionSource<bool>();
370+
371+
protected override void Given()
372+
{
373+
base.Given();
374+
375+
AllWritesSucceed();
376+
377+
_parked = new TaskCompletionSource<bool>();
378+
_messageParker = new PersistentSubscriptionMessageParker(_streamId, _ioDispatcher);
379+
_messageParker.BeginParkMessage(CreateResolvedEvent(0, 0), "testing", (_, __) =>
380+
{
381+
_messageParker.BeginParkMessage(CreateResolvedEvent(1, 100), "testing", (_, __) =>
382+
{
383+
_parked.SetResult(true);
384+
});
385+
});
386+
}
387+
388+
[Test]
389+
public async Task should_have_no_parked_messages()
390+
{
391+
await _parked.Task;
392+
_messageParker.BeginMarkParkedMessagesTruncated(2, _ =>
393+
{
394+
Assert.Zero(_messageParker.ParkedMessageCount);
395+
Assert.Null(_messageParker.GetOldestParkedMessage);
396+
Assert.Zero(_messageParker.ParkedMessageReplays);
397+
_done.TrySetResult(true);
398+
});
399+
await _done.Task.WithTimeout();
400+
}
401+
402+
[Test]
403+
public async Task should_not_lower_the_truncate_before_watermark()
404+
{
405+
await _parked.Task;
406+
_messageParker.BeginMarkParkedMessagesTruncated(2, _ =>
407+
{
408+
_messageParker.BeginMarkParkedMessagesTruncated(0, __ =>
409+
{
410+
Assert.Zero(_messageParker.ParkedMessageCount);
411+
Assert.Null(_messageParker.GetOldestParkedMessage);
412+
_done.TrySetResult(true);
413+
});
414+
});
415+
await _done.Task.WithTimeout();
416+
}
417+
}
418+
419+
363420
[TestFixture(typeof(LogFormat.V2), typeof(string))]
364421
public class given_read_backwards_fails_when_getting_stats<TLogFormat, TStreamId> : TestFixtureWithExistingEvents<TLogFormat, TStreamId>
365422
{

src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionServiceNotReadyTests.cs

Lines changed: 12 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -172,6 +172,18 @@ public void replay_parked_replies_not_ready()
172172
AssertNotReady(envelope, correlationId);
173173
}
174174

175+
[Test]
176+
public void truncate_parked_replies_not_ready()
177+
{
178+
var envelope = new FakeEnvelope();
179+
var correlationId = Guid.NewGuid();
180+
181+
_sut.Handle(new ClientMessage.TruncateParkedMessages(
182+
Guid.NewGuid(), correlationId, envelope, "stream", "group", null, ClaimsPrincipal.Current));
183+
184+
AssertNotReady(envelope, correlationId);
185+
}
186+
175187
private static void AssertNotReady(FakeEnvelope envelope, Guid correlationId)
176188
{
177189
Assert.That(envelope.Replies, Has.Count.EqualTo(1));

src/EventStore.Core.Tests/Services/PersistentSubscription/PersistentSubscriptionTests.cs

Lines changed: 22 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -2980,6 +2980,16 @@ public void RecordParkMessageRequest(ParkReason parkReason)
29802980
}
29812981
}
29822982

2983+
public void RecordParkedMessageReplay()
2984+
{
2985+
ParkedMessageReplays++;
2986+
}
2987+
2988+
public void RecordParkedMessageTruncate()
2989+
{
2990+
ParkedMessageTruncates++;
2991+
}
2992+
29832993
public void BeginParkMessage(ResolvedEvent ev, string reason,
29842994
Action<ResolvedEvent, OperationResult> completed)
29852995
{
@@ -2988,23 +2998,30 @@ public void BeginParkMessage(ResolvedEvent ev, string reason,
29882998
_parkMessageCompleted = completed;
29892999
}
29903000

2991-
public void BeginReadEndSequence(Action<long?> completed)
3001+
public void BeginReadEndSequence(Action<ParkedStreamEndReadResult> completed)
29923002
{
29933003
BeginReadEndSequenceCount++;
2994-
ParkedMessageReplays++;
29953004
if (_lastParkedEventNumber == -1)
29963005
{
2997-
completed(null); //NoStream
3006+
completed(ParkedStreamEndReadResult.Success(null)); //NoStream
29983007
}
29993008
else
30003009
{
3001-
completed(_lastParkedEventNumber);
3010+
completed(ParkedStreamEndReadResult.Success(_lastParkedEventNumber));
30023011
}
30033012
}
30043013

30053014
public void BeginMarkParkedMessagesReprocessed(long sequence, DateTime? dateTime, bool updateOldestParkedMessage)
30063015
{
30073016
MarkedAsProcessed = sequence;
3017+
_lastTruncateBefore = sequence;
3018+
}
3019+
3020+
public void BeginMarkParkedMessagesTruncated(long sequence, Action<OperationResult> completed)
3021+
{
3022+
MarkedAsProcessed = sequence;
3023+
_lastTruncateBefore = sequence;
3024+
completed?.Invoke(OperationResult.Success);
30083025
}
30093026

30103027
public void BeginDelete(Action<IPersistentSubscriptionMessageParker> completed)
@@ -3031,6 +3048,7 @@ public void BeginLoadStats(Action completed)
30313048
public long ParkedDueToClientNak { get; private set; }
30323049
public long ParkedDueToMaxRetries { get; private set; }
30333050
public long ParkedMessageReplays { get; private set; }
3051+
public long ParkedMessageTruncates { get; private set; }
30343052
}
30353053

30363054

0 commit comments

Comments
 (0)