|
58 | 58 | import org.apache.bookkeeper.mledger.ScanOutcome; |
59 | 59 | import org.apache.bookkeeper.mledger.util.ManagedLedgerUtils; |
60 | 60 | import org.apache.bookkeeper.test.MockedBookKeeperTestCase; |
| 61 | +import org.testng.annotations.DataProvider; |
61 | 62 | import org.testng.annotations.Test; |
62 | 63 |
|
63 | 64 | public class ReadCompletionAffinityTest extends MockedBookKeeperTestCase { |
@@ -205,7 +206,7 @@ public void readEntriesComplete(List<Entry> entries, Object ctx) { |
205 | 206 | queuedPool.set(ForkJoinTask.getPool()); |
206 | 207 | queuedCompletionStarted.complete(Thread.currentThread()); |
207 | 208 | try { |
208 | | - if (!releaseQueuedCompletion.await(20, TimeUnit.SECONDS)) { |
| 209 | + if (!awaitQueuedCompletion(releaseQueuedCompletion)) { |
209 | 210 | queuedCompletion.completeExceptionally( |
210 | 211 | new IllegalStateException("Timed out waiting to release queued completion")); |
211 | 212 | return; |
@@ -275,6 +276,88 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { |
275 | 276 | } |
276 | 277 | } |
277 | 278 |
|
| 279 | + @Test(timeOut = 30000) |
| 280 | + public void testNestedCachedChainReturnsToLedgerExecutorAcrossLedgers() throws Exception { |
| 281 | + int firstLedgerCount = OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 2; |
| 282 | + ManagedLedgerConfig config = inlineConfig().setMaxEntriesPerLedger(firstLedgerCount); |
| 283 | + config.setMinimumRolloverTime(0, TimeUnit.SECONDS); |
| 284 | + ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("completion-cached-chain-rollover", config); |
| 285 | + AtomicInteger storageReads = new AtomicInteger(); |
| 286 | + try { |
| 287 | + ManagedCursor cursor = ledger.openCursor("cursor"); |
| 288 | + List<Position> positions = new ArrayList<>(); |
| 289 | + for (int i = 0; i <= firstLedgerCount; i++) { |
| 290 | + positions.add(ledger.addEntry(new byte[] {(byte) i})); |
| 291 | + } |
| 292 | + assertThat(positions.get(firstLedgerCount).getLedgerId()) |
| 293 | + .isNotEqualTo(positions.get(0).getLedgerId()); |
| 294 | + // Resolve the closed ledger handle before starting the chain so opening it cannot cause a handoff. |
| 295 | + ledger.getLedgerHandle(positions.get(0).getLedgerId()).get(10, TimeUnit.SECONDS); |
| 296 | + ledger.entryCache.clear(); |
| 297 | + for (int i = 0; i < positions.size(); i++) { |
| 298 | + cacheEntry(ledger, positions.get(i), (byte) i); |
| 299 | + } |
| 300 | + bkc.setReadHandleInterceptor((ledgerId, first, last, entries) -> { |
| 301 | + storageReads.incrementAndGet(); |
| 302 | + return CompletableFuture.completedFuture(entries); |
| 303 | + }); |
| 304 | + CompletableFuture<Thread> ledgerWorker = new CompletableFuture<>(); |
| 305 | + ledger.getExecutor().execute(() -> ledgerWorker.complete(Thread.currentThread())); |
| 306 | + Thread worker = ledgerWorker.get(10, TimeUnit.SECONDS); |
| 307 | + Thread caller = Thread.currentThread(); |
| 308 | + |
| 309 | + AtomicInteger completions = new AtomicInteger(); |
| 310 | + AtomicReference<Thread> overflowThread = new AtomicReference<>(); |
| 311 | + AtomicReference<ForkJoinPool> overflowPool = new AtomicReference<>(); |
| 312 | + AtomicReference<Thread> boundaryThread = new AtomicReference<>(); |
| 313 | + AtomicReference<ForkJoinPool> boundaryPool = new AtomicReference<>(); |
| 314 | + CompletableFuture<List<Position>> crossedEntries = new CompletableFuture<>(); |
| 315 | + cursor.asyncReadEntries(1, new ReadEntriesCallback() { |
| 316 | + @Override |
| 317 | + public void readEntriesComplete(List<Entry> entries, Object ctx) { |
| 318 | + List<Position> returnedPositions = entries.stream().map(Entry::getPosition).toList(); |
| 319 | + entries.forEach(Entry::release); |
| 320 | + int completion = completions.incrementAndGet(); |
| 321 | + if (completion <= OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS) { |
| 322 | + cursor.asyncReadEntries(1, this, null, PositionFactory.LATEST); |
| 323 | + } else if (completion == OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 1) { |
| 324 | + overflowThread.set(Thread.currentThread()); |
| 325 | + overflowPool.set(ForkJoinTask.getPool()); |
| 326 | + // One entry remains in the first ledger. Reading two forces checkReadCompletion to |
| 327 | + // schedule the continuation on the ledger executor, even though both entries are cached. |
| 328 | + cursor.asyncReadEntries(2, this, null, PositionFactory.LATEST); |
| 329 | + } else { |
| 330 | + boundaryThread.set(Thread.currentThread()); |
| 331 | + boundaryPool.set(ForkJoinTask.getPool()); |
| 332 | + crossedEntries.complete(returnedPositions); |
| 333 | + } |
| 334 | + } |
| 335 | + |
| 336 | + @Override |
| 337 | + public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { |
| 338 | + crossedEntries.completeExceptionally(exception); |
| 339 | + } |
| 340 | + }, null, PositionFactory.LATEST); |
| 341 | + |
| 342 | + assertThat(crossedEntries.get(10, TimeUnit.SECONDS)) |
| 343 | + .containsExactly(positions.get(firstLedgerCount - 1), positions.get(firstLedgerCount)); |
| 344 | + assertThat(overflowThread.get()).isNotSameAs(caller); |
| 345 | + if (ForkJoinPool.getCommonPoolParallelism() > 1) { |
| 346 | + assertThat(overflowPool.get()).isSameAs(ForkJoinPool.commonPool()); |
| 347 | + assertThat(overflowThread.get()).isNotSameAs(worker); |
| 348 | + } else { |
| 349 | + assertThat(overflowThread.get()).isSameAs(worker); |
| 350 | + } |
| 351 | + assertThat(boundaryThread.get()).isSameAs(worker); |
| 352 | + assertThat(boundaryPool.get()).isNull(); |
| 353 | + assertThat(completions.get()).isEqualTo(OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 2); |
| 354 | + assertThat(storageReads.get()).isZero(); |
| 355 | + } finally { |
| 356 | + bkc.setReadHandleInterceptor(null); |
| 357 | + ledger.close(); |
| 358 | + } |
| 359 | + } |
| 360 | + |
278 | 361 | @Test |
279 | 362 | public void testLegacyCachedReadQueuesCompletionOffLedgerAndRunsItInlineOnLedger() throws Exception { |
280 | 363 | ManagedLedgerImpl ledger = (ManagedLedgerImpl) factory.open("completion-default-legacy", rawEntryConfig()); |
@@ -419,11 +502,18 @@ public void readEntriesFailed(ManagedLedgerException exception, Object ctx) { |
419 | 502 | } |
420 | 503 | } |
421 | 504 |
|
422 | | - @Test |
423 | | - public void testRejectedLedgerExecutorAfterDepthLimitReleasesEntriesAndCompletesFailureOnce() |
| 505 | + @DataProvider |
| 506 | + public Object[][] rejectedDepthLimitCompletionModes() { |
| 507 | + // Run in a fresh JVM with common-pool parallelism at most one to also cover the inline fallback. |
| 508 | + return ForkJoinPool.getCommonPoolParallelism() > 1 |
| 509 | + ? new Object[][] {{false}} : new Object[][] {{false}, {true}}; |
| 510 | + } |
| 511 | + |
| 512 | + @Test(dataProvider = "rejectedDepthLimitCompletionModes") |
| 513 | + public void testRejectedLedgerExecutorAfterDepthLimitReleasesEntriesAndCompletesFailureOnce(boolean inline) |
424 | 514 | throws Exception { |
425 | | - ManagedLedgerImpl ledger = spy((ManagedLedgerImpl) factory.open("completion-depth-rejection", |
426 | | - rawEntryConfig())); |
| 515 | + ManagedLedgerImpl ledger = spy((ManagedLedgerImpl) factory.open("completion-depth-rejection-" + inline, |
| 516 | + rawEntryConfig().setReadEntriesCallbackInline(inline))); |
427 | 517 | try { |
428 | 518 | ManagedCursorImpl cursor = (ManagedCursorImpl) ledger.openCursor("cursor"); |
429 | 519 | int count = OpReadEntry.MAX_NESTED_INLINE_COMPLETIONS + 1; |
@@ -509,6 +599,31 @@ public void testReadCompletionPolicyIsCapturedWhenLedgerOpens() throws Exception |
509 | 599 | } |
510 | 600 | } |
511 | 601 |
|
| 602 | + private static boolean awaitQueuedCompletion(CountDownLatch latch) throws InterruptedException { |
| 603 | + if (!ForkJoinTask.inForkJoinPool()) { |
| 604 | + return latch.await(20, TimeUnit.SECONDS); |
| 605 | + } |
| 606 | + ForkJoinPool.ManagedBlocker blocker = new ForkJoinPool.ManagedBlocker() { |
| 607 | + private final long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(20); |
| 608 | + |
| 609 | + @Override |
| 610 | + public boolean block() throws InterruptedException { |
| 611 | + long remaining = deadline - System.nanoTime(); |
| 612 | + if (remaining > 0) { |
| 613 | + latch.await(remaining, TimeUnit.NANOSECONDS); |
| 614 | + } |
| 615 | + return true; |
| 616 | + } |
| 617 | + |
| 618 | + @Override |
| 619 | + public boolean isReleasable() { |
| 620 | + return latch.getCount() == 0 || System.nanoTime() >= deadline; |
| 621 | + } |
| 622 | + }; |
| 623 | + ForkJoinPool.managedBlock(blocker); |
| 624 | + return latch.getCount() == 0; |
| 625 | + } |
| 626 | + |
512 | 627 | private static ByteBuf cacheEntry(ManagedLedgerImpl ledger, Position position, byte value) { |
513 | 628 | ByteBuf data = Unpooled.wrappedBuffer(new byte[] {value}); |
514 | 629 | EntryImpl entry = EntryImpl.create(position, data, 0); |
|
0 commit comments