From 2fdf4c80897604b34db54dad4186d17aedeabcd1 Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Tue, 4 Jul 2023 20:26:24 +0800 Subject: [PATCH 1/4] [fix][test]Fix resource not close after method --- .../pulsar/broker/service/PublishRateLimiterTest.java | 9 +++++++++ 1 file changed, 9 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PublishRateLimiterTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PublishRateLimiterTest.java index b934ced08c5db..f3cb25e789f08 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PublishRateLimiterTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/service/PublishRateLimiterTest.java @@ -21,6 +21,7 @@ import org.apache.pulsar.common.policies.data.Policies; import org.apache.pulsar.common.policies.data.PublishRate; import org.apache.pulsar.common.util.RateLimiter; +import org.testng.annotations.AfterMethod; import org.testng.annotations.BeforeMethod; import org.testng.annotations.Test; @@ -51,6 +52,14 @@ public void setup() throws Exception { publishRateLimiter = new PublishRateLimiterImpl(policies, CLUSTER_NAME); } + @AfterMethod + public void cleanup() throws Exception { + policies.publishMaxMessageRate.clear(); + policies.publishMaxMessageRate = null; + precisePublishLimiter.close(); + publishRateLimiter.close(); + } + @Test public void testPublishRateLimiterImplExceed() throws Exception { // increment not exceed From 33e9158a54fa6e703103781b757233dd9a47743b Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Tue, 2 Jan 2024 20:18:16 +0800 Subject: [PATCH 2/4] [fix][broker] Fix messages could not expire due to incorrect client clock --- .../apache/pulsar/common/protocol/Commands.java | 3 +++ .../pulsar/common/protocol/CommandUtilsTests.java | 14 ++++++++++++++ 2 files changed, 17 insertions(+) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index 34d47e2836bb2..f0f55ade799bf 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -1704,6 +1704,9 @@ public static ByteBuf addBrokerEntryMetadata(ByteBuf headerAndPayload, interceptor.interceptWithNumberOfMessages(brokerEntryMetadata, numberOfMessages); } } + if (!brokerEntryMetadata.hasBrokerTimestamp()) { + brokerEntryMetadata.setBrokerTimestamp(System.currentTimeMillis()); + } int brokerMetaSize = brokerEntryMetadata.getSerializedSize(); ByteBuf brokerMeta = diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java index 4523a5cc97b37..0c06ec01f8d44 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java @@ -20,6 +20,7 @@ import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; import io.netty.buffer.ByteBuf; @@ -151,6 +152,19 @@ public void testAddBrokerEntryMetadata() throws Exception { assertTrue(new String(content, StandardCharsets.UTF_8).endsWith(data)); } + @Test + public void testAlwaysSetBrokerTimestamp() { + int MOCK_BATCH_SIZE = 10; + String data = "test-message"; + ByteBuf byteBuf = PulsarByteBufAllocator.DEFAULT.buffer(data.length(), data.length()); + byteBuf.writeBytes(data.getBytes(StandardCharsets.UTF_8)); + ByteBuf dataWithBrokerEntryMetadata = + Commands.addBrokerEntryMetadata(byteBuf, new HashSet<>(), MOCK_BATCH_SIZE); + BrokerEntryMetadata entryMetadata = Commands.parseBrokerEntryMetadataIfExist(dataWithBrokerEntryMetadata); + assertNotNull(entryMetadata); + assertTrue(entryMetadata.hasBrokerTimestamp()); + } + @Test public void testSkipBrokerEntryMetadata() throws Exception { String data = "test-message"; From b9bf55a74a2bb7c082603d60922ae1117eef46b7 Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Wed, 3 Jan 2024 17:29:41 +0800 Subject: [PATCH 3/4] [fix][broker] Fix messages could not expire due to incorrect client clock --- .../java/org/apache/pulsar/common/protocol/Commands.java | 4 ++-- .../apache/pulsar/common/protocol/CommandUtilsTests.java | 7 +++++++ 2 files changed, 9 insertions(+), 2 deletions(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index f0f55ade799bf..7c59bff717328 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -470,7 +470,7 @@ public static void skipMessageMetadata(ByteBuf buffer) { } public static long getEntryTimestamp(ByteBuf headersAndPayloadWithBrokerEntryMetadata) throws IOException { - // get broker timestamp first if BrokerEntryMetadata is enabled with AppendBrokerTimestampMetadataInterceptor + // get broker timestamp first if exists BrokerEntryMetadata brokerEntryMetadata = Commands.parseBrokerEntryMetadataIfExist(headersAndPayloadWithBrokerEntryMetadata); if (brokerEntryMetadata != null && brokerEntryMetadata.hasBrokerTimestamp()) { @@ -1697,7 +1697,7 @@ public static ByteBuf addBrokerEntryMetadata(ByteBuf headerAndPayload, // | BROKER_ENTRY_METADATA_MAGIC_NUMBER | BROKER_ENTRY_METADATA_SIZE | BROKER_ENTRY_METADATA | // | 2 bytes | 4 bytes | BROKER_ENTRY_METADATA_SIZE bytes | - BrokerEntryMetadata brokerEntryMetadata = BROKER_ENTRY_METADATA.get(); + BrokerEntryMetadata brokerEntryMetadata = BROKER_ENTRY_METADATA.get().clearBrokerTimestamp(); for (BrokerEntryMetadataInterceptor interceptor : brokerInterceptors) { interceptor.intercept(brokerEntryMetadata); if (numberOfMessages >= 0) { diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java index 0c06ec01f8d44..e306149259c0a 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java @@ -20,6 +20,7 @@ import static org.apache.pulsar.common.protocol.Commands.serializeMetadataAndPayload; import static org.testng.Assert.assertEquals; +import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertTrue; @@ -163,6 +164,12 @@ public void testAlwaysSetBrokerTimestamp() { BrokerEntryMetadata entryMetadata = Commands.parseBrokerEntryMetadataIfExist(dataWithBrokerEntryMetadata); assertNotNull(entryMetadata); assertTrue(entryMetadata.hasBrokerTimestamp()); + long timestamp = entryMetadata.getBrokerTimestamp(); + // add again, the timestamp should change + dataWithBrokerEntryMetadata = + Commands.addBrokerEntryMetadata(byteBuf, new HashSet<>(), MOCK_BATCH_SIZE); + entryMetadata = Commands.parseBrokerEntryMetadataIfExist(dataWithBrokerEntryMetadata); + assertNotEquals(timestamp, entryMetadata.getBrokerTimestamp()); } @Test From 523ac7679cb4e31d9c9e22f6edf83e3b853b55c7 Mon Sep 17 00:00:00 2001 From: feynmanlin <315157973@qq.com> Date: Thu, 4 Jan 2024 14:20:06 +0800 Subject: [PATCH 4/4] Address comment --- .../main/java/org/apache/pulsar/common/protocol/Commands.java | 2 +- .../org/apache/pulsar/common/protocol/CommandUtilsTests.java | 1 + 2 files changed, 2 insertions(+), 1 deletion(-) diff --git a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java index 7c59bff717328..c40923594e0ff 100644 --- a/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java +++ b/pulsar-common/src/main/java/org/apache/pulsar/common/protocol/Commands.java @@ -1697,7 +1697,7 @@ public static ByteBuf addBrokerEntryMetadata(ByteBuf headerAndPayload, // | BROKER_ENTRY_METADATA_MAGIC_NUMBER | BROKER_ENTRY_METADATA_SIZE | BROKER_ENTRY_METADATA | // | 2 bytes | 4 bytes | BROKER_ENTRY_METADATA_SIZE bytes | - BrokerEntryMetadata brokerEntryMetadata = BROKER_ENTRY_METADATA.get().clearBrokerTimestamp(); + BrokerEntryMetadata brokerEntryMetadata = BROKER_ENTRY_METADATA.get().clear(); for (BrokerEntryMetadataInterceptor interceptor : brokerInterceptors) { interceptor.intercept(brokerEntryMetadata); if (numberOfMessages >= 0) { diff --git a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java index e306149259c0a..f508eb3d216eb 100644 --- a/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java +++ b/pulsar-common/src/test/java/org/apache/pulsar/common/protocol/CommandUtilsTests.java @@ -169,6 +169,7 @@ public void testAlwaysSetBrokerTimestamp() { dataWithBrokerEntryMetadata = Commands.addBrokerEntryMetadata(byteBuf, new HashSet<>(), MOCK_BATCH_SIZE); entryMetadata = Commands.parseBrokerEntryMetadataIfExist(dataWithBrokerEntryMetadata); + assertNotNull(entryMetadata); assertNotEquals(timestamp, entryMetadata.getBrokerTimestamp()); }