Skip to content
Merged
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 @@ -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.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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();
Expand All @@ -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();
Expand Down
Loading