From 0a60783da322fa1fb749349b7751c7a9d742c4d7 Mon Sep 17 00:00:00 2001 From: Jiwei Guo Date: Sat, 22 Aug 2026 10:53:38 +0800 Subject: [PATCH 1/4] Fix test --- .../amqp/rabbitmq/functional/ExchangeDeclareTest.java | 10 +++++----- 1 file changed, 5 insertions(+), 5 deletions(-) diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java index 57658254c..07013a66b 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java @@ -16,8 +16,6 @@ package io.streamnative.pulsar.handlers.amqp.rabbitmq.functional; -import static org.junit.Assert.assertEquals; - import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; @@ -118,9 +116,11 @@ public void exchangeDeclaredWithEnumerationEquivalentOnRecoverableConnection() private void doTestExchangeDeclaredWithEnumerationEquivalent(Channel channel) throws IOException, InterruptedException { - assertEquals("There are 4 standard exchange types", - 4, BuiltinExchangeType.values().length); - for (BuiltinExchangeType exchangeType : BuiltinExchangeType.values()) { + for (BuiltinExchangeType exchangeType : new BuiltinExchangeType[]{ + BuiltinExchangeType.DIRECT, + BuiltinExchangeType.FANOUT, + BuiltinExchangeType.TOPIC, + BuiltinExchangeType.HEADERS}) { channel.exchangeDeclare(NAME, exchangeType); verifyEquivalent(NAME, exchangeType.getType(), false, false, null); deleteExchangeWithRetry(); From f45a6eb704ee0e3cf2ccf8bc95ffdb04496d3d23 Mon Sep 17 00:00:00 2001 From: Jiwei Guo Date: Sat, 22 Aug 2026 11:02:38 +0800 Subject: [PATCH 2/4] Fix upgrade rabbitmq version cause BuiltinExchangeType has 4 new types. --- .../pulsar/handlers/amqp/AmqpExchange.java | 31 ++++++++++++++++--- .../handlers/amqp/ExchangeMessageRouter.java | 1 + .../functional/ExchangeDeclareTest.java | 9 +++--- 3 files changed, 32 insertions(+), 9 deletions(-) diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpExchange.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpExchange.java index 31c7eebc4..f9ed25706 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpExchange.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/AmqpExchange.java @@ -34,10 +34,20 @@ public interface AmqpExchange { * */ enum Type{ - Direct, - Fanout, - Topic, - Headers; + Direct("direct"), + Fanout("fanout"), + Topic("topic"), + Headers("headers"), + ConsistentHash("x-consistent-hash"), + ModulusHash("x-modulus-hash"), + LocalRandom("x-local-random"), + Random("x-random"); + + private final String type; + + Type(String type) { + this.type = type; + } public static Type value(String type) { if (type == null || type.length() == 0) { @@ -53,11 +63,24 @@ public static Type value(String type) { return Topic; case "headers": return Headers; + case "x-consistent-hash": + return ConsistentHash; + case "x-modulus-hash": + return ModulusHash; + case "x-local-random": + return LocalRandom; + case "x-random": + return Random; default: return null; } } + @Override + public String toString() { + return type; + } + } /** diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java index 028adda5a..3586a774c 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java @@ -301,6 +301,7 @@ public static ExchangeMessageRouter getInstance(PersistentExchange exchange, Exe case Direct -> new DirectExchangeMessageRouter(exchange, routeExecutor); case Topic -> new TopicExchangeMessageRouter(exchange, routeExecutor); case Headers -> new HeadersExchangeMessageRouter(exchange, routeExecutor); + default -> null; }; } diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java index 07013a66b..24ad0a7f5 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java @@ -16,6 +16,7 @@ package io.streamnative.pulsar.handlers.amqp.rabbitmq.functional; +import static org.junit.Assert.assertEquals; import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; @@ -116,11 +117,9 @@ public void exchangeDeclaredWithEnumerationEquivalentOnRecoverableConnection() private void doTestExchangeDeclaredWithEnumerationEquivalent(Channel channel) throws IOException, InterruptedException { - for (BuiltinExchangeType exchangeType : new BuiltinExchangeType[]{ - BuiltinExchangeType.DIRECT, - BuiltinExchangeType.FANOUT, - BuiltinExchangeType.TOPIC, - BuiltinExchangeType.HEADERS}) { + assertEquals("There are 8 standard exchange types", + 8, BuiltinExchangeType.values().length); + for (BuiltinExchangeType exchangeType : BuiltinExchangeType.values()) { channel.exchangeDeclare(NAME, exchangeType); verifyEquivalent(NAME, exchangeType.getType(), false, false, null); deleteExchangeWithRetry(); From fa346674d78c99c04859591a9c2e031a09dacc88 Mon Sep 17 00:00:00 2001 From: Jiwei Guo Date: Sat, 22 Aug 2026 11:08:55 +0800 Subject: [PATCH 3/4] fix --- .../handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java b/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java index 24ad0a7f5..708d5c5c4 100644 --- a/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java +++ b/tests/src/test/java/io/streamnative/pulsar/handlers/amqp/rabbitmq/functional/ExchangeDeclareTest.java @@ -17,6 +17,7 @@ package io.streamnative.pulsar.handlers.amqp.rabbitmq.functional; import static org.junit.Assert.assertEquals; + import com.rabbitmq.client.BuiltinExchangeType; import com.rabbitmq.client.Channel; import com.rabbitmq.client.Connection; From 49ba88346d56e10866f1c9e72bb4f257c2074bdb Mon Sep 17 00:00:00 2001 From: "gaoran_10@126.com" Date: Sat, 22 Aug 2026 12:19:43 +0800 Subject: [PATCH 4/4] add unsupported exception --- .../handlers/amqp/ExchangeMessageRouter.java | 3 +- .../amqp/test/ExchangeMessageRouterTest.java | 51 +++++++++++++++++++ 2 files changed, 53 insertions(+), 1 deletion(-) create mode 100644 amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/ExchangeMessageRouterTest.java diff --git a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java index 3586a774c..c10cc3fb6 100644 --- a/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java +++ b/amqp-impl/src/main/java/io/streamnative/pulsar/handlers/amqp/ExchangeMessageRouter.java @@ -301,7 +301,8 @@ public static ExchangeMessageRouter getInstance(PersistentExchange exchange, Exe case Direct -> new DirectExchangeMessageRouter(exchange, routeExecutor); case Topic -> new TopicExchangeMessageRouter(exchange, routeExecutor); case Headers -> new HeadersExchangeMessageRouter(exchange, routeExecutor); - default -> null; + default -> throw new AoPServiceRuntimeException.NotSupportedOperationException( + "Exchange router is not supported for type " + exchange.getType() + "."); }; } diff --git a/amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/ExchangeMessageRouterTest.java b/amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/ExchangeMessageRouterTest.java new file mode 100644 index 000000000..a084e9472 --- /dev/null +++ b/amqp-impl/src/test/java/io/streamnative/pulsar/handlers/amqp/test/ExchangeMessageRouterTest.java @@ -0,0 +1,51 @@ +/** + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package io.streamnative.pulsar.handlers.amqp.test; + +import io.streamnative.pulsar.handlers.amqp.AmqpExchange; +import io.streamnative.pulsar.handlers.amqp.ExchangeMessageRouter; +import io.streamnative.pulsar.handlers.amqp.common.exception.AoPServiceRuntimeException.NotSupportedOperationException; +import io.streamnative.pulsar.handlers.amqp.impl.PersistentExchange; +import java.util.concurrent.ExecutorService; +import org.mockito.Mockito; +import org.testng.Assert; +import org.testng.annotations.DataProvider; +import org.testng.annotations.Test; + +/** + * Unit tests for exchange router creation. + */ +public class ExchangeMessageRouterTest { + + @DataProvider(name = "unsupportedExchangeTypes") + public Object[][] unsupportedExchangeTypes() { + return new Object[][]{ + {AmqpExchange.Type.ConsistentHash}, + {AmqpExchange.Type.ModulusHash}, + {AmqpExchange.Type.LocalRandom}, + {AmqpExchange.Type.Random} + }; + } + + @Test(dataProvider = "unsupportedExchangeTypes") + public void shouldThrowWhenExchangeTypeHasNoRouter(AmqpExchange.Type exchangeType) { + PersistentExchange exchange = Mockito.mock(PersistentExchange.class); + ExecutorService routeExecutor = Mockito.mock(ExecutorService.class); + Mockito.when(exchange.getType()).thenReturn(exchangeType); + + NotSupportedOperationException exception = Assert.expectThrows(NotSupportedOperationException.class, + () -> ExchangeMessageRouter.getInstance(exchange, routeExecutor)); + Assert.assertTrue(exception.getMessage().contains(exchangeType.toString())); + } +}