From d626d6030ab0cbbceca939a5b81a921ae1f7fec3 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Sun, 4 Oct 2026 01:02:11 +0300 Subject: [PATCH] [fix][ml] Decode the shared cached message metadata before publishing it Assisted-by: Claude Code (claude-opus-5-5) --- .../bookkeeper/mledger/impl/EntryImpl.java | 5 ++ .../impl/cache/RangeEntryCacheImpl.java | 2 +- .../mledger/impl/EntryImplTest.java | 51 +++++++++++++++++++ .../impl/cache/RangeEntryCacheImplTest.java | 13 ++--- 4 files changed, 62 insertions(+), 9 deletions(-) diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java index bcfc22c4d9784..a04b9eb541000 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/EntryImpl.java @@ -351,6 +351,11 @@ public synchronized void initializeMessageMetadataIfNeeded(String managedLedgerN try { MessageMetadata msgMetadata = new MessageMetadata(); Commands.parseMessageMetadata(data.duplicate(), msgMetadata); + // Copies of this entry share the instance across threads. Decode its lazily decoded fields now, so + // that readers don't decode them concurrently: a lazily decoded field is cached in a plain field, and + // LightProto creates an ASCII string without a constructor, so another thread could see the cached + // string before its contents. The volatile write below publishes the decoded fields. + msgMetadata.materialize(); this.messageMetadata = msgMetadata; } catch (Throwable t) { // The entry bytes are immutable; another cache reader cannot make a failed parse succeed. diff --git a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java index 8e43221b231e7..8f2943c0cc2cf 100644 --- a/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java +++ b/managed-ledger/src/main/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImpl.java @@ -562,7 +562,7 @@ public void accept(ReferenceCountedEntry entry) { } } // The visitor retains the cached entry while parsing. Initialize on the shared cached entry - // before copying, so fanout readers reuse one instance backed by the cache-owned buffer. + // before copying, so fanout readers reuse one instance, which is decoded when it's parsed. if (managedLedgerName != null && entry.getMessageMetadata() == null) { ((EntryImpl) entry).initializeMessageMetadataIfNeeded(managedLedgerName); } diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java index 38d7d34e9baeb..777bdd374835f 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/EntryImplTest.java @@ -31,9 +31,12 @@ import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; import io.netty.buffer.Unpooled; +import java.nio.charset.StandardCharsets; import org.apache.bookkeeper.mledger.Entry; import org.apache.bookkeeper.mledger.Position; import org.apache.bookkeeper.mledger.PositionFactory; +import org.apache.pulsar.common.api.proto.MessageMetadata; +import org.apache.pulsar.common.protocol.Commands; import org.testng.annotations.Test; public class EntryImplTest { @@ -56,6 +59,54 @@ public void testFailedMetadataInitializationIsNotRetried() { } } + @Test + public void testInitializedMessageMetadataDoesNotDecodeFromTheBufferLater() { + MessageMetadata metadata = new MessageMetadata() + .setProducerName("producer") + .setSequenceId(1) + .setPublishTime(2) + .setPartitionKey("cGFydGl0aW9uLWtleQ==") + .setPartitionKeyB64Encoded(true) + .setOrderingKey("ordering-key".getBytes(StandardCharsets.UTF_8)) + .setReplicatedFrom("cluster") + .setUuid("uuid-\u00e4") + .setSchemaVersion(new byte[] {1, 2}) + .setSchemaId(new byte[] {3, 4}) + .setEncryptionAlgo("algo") + .setEncryptionParam(new byte[] {5, 6}); + metadata.addReplicateTo("other-cluster"); + metadata.addProperty().setKey("key").setValue("value"); + metadata.addEncryptionKey().setKey("encryption-key").setValue(new byte[] {7, 8}) + .addMetadata().setKey("metadata-key").setValue("metadata-value"); + byte[] expected = metadata.toByteArray(); + ByteBuf bytes = Commands.serializeMetadataAndPayload(Commands.ChecksumType.None, metadata, + Unpooled.wrappedBuffer("payload".getBytes(StandardCharsets.UTF_8))); + EntryImpl entry = EntryImpl.create(1, 0, bytes); + bytes.release(); + try { + entry.initializeMessageMetadataIfNeeded("ledger"); + // Copies share the metadata across threads, so it must not decode fields lazily from the buffer, which + // concurrent readers would race on. Overwrite the buffer to show that it no longer reads it. + ByteBuf data = entry.getDataBuffer(); + data.setZero(0, data.capacity()); + MessageMetadata shared = entry.getMessageMetadata(); + assertEquals(shared.toByteArray(), expected); + assertEquals(shared.getProducerName(), "producer"); + assertEquals(shared.getPartitionKey(), "cGFydGl0aW9uLWtleQ=="); + assertEquals(new String(shared.getOrderingKey(), StandardCharsets.UTF_8), "ordering-key"); + assertEquals(shared.getUuid(), "uuid-\u00e4"); + assertEquals(shared.getReplicateToAt(0), "other-cluster"); + assertEquals(shared.getPropertyAt(0).getValue(), "value"); + assertEquals(shared.getEncryptionKeyAt(0).getMetadataAt(0).getValue(), "metadata-value"); + // a copy shares the same instance + EntryImpl copy = EntryImpl.create(entry); + assertSame(copy.getMessageMetadata(), shared); + copy.release(); + } finally { + entry.release(); + } + } + @Test public void testCreateWithLedgerIdEntryIdAndByteBuf() { // Given diff --git a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java index 708aff8c5a4aa..8e3cb415bd5cc 100644 --- a/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java +++ b/managed-ledger/src/test/java/org/apache/bookkeeper/mledger/impl/cache/RangeEntryCacheImplTest.java @@ -518,9 +518,7 @@ public void testReadFromStorageDoesNotShareSourceMetadataWithTheCopiedCacheEntry Entry sourceEntry = readEntries.get(0); // unlike the write path, the read path parses the metadata of the entry it returns before inserting it, - // so this is the case where insert receives an entry that already carries a MessageMetadata. Only the - // eagerly decoded sequenceId is read from it here, since reading a string field would decode it from - // the buffer and keep the decoded value, hiding the very problem this test is about + // so this is the case where insert receives an entry that already carries a MessageMetadata MessageMetadata sourceMetadata = sourceEntry.getMessageMetadata(); assertThat(sourceMetadata).isNotNull(); assertThat(sourceMetadata.getSequenceId()).isEqualTo(7); @@ -535,9 +533,9 @@ public void testReadFromStorageDoesNotShareSourceMetadataWithTheCopiedCacheEntry // overwritten once it has been recycled. This turns a leftover dependency on the source buffer into a // wrong value rather than into a read that only fails when the released memory happens to be reused headersAndPayload.setZero(headersAndPayload.readerIndex(), headersAndPayload.readableBytes()); - // metadata parsed from the source buffer does decode the overwritten bytes, which keeps the assertions - // below from turning vacuous should MessageMetadata ever stop decoding these fields lazily - assertThat(sourceMetadata.getProducerName()).isNotEqualTo("producer"); + // the metadata parsed from the source buffer was decoded when it was parsed, since copies of an entry share + // it across threads, so it doesn't depend on the overwritten bytes either + assertThat(sourceMetadata.getProducerName()).isEqualTo("producer"); // the read path releases the entries it returned once dispatch is done, while the cached copy stays sourceEntry.release(); @@ -546,8 +544,7 @@ public void testReadFromStorageDoesNotShareSourceMetadataWithTheCopiedCacheEntry ledgerEntry.close(); assertThat(headersAndPayload.refCnt()).isZero(); - // MessageMetadata decodes its string and bytes fields lazily from the buffer it was parsed from, so the - // cached entry stays readable only because its metadata is parsed from the buffer the cache owns + // the cached entry parses its own metadata from the buffer the cache owns Entry readBack = readSingleEntryFromCache(copyingCache, 1, 0); assertThat(readBack.getMessageMetadata()).isNotNull().isNotSameAs(sourceMetadata); readBack.release();