From 5275b527cf3c66961d412580cd1cc5d703cbdad0 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Tue, 13 Dec 2022 01:34:53 +0100 Subject: [PATCH 1/3] [fix][misc] do not require encryption on system topics --- .../pulsar/broker/service/Producer.java | 4 ++- .../pulsar/broker/service/ServerCnx.java | 5 +++- .../service/persistent/SystemTopic.java | 6 ++++ .../broker/admin/TopicPoliciesTest.java | 14 ++++++++++ .../broker/transaction/TransactionTest.java | 28 +++++++++++++++++++ 5 files changed, 55 insertions(+), 2 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index e7034391b1461..3dbf503a282a5 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -800,7 +800,9 @@ public void checkEncryption() { public void publishTxnMessage(TxnID txnID, long producerId, long sequenceId, long highSequenceId, ByteBuf headersAndPayload, long batchSize, boolean isChunked, boolean isMarker) { - checkAndStartPublish(producerId, sequenceId, headersAndPayload, batchSize, null); + if (!checkAndStartPublish(producerId, sequenceId, headersAndPayload, batchSize, null)) { + // return; + } MessagePublishContext messagePublishContext = MessagePublishContext.get(this, sequenceId, highSequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, isChunked, System.nanoTime(), isMarker, null); diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java index 15842f4598549..f6764b5b69f02 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/ServerCnx.java @@ -142,6 +142,7 @@ import org.apache.pulsar.common.intercept.InterceptException; import org.apache.pulsar.common.naming.Metadata; import org.apache.pulsar.common.naming.NamespaceName; +import org.apache.pulsar.common.naming.SystemTopicNames; import org.apache.pulsar.common.naming.TopicName; import org.apache.pulsar.common.policies.data.BacklogQuota; import org.apache.pulsar.common.policies.data.BacklogQuota.BacklogQuotaType; @@ -1326,7 +1327,9 @@ protected void handleProducer(final CommandProducer cmdProducer) { backlogQuotaCheckFuture.thenRun(() -> { // Check whether the producer will publish encrypted messages or not - if ((topic.isEncryptionRequired() || encryptionRequireOnProducer) && !isEncrypted) { + if ((topic.isEncryptionRequired() || encryptionRequireOnProducer) + && !isEncrypted + && !SystemTopicNames.isSystemTopic(topicName)) { String msg = String.format("Encryption is required in %s", topicName); log.warn("[{}] {}", remoteAddress, msg); if (producerFuture.completeExceptionally(new ServerMetadataException(msg))) { diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/SystemTopic.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/SystemTopic.java index 4a9d9b0b2d5bc..395a8c9075eb3 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/SystemTopic.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/SystemTopic.java @@ -76,4 +76,10 @@ public boolean isCompactionEnabled() { // even though is not explicitly set in the policies. return !NamespaceService.isHeartbeatNamespace(TopicName.get(topic)); } + + @Override + public boolean isEncryptionRequired() { + // System topics are only written by the broker that can't know the encryption context. + return false; + } } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java index 94cd299e47542..db26f250067fe 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/TopicPoliciesTest.java @@ -20,6 +20,7 @@ import static org.testng.Assert.assertEquals; import static org.testng.Assert.assertFalse; +import static org.testng.Assert.assertNotEquals; import static org.testng.Assert.assertNotNull; import static org.testng.Assert.assertNull; import static org.testng.Assert.assertTrue; @@ -3109,4 +3110,17 @@ public void testGetTopicPoliciesWhenDeleteTopicPolicy() throws Exception { assertNull(topicPolicies); } + @Test + public void testProduceChangesWithEncryptionRequired() throws Exception { + final String beforeLac = admin.topics().getInternalStats(topicPolicyEventsTopic).lastConfirmedEntry; + admin.namespaces().setEncryptionRequiredStatus(myNamespace, true); + // just an update to trigger writes on __change_events + admin.topicPolicies().setMaxConsumers(testTopic, 5); + Awaitility.await() + .untilAsserted(() -> { + final PersistentTopicInternalStats newLac = admin.topics().getInternalStats(topicPolicyEventsTopic); + assertNotEquals(newLac, beforeLac); + }); + } + } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index f075b80ef7aba..3e3c71b20786f 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -115,13 +115,16 @@ import org.apache.pulsar.client.api.ReaderBuilder; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; +import org.apache.pulsar.client.api.interceptor.ProducerInterceptor; import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.client.impl.ClientCnx; import org.apache.pulsar.client.impl.ConsumerBase; +import org.apache.pulsar.client.impl.DefaultCryptoKeyReader; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.MessagesImpl; import org.apache.pulsar.client.util.ExecutorProvider; +import org.apache.pulsar.common.api.EncryptionContext; import org.apache.pulsar.common.api.proto.CommandSubscribe; import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.common.events.EventType; @@ -1617,4 +1620,29 @@ public void testGetTxnState() throws Exception { Transaction abortingTxn = transaction; Awaitility.await().until(() -> abortingTxn.getState() == Transaction.State.ABORTING); } + + + @Test + public void testEncryptionRequired() throws Exception { + final String namespace = "tnx/ns-prechecks"; + final String topic = "persistent://" + namespace + "/test_transaction_topic"; + admin.namespaces().createNamespace(namespace); + admin.namespaces().setEncryptionRequiredStatus(namespace, true); + admin.topics().createNonPartitionedTopic(topic); + + @Cleanup + Producer producer = this.pulsarClient.newProducer() + .topic(topic) + .sendTimeout(5, TimeUnit.SECONDS) + .addEncryptionKey("my-app-key") + .defaultCryptoKeyReader("file:./src/test/resources/certificate/public-key.client-rsa.pem") + .create(); + + Transaction txn = pulsarClient.newTransaction() + .withTransactionTimeout(5, TimeUnit.SECONDS).build().get(); + producer.newMessage(txn) + .value(UUID.randomUUID().toString().getBytes(StandardCharsets.UTF_8)) + .send(); + txn.commit(); + } } From 65138726b9c8f381b8c06fbe1d00340f8cfb2e8f Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Tue, 13 Dec 2022 09:19:28 +0100 Subject: [PATCH 2/3] style and fix return --- .../main/java/org/apache/pulsar/broker/service/Producer.java | 2 +- .../org/apache/pulsar/broker/transaction/TransactionTest.java | 3 --- 2 files changed, 1 insertion(+), 4 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index 3dbf503a282a5..bc101e31d2723 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -801,7 +801,7 @@ public void checkEncryption() { public void publishTxnMessage(TxnID txnID, long producerId, long sequenceId, long highSequenceId, ByteBuf headersAndPayload, long batchSize, boolean isChunked, boolean isMarker) { if (!checkAndStartPublish(producerId, sequenceId, headersAndPayload, batchSize, null)) { - // return; + return; } MessagePublishContext messagePublishContext = MessagePublishContext.get(this, sequenceId, highSequenceId, msgIn, diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java index 3e3c71b20786f..c237b5024c433 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/broker/transaction/TransactionTest.java @@ -115,16 +115,13 @@ import org.apache.pulsar.client.api.ReaderBuilder; import org.apache.pulsar.client.api.Schema; import org.apache.pulsar.client.api.SubscriptionType; -import org.apache.pulsar.client.api.interceptor.ProducerInterceptor; import org.apache.pulsar.client.api.transaction.Transaction; import org.apache.pulsar.client.api.transaction.TxnID; import org.apache.pulsar.client.impl.ClientCnx; import org.apache.pulsar.client.impl.ConsumerBase; -import org.apache.pulsar.client.impl.DefaultCryptoKeyReader; import org.apache.pulsar.client.impl.MessageIdImpl; import org.apache.pulsar.client.impl.MessagesImpl; import org.apache.pulsar.client.util.ExecutorProvider; -import org.apache.pulsar.common.api.EncryptionContext; import org.apache.pulsar.common.api.proto.CommandSubscribe; import org.apache.pulsar.client.impl.transaction.TransactionImpl; import org.apache.pulsar.common.events.EventType; From e336ff6558c3dbaa006f5adf98538cbbd153bdee Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Nicol=C3=B2=20Boschi?= Date: Tue, 13 Dec 2022 20:32:38 +0100 Subject: [PATCH 3/3] rm unrelated change --- .../main/java/org/apache/pulsar/broker/service/Producer.java | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java index bc101e31d2723..e7034391b1461 100644 --- a/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java +++ b/pulsar-broker/src/main/java/org/apache/pulsar/broker/service/Producer.java @@ -800,9 +800,7 @@ public void checkEncryption() { public void publishTxnMessage(TxnID txnID, long producerId, long sequenceId, long highSequenceId, ByteBuf headersAndPayload, long batchSize, boolean isChunked, boolean isMarker) { - if (!checkAndStartPublish(producerId, sequenceId, headersAndPayload, batchSize, null)) { - return; - } + checkAndStartPublish(producerId, sequenceId, headersAndPayload, batchSize, null); MessagePublishContext messagePublishContext = MessagePublishContext.get(this, sequenceId, highSequenceId, msgIn, headersAndPayload.readableBytes(), batchSize, isChunked, System.nanoTime(), isMarker, null);