From cb9251b54bb9bb750456dfbb0712f06f23fb2c6a Mon Sep 17 00:00:00 2001 From: void-ptr974 Date: Sat, 1 Aug 2026 22:14:02 +0800 Subject: [PATCH] [fix][broker] Wait for delayed delivery snapshot creation before trimming --- .../bucket/BucketDelayedDeliveryTracker.java | 41 ++++++--- .../BucketDelayedDeliveryTrackerTest.java | 83 ++++++++++++++++++- 2 files changed, 110 insertions(+), 14 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java index 1bcfcd5eb986e..9b2e659ec7e02 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTracker.java @@ -126,6 +126,11 @@ public static record SnapshotKey(long ledgerId, long entryId) {} private volatile CompletableFuture trimFuture; + @VisibleForTesting + CompletableFuture getTrimFuture() { + return trimFuture; + } + public BucketDelayedDeliveryTracker(AbstractPersistentDispatcherMultipleConsumers dispatcher, Timer timer, long tickTimeMillis, boolean isDelayedDeliveryDeliverAtTimeStrict, @@ -918,20 +923,32 @@ private synchronized CompletableFuture asyncTrimImmutableBuckets() { private CompletableFuture deleteBucketSnapshot(String ledgerName, Range range, ImmutableBucket bucket) { - return bucket.asyncDeleteBucketSnapshot(stats) - .handle((__, t) -> { - if (t != null) { - log.warn().attr("LedgerName", ledgerName) - .attr("BucketKey", bucket.bucketKey()) - .log("Failed to delete bucket snapshot"); - throw new CompletionException(t); - } + return bucket.getSnapshotCreateFuture().orElse(NULL_LONG_PROMISE) + .thenCompose(bucketId -> { synchronized (this) { - snapshotSegmentLastIndexMap.entrySet().removeIf(entry -> entry.getValue() == bucket); - removeBucket(range); - numberDelayedMessages.addAndGet(-bucket.getNumberBucketDelayedMessages()); + Long firstLedgerId = firstActiveLedgerId(); + if (INVALID_BUCKET_ID.equals(bucketId) + || firstLedgerId == null + || range.upperEndpoint() >= firstLedgerId + || immutableBuckets.asMapOfRanges().get(range) != bucket) { + return CompletableFuture.completedFuture(null); + } } - return null; + return bucket.asyncDeleteBucketSnapshot(stats).thenRun(() -> { + synchronized (this) { + if (immutableBuckets.asMapOfRanges().get(range) != bucket) { + return; + } + snapshotSegmentLastIndexMap.entrySet().removeIf(entry -> entry.getValue() == bucket); + removeBucket(range); + numberDelayedMessages.addAndGet(-bucket.getNumberBucketDelayedMessages()); + } + }); + }).exceptionally(t -> { + log.warn().attr("LedgerName", ledgerName) + .attr("BucketKey", bucket.bucketKey()) + .log("Failed to delete bucket snapshot"); + throw new CompletionException(t); }); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java index 28626496e49ae..2483063afaeb2 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/delayed/bucket/BucketDelayedDeliveryTrackerTest.java @@ -57,6 +57,8 @@ import org.apache.pulsar.broker.delayed.AbstractDeliveryTrackerTest; import org.apache.pulsar.broker.delayed.MockBucketSnapshotStorage; import org.apache.pulsar.broker.delayed.MockManagedCursor; +import org.apache.pulsar.broker.delayed.proto.SnapshotMetadata; +import org.apache.pulsar.broker.delayed.proto.SnapshotSegment; import org.apache.pulsar.broker.service.persistent.AbstractPersistentDispatcherMultipleConsumers; import org.awaitility.Awaitility; import org.roaringbitmap.RoaringBitmap; @@ -556,6 +558,28 @@ public CompletableFuture deleteBucketSnapshot(long bucketId) { } } + private static class BlockingCreateStorage extends MockBucketSnapshotStorage { + final CompletableFuture allowCreate = new CompletableFuture<>(); + final AtomicLong createCalls = new AtomicLong(); + final AtomicLong deleteCalls = new AtomicLong(); + + @Override + public CompletableFuture createBucketSnapshot( + SnapshotMetadata snapshotMetadata, List bucketSnapshotSegments, String bucketKey, + String topicName, String cursorName) { + createCalls.incrementAndGet(); + return super.createBucketSnapshot(snapshotMetadata, bucketSnapshotSegments, bucketKey, topicName, + cursorName) + .thenCompose(bucketId -> allowCreate.thenApply(ignored -> bucketId)); + } + + @Override + public CompletableFuture deleteBucketSnapshot(long bucketId) { + deleteCalls.incrementAndGet(); + return super.deleteBucketSnapshot(bucketId); + } + } + private static class RecordingDeleteStorage extends MockBucketSnapshotStorage { final AtomicLong deleteCalls = new AtomicLong(); @@ -586,11 +610,17 @@ private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, int m private TrackerWithStorage createTrackerWithMockLedger(long firstLedgerId, int maxNumBuckets, MockBucketSnapshotStorage storage) throws Exception { + return createTrackerWithMockLedger(new AtomicLong(firstLedgerId), maxNumBuckets, storage); + } + + private TrackerWithStorage createTrackerWithMockLedger(AtomicLong firstLedgerId, int maxNumBuckets, + MockBucketSnapshotStorage storage) + throws Exception { storage.start(); ManagedLedger mockLedger = mock(ManagedLedger.class); NavigableMap ledgerInfo = new TreeMap<>(); - ledgerInfo.put(firstLedgerId, mock(LedgerInfo.class)); + ledgerInfo.put(firstLedgerId.get(), mock(LedgerInfo.class)); when(mockLedger.getLedgersInfo()).thenReturn(ledgerInfo); when(mockLedger.getName()).thenReturn("test-ledger"); @@ -602,7 +632,7 @@ public ManagedLedger getManagedLedger() { @Override public Position getMarkDeletedPosition() { - return PositionFactory.create(firstLedgerId, -1); + return PositionFactory.create(firstLedgerId.get(), -1); } }; @@ -830,6 +860,55 @@ public void testClearRunsAfterInFlightTrimFailure() throws Exception { ts.close(); } + @Test + public void testTrimWaitsForSnapshotCreation() throws Exception { + BlockingCreateStorage storage = new BlockingCreateStorage(); + TrackerWithStorage ts = createTrackerWithMockLedger(50L, 5, storage); + try { + for (int i = 1; i <= 31; i++) { + ts.tracker.addMessage(i, i, i * 10); + } + + assertEquals(storage.createCalls.get(), 6L); + assertEquals(storage.deleteCalls.get(), 0L, + "Trim must not delete a snapshot while its creation is in flight"); + + storage.allowCreate.complete(null); + ts.tracker.getTrimFuture().get(1, TimeUnit.MINUTES); + + assertEquals(storage.deleteCalls.get(), 6L); + assertTrue(ts.tracker.getImmutableBuckets().asMapOfRanges().isEmpty()); + } finally { + storage.allowCreate.complete(null); + ts.close(); + } + } + + @Test + public void testTrimRevalidatesBucketAfterSnapshotCreation() throws Exception { + AtomicLong firstLedgerId = new AtomicLong(50L); + BlockingCreateStorage storage = new BlockingCreateStorage(); + TrackerWithStorage ts = createTrackerWithMockLedger(firstLedgerId, 5, storage); + try { + for (int i = 1; i <= 31; i++) { + ts.tracker.addMessage(i, i, 1000L); + } + + assertEquals(storage.createCalls.get(), 6L); + firstLedgerId.set(0L); + storage.allowCreate.complete(null); + ts.tracker.getTrimFuture().get(1, TimeUnit.MINUTES); + + assertEquals(storage.deleteCalls.get(), 0L, + "Trim must revalidate that a bucket is still orphaned after snapshot creation"); + assertEquals(ts.tracker.getImmutableBuckets().asMapOfRanges().size(), 6); + assertEquals(ts.tracker.getNumberOfDelayedMessages(), 31L); + } finally { + storage.allowCreate.complete(null); + ts.close(); + } + } + @Test public void testTrimWithNoOrphanedBuckets() throws Exception { TrackerWithStorage ts = createTrackerWithMockLedger(0L, 5);