From e89e7b31403055460926ca473b925ec7cd6f4466 Mon Sep 17 00:00:00 2001 From: technoboy Date: Wed, 9 Mar 2022 20:33:22 +0800 Subject: [PATCH 1/5] Fix inconsistent prompt message when schema version is empty using AVRO. --- .../pulsar/broker/service/ServerCnx.java | 4 + .../org/apache/pulsar/schema/SchemaTest.java | 106 +++++++++++++++++- .../client/impl/BinaryProtoLookupService.java | 8 +- .../pulsar/client/impl/HttpLookupService.java | 5 + .../pulsar/client/impl/ProducerImpl.java | 2 +- 5 files changed, 119 insertions(+), 6 deletions(-) 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 83f80e2706476..c47c9427eae4e 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 @@ -2027,6 +2027,10 @@ remoteAddress, new String(commandGetSchema.getSchemaVersion()), long requestId = commandGetSchema.getRequestId(); SchemaVersion schemaVersion = SchemaVersion.Latest; if (commandGetSchema.hasSchemaVersion()) { + if (commandGetSchema.getSchemaVersion().length == 0) { + commandSender.sendGetSchemaErrorResponse(requestId, ServerError.IncompatibleSchema, "Empty schema version"); + return; + } schemaVersion = schemaService.versionFromBytes(commandGetSchema.getSchemaVersion()); } diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java index 3c6f7c73fbc9b..4ffe5e45b5fe7 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java @@ -29,12 +29,13 @@ import static org.testng.Assert.fail; import static org.testng.internal.junit.ArrayAsserts.assertArrayEquals; +import com.google.common.base.Throwables; +import lombok.EqualsAndHashCode; import org.apache.avro.Schema.Parser; - import com.fasterxml.jackson.databind.JsonNode; import com.google.common.collect.Sets; - import java.io.ByteArrayInputStream; +import java.io.Serializable; import java.nio.charset.StandardCharsets; import java.util.Collections; import java.util.HashMap; @@ -44,7 +45,6 @@ import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.stream.Collectors; - import lombok.Cleanup; import lombok.extern.slf4j.Slf4j; import org.apache.bookkeeper.client.BKException; @@ -59,7 +59,9 @@ import org.apache.pulsar.client.api.Producer; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.Schema; +import org.apache.pulsar.client.api.SchemaSerializationException; import org.apache.pulsar.client.api.SubscriptionInitialPosition; +import org.apache.pulsar.client.api.SubscriptionType; import org.apache.pulsar.client.api.TypedMessageBuilder; import org.apache.pulsar.client.api.schema.GenericRecord; import org.apache.pulsar.client.api.schema.SchemaDefinition; @@ -1112,4 +1114,102 @@ private void checkSchemaForAutoSchema(Message message) { } } + @Test + public void testAvroSchemaWithHttpLookup() throws Exception { + final String namespace = "test-namespace-" + randomName(16); + String ns = PUBLIC_TENANT + "/" + namespace; + admin.namespaces().createNamespace(ns, Sets.newHashSet(CLUSTER_NAME)); + final String autoProducerTopic = getTopicName(ns, "testAvroSchemaWithHttpLookup"); + + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.AVRO(User.class)) + .topic(autoProducerTopic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub-1") + .subscribe(); + + @Cleanup + Producer userProducer = pulsarClient + .newProducer(Schema.AVRO(User.class)) + .topic(autoProducerTopic) + .enableBatching(false) + .create(); + + @Cleanup + Producer producer = pulsarClient + .newProducer() + .topic(autoProducerTopic) + .enableBatching(false) + .create(); + User test = new User("test"); + userProducer.send(test); + producer.send("test".getBytes(StandardCharsets.UTF_8)); + Message message1 = consumer.receive(); + Assert.assertEquals(test, message1.getValue()); + try { + Message message2 = consumer.receive(); + message2.getValue(); + } catch (Throwable ex) { + Assert.assertTrue(Throwables.getRootCause(ex) instanceof SchemaSerializationException); + Assert.assertEquals(Throwables.getRootCause(ex).getMessage(),"Empty schema version"); + } + } + + @Test + public void testAvroSchemaWithTcpLookup() throws Exception { + stopBroker(); + isTcpLookup = true; + setup(); + final String namespace = "test-namespace-" + randomName(16); + String ns = PUBLIC_TENANT + "/" + namespace; + admin.namespaces().createNamespace(ns, Sets.newHashSet(CLUSTER_NAME)); + + final String autoProducerTopic = getTopicName(ns, "testAvroSchemaWithTcpLookup"); + + @Cleanup + Consumer consumer = pulsarClient + .newConsumer(Schema.AVRO(User.class)) + .topic(autoProducerTopic) + .subscriptionType(SubscriptionType.Shared) + .subscriptionName("sub-1") + .subscribe(); + + @Cleanup + Producer userProducer = pulsarClient + .newProducer(Schema.AVRO(User.class)) + .topic(autoProducerTopic) + .enableBatching(false) + .create(); + + @Cleanup + Producer producer = pulsarClient + .newProducer() + .topic(autoProducerTopic) + .enableBatching(false) + .create(); + + User test = new User("test"); + userProducer.send(test); + producer.send("test".getBytes(StandardCharsets.UTF_8)); + Message message1 = consumer.receive(); + Assert.assertEquals(test, message1.getValue()); + try { + Message message2 = consumer.receive(); + message2.getValue(); + } catch (Throwable ex) { + Assert.assertTrue(Throwables.getRootCause(ex) instanceof SchemaSerializationException); + Assert.assertEquals(Throwables.getRootCause(ex).getMessage(),"Empty schema version"); + } + } + + @EqualsAndHashCode + static class User implements Serializable { + private String name; + public User() {} + public User(String name) { + this.name = name; + } + } + } diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BinaryProtoLookupService.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BinaryProtoLookupService.java index 274c83c8c1fc8..190599dfffdbc 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BinaryProtoLookupService.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/BinaryProtoLookupService.java @@ -32,6 +32,7 @@ import java.util.concurrent.atomic.AtomicLong; import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.client.api.PulsarClientException; +import org.apache.pulsar.client.api.SchemaSerializationException; import org.apache.pulsar.common.api.proto.CommandGetTopicsOfNamespace.Mode; import org.apache.pulsar.common.api.proto.CommandLookupTopicResponse; import org.apache.pulsar.common.api.proto.CommandLookupTopicResponse.LookupType; @@ -222,9 +223,12 @@ public CompletableFuture> getSchema(TopicName topicName) { @Override public CompletableFuture> getSchema(TopicName topicName, byte[] version) { - InetSocketAddress socketAddress = serviceNameResolver.resolveHost(); CompletableFuture> schemaFuture = new CompletableFuture<>(); - + if (version != null && version.length == 0) { + schemaFuture.completeExceptionally(new SchemaSerializationException("Empty schema version")); + return schemaFuture; + } + InetSocketAddress socketAddress = serviceNameResolver.resolveHost(); client.getCnxPool().getConnection(socketAddress).thenAccept(clientCnx -> { long requestId = client.newRequestId(); ByteBuf request = Commands.newGetSchema(requestId, topicName.toString(), diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java index e05486cd95dd7..ced0aba9bba10 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/HttpLookupService.java @@ -34,6 +34,7 @@ import org.apache.commons.lang3.tuple.Pair; import org.apache.pulsar.client.api.PulsarClientException; import org.apache.pulsar.client.api.PulsarClientException.NotFoundException; +import org.apache.pulsar.client.api.SchemaSerializationException; import org.apache.pulsar.client.impl.conf.ClientConfigurationData; import org.apache.pulsar.client.impl.schema.SchemaInfoUtil; import org.apache.pulsar.client.impl.schema.SchemaUtils; @@ -159,6 +160,10 @@ public CompletableFuture> getSchema(TopicName topicName, by String schemaName = topicName.getSchemaName(); String path = String.format("admin/v2/schemas/%s/schema", schemaName); if (version != null) { + if (version.length == 0) { + future.completeExceptionally(new SchemaSerializationException("Empty schema version")); + return future; + } path = String.format("admin/v2/schemas/%s/schema/%s", schemaName, ByteBuffer.wrap(version).getLong()); diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index 67f6011eb6315..f60c69eaf8608 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -1562,7 +1562,7 @@ public void connectionOpened(final ClientCnx cnx) { schemaInfo = schema.getSchemaInfo(); } } else if (schema.getSchemaInfo().getType() == SchemaType.BYTES - || schema.getSchemaInfo().getType() == SchemaType.NONE) { + || schema.getSchemaInfo().getType() == SchemaType.NONE) { // don't set schema info for Schema.BYTES schemaInfo = null; } else { From 200b85ff7594648ba72ee01fdaa71deacb5795e6 Mon Sep 17 00:00:00 2001 From: technoboy Date: Wed, 9 Mar 2022 20:35:36 +0800 Subject: [PATCH 2/5] update --- .../main/java/org/apache/pulsar/client/impl/ProducerImpl.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java index f60c69eaf8608..67f6011eb6315 100644 --- a/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java +++ b/pulsar-client/src/main/java/org/apache/pulsar/client/impl/ProducerImpl.java @@ -1562,7 +1562,7 @@ public void connectionOpened(final ClientCnx cnx) { schemaInfo = schema.getSchemaInfo(); } } else if (schema.getSchemaInfo().getType() == SchemaType.BYTES - || schema.getSchemaInfo().getType() == SchemaType.NONE) { + || schema.getSchemaInfo().getType() == SchemaType.NONE) { // don't set schema info for Schema.BYTES schemaInfo = null; } else { From 820c17b343b32e10df7e568bbe65717c3bc67706 Mon Sep 17 00:00:00 2001 From: technoboy Date: Thu, 10 Mar 2022 10:21:39 +0800 Subject: [PATCH 3/5] fix checkstyle. --- .../main/java/org/apache/pulsar/broker/service/ServerCnx.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) 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 c47c9427eae4e..5f4d5f5c0e082 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 @@ -2028,7 +2028,8 @@ remoteAddress, new String(commandGetSchema.getSchemaVersion()), SchemaVersion schemaVersion = SchemaVersion.Latest; if (commandGetSchema.hasSchemaVersion()) { if (commandGetSchema.getSchemaVersion().length == 0) { - commandSender.sendGetSchemaErrorResponse(requestId, ServerError.IncompatibleSchema, "Empty schema version"); + commandSender.sendGetSchemaErrorResponse(requestId, ServerError.IncompatibleSchema, + "Empty schema version"); return; } schemaVersion = schemaService.versionFromBytes(commandGetSchema.getSchemaVersion()); From 168cf2daa2665b3de7f9466eb3aa7d6b05899e97 Mon Sep 17 00:00:00 2001 From: technoboy Date: Fri, 11 Mar 2022 12:55:32 +0800 Subject: [PATCH 4/5] update. --- .../src/test/java/org/apache/pulsar/schema/SchemaTest.java | 3 +++ 1 file changed, 3 insertions(+) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java index 4ffe5e45b5fe7..0a25b29038eb4 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java @@ -1116,6 +1116,9 @@ private void checkSchemaForAutoSchema(Message message) { @Test public void testAvroSchemaWithHttpLookup() throws Exception { + stopBroker(); + isTcpLookup = false; + setup(); final String namespace = "test-namespace-" + randomName(16); String ns = PUBLIC_TENANT + "/" + namespace; admin.namespaces().createNamespace(ns, Sets.newHashSet(CLUSTER_NAME)); From cd649de7206205976ce3edd650f012c664cb9cbc Mon Sep 17 00:00:00 2001 From: technoboy Date: Fri, 11 Mar 2022 15:34:01 +0800 Subject: [PATCH 5/5] combine test. --- .../org/apache/pulsar/schema/SchemaTest.java | 45 +++---------------- 1 file changed, 6 insertions(+), 39 deletions(-) diff --git a/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java b/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java index 0a25b29038eb4..217e76820f15c 100644 --- a/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java +++ b/pulsar-broker/src/test/java/org/apache/pulsar/schema/SchemaTest.java @@ -1119,44 +1119,7 @@ public void testAvroSchemaWithHttpLookup() throws Exception { stopBroker(); isTcpLookup = false; setup(); - final String namespace = "test-namespace-" + randomName(16); - String ns = PUBLIC_TENANT + "/" + namespace; - admin.namespaces().createNamespace(ns, Sets.newHashSet(CLUSTER_NAME)); - final String autoProducerTopic = getTopicName(ns, "testAvroSchemaWithHttpLookup"); - - @Cleanup - Consumer consumer = pulsarClient - .newConsumer(Schema.AVRO(User.class)) - .topic(autoProducerTopic) - .subscriptionType(SubscriptionType.Shared) - .subscriptionName("sub-1") - .subscribe(); - - @Cleanup - Producer userProducer = pulsarClient - .newProducer(Schema.AVRO(User.class)) - .topic(autoProducerTopic) - .enableBatching(false) - .create(); - - @Cleanup - Producer producer = pulsarClient - .newProducer() - .topic(autoProducerTopic) - .enableBatching(false) - .create(); - User test = new User("test"); - userProducer.send(test); - producer.send("test".getBytes(StandardCharsets.UTF_8)); - Message message1 = consumer.receive(); - Assert.assertEquals(test, message1.getValue()); - try { - Message message2 = consumer.receive(); - message2.getValue(); - } catch (Throwable ex) { - Assert.assertTrue(Throwables.getRootCause(ex) instanceof SchemaSerializationException); - Assert.assertEquals(Throwables.getRootCause(ex).getMessage(),"Empty schema version"); - } + testEmptySchema(); } @Test @@ -1164,11 +1127,15 @@ public void testAvroSchemaWithTcpLookup() throws Exception { stopBroker(); isTcpLookup = true; setup(); + testEmptySchema(); + } + + private void testEmptySchema() throws Exception { final String namespace = "test-namespace-" + randomName(16); String ns = PUBLIC_TENANT + "/" + namespace; admin.namespaces().createNamespace(ns, Sets.newHashSet(CLUSTER_NAME)); - final String autoProducerTopic = getTopicName(ns, "testAvroSchemaWithTcpLookup"); + final String autoProducerTopic = getTopicName(ns, "testEmptySchema"); @Cleanup Consumer consumer = pulsarClient