Skip to content
Open
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 @@ -72,9 +72,14 @@ public static EntryImpl create(LedgerEntry ledgerEntry, int expectedReadCount) {
entry.data.retain();
entry.readCountHandler = EntryReadCountHandlerImpl.maybeCreate(expectedReadCount);
// Reset the lazily-cached position LAST, after the id assignments: a recycled object can
// carry a stale Position materialized by a getPosition() call that raced past the recycle
// (deallocation nulls the field, but a late reader re-materializes it from the reset ids
// as (-1, -1)), and any racy lazy rebuild must observe the fresh legitimate ids.
// carry a Position poisoned by a getPosition() call that slipped in after the release
// (deallocation nulls the field, but the late reader re-materializes it from the reset
// ids as (-1, -1) and caches it in the pooled object), so the next create() would deliver
// the poisoned position with otherwise legitimate ids. This closes the SEQUENTIAL
// poison-and-reuse case; it is best-effort hardening only — a getPosition() concurrent
// with this create() can still read the old ids and publish its stale value after the
// reset (unsynchronized fields, no generation check). Post-release access is a caller
// bug and must be fixed at the call site, not here.
entry.position = null;
entry.setRefCnt(1);
return entry;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,12 @@
import static org.testng.Assert.assertTrue;
import io.netty.buffer.ByteBuf;
import io.netty.buffer.Unpooled;
import io.netty.util.concurrent.FastThreadLocalThread;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import org.apache.bookkeeper.client.api.LedgerEntry;
import org.apache.bookkeeper.client.impl.LedgerEntryImpl;
import org.apache.bookkeeper.mledger.Entry;
import org.apache.bookkeeper.mledger.Position;
import org.apache.bookkeeper.mledger.PositionFactory;
Expand Down Expand Up @@ -260,29 +266,100 @@ public void testCreateWithPositionThatIsImmutable() {
}

@Test
public void testRecycledObjectDoesNotInheritPoisonedPosition() {
public void testRecycledObjectDoesNotInheritPoisonedPosition() throws Exception {
// Netty's Recycler only pools instances for FastThreadLocalThreads; a plain test worker
// thread receives a no-op handle, every create() would return a fresh object, and this
// scenario would silently degrade to asserting fresh instances (see the review on #26707).
// Run it on a FastThreadLocalThread and prove instance identity with assertSame.
runOnFastThreadLocalThread(() -> {
assertTrue(warmUpRecyclerUntilReused(),
"recycler never handed the same instance back; the scenario cannot exercise recycling");

// The managed-ledger read path wraps a BookKeeper LedgerEntry.
assertRecycledPoisonedPositionIsNotInherited("create(LedgerEntry, int)",
() -> createFromLedgerEntry(5L, 10L),
() -> createFromLedgerEntry(6L, 20L), 6L, 20L);
// byte[] variant.
assertRecycledPoisonedPositionIsNotInherited("byte[] variant",
() -> EntryImpl.create(5L, 10L, new byte[]{1, 2, 3}),
() -> EntryImpl.create(6L, 20L, new byte[]{4, 5, 6}), 6L, 20L);
// ByteBuf variant.
assertRecycledPoisonedPositionIsNotInherited("ByteBuf variant",
() -> createFromRetainedBuffer(5L, 10L),
() -> createFromRetainedBuffer(6L, 20L), 6L, 20L);
});
}

private static void assertRecycledPoisonedPositionIsNotInherited(String variant,
Supplier<EntryImpl> firstEntry,
Supplier<EntryImpl> secondEntry,
long expectedLedgerId,
long expectedEntryId) {
// Given a legitimate entry that is released normally
EntryImpl first = EntryImpl.create(5L, 10L, new byte[]{1, 2, 3});
EntryImpl first = firstEntry.get();
first.release();

// When a getPosition() call slips in AFTER the release: deallocation nulls the lazy
// position field, so this late reader re-materializes it from the reset ids as (-1, -1)
// and leaves the poisoned value cached inside the pooled object.
first.getPosition();

// Then the next create() (the recycler hands back the most recently released object on
// the same thread) must not report that stale (-1, -1) position as its own — through
// both the byte[] and the ByteBuf variants, which own the lazy field.
EntryImpl second = EntryImpl.create(6L, 20L, new byte[]{4, 5, 6});
assertTrue(second.getPosition().compareTo(PositionFactory.create(6L, 20L)) == 0,
"byte[] variant: a recycled entry must not inherit the poisoned (-1, -1) position");
// Then the very next create() on this thread reuses that exact instance and must not
// report the poisoned (-1, -1) position as its own.
EntryImpl second = secondEntry.get();
assertSame(second, first, variant + ": expected the recycler to reuse the same instance");
assertTrue(second.getPosition().compareTo(PositionFactory.create(expectedLedgerId, expectedEntryId)) == 0,

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

assertTrue(x.compareTo(...) == 0, msg) produces a failure message like expected true but got false. Swapping to assertEquals would print the actual and expected positions on failure, which is much more useful when debugging:

assertEquals(second.getPosition(), PositionFactory.create(expectedLedgerId, expectedEntryId),
        variant + ": a recycled entry must not inherit the poisoned (-1, -1) position");

The compareTo-based form is already used elsewhere in this file, so this isn't blocking, but it's worth changing in a test whose whole point is catching bad values.

variant + ": a recycled entry must not inherit the poisoned (-1, -1) position");
second.release();
}

private static EntryImpl createFromLedgerEntry(long ledgerId, long entryId) {
ByteBuf buffer = Unpooled.wrappedBuffer(new byte[]{1, 2, 3});
LedgerEntry ledgerEntry = LedgerEntryImpl.create(ledgerId, entryId, buffer.readableBytes(), buffer);
// LedgerEntryImpl.create adopts the caller's buffer reference (released again on close()),
// while EntryImpl.create retains its own, so this ordering keeps the counts balanced.
EntryImpl entry = EntryImpl.create(ledgerEntry, 0);
ledgerEntry.close();
return entry;
}

private static EntryImpl createFromRetainedBuffer(long ledgerId, long entryId) {
ByteBuf buffer = Unpooled.wrappedBuffer(new byte[]{1, 2, 3});
EntryImpl entry = EntryImpl.create(ledgerId, entryId, buffer);
// create(long, long, ByteBuf, int) retains the buffer; drop the caller-owned reference.
buffer.release();
return entry;
}

second.getPosition(); // re-poison the recycled object
EntryImpl third = EntryImpl.create(7L, 30L, Unpooled.wrappedBuffer(new byte[]{7, 8}));
assertTrue(third.getPosition().compareTo(PositionFactory.create(7L, 30L)) == 0,
"ByteBuf variant: a recycled entry must not inherit the poisoned (-1, -1) position");
third.release();
private static boolean warmUpRecyclerUntilReused() {
// Absorb Netty's io.netty.recycler.ratio warm-up and prove the per-thread pool hands the
// same instance back; without it the assertions above could be checking fresh objects.
for (int i = 0; i < 64; i++) {
EntryImpl warm = EntryImpl.create(0L, 0L, new byte[]{0});
warm.release();
EntryImpl again = EntryImpl.create(0L, 0L, new byte[]{0});
boolean reused = again == warm;
again.release();
if (reused) {
return true;
}
}
return false;
}

private static void runOnFastThreadLocalThread(Runnable scenario) throws Exception {
CompletableFuture<Void> completion = new CompletableFuture<>();
Thread thread = new FastThreadLocalThread(() -> {
try {
scenario.run();
completion.complete(null);
} catch (Throwable error) {
completion.completeExceptionally(error);
}
});
thread.start();
// Rethrow assertion failures on the test thread.
completion.get(30, TimeUnit.SECONDS);

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

completion.get(30, TimeUnit.SECONDS) returns as soon as the future is resolved, but the thread is technically still alive until its run() returns. The gap is tiny, but adding thread.join() after completion.get(...) makes the lifecycle explicit and prevents any thread-leak warnings from strict test frameworks:

completion.get(30, TimeUnit.SECONDS);
thread.join();

}

private void assertEntryFields(EntryImpl entry, long expectedLedgerId, long expectedEntryId) {
Expand Down
Loading