From f99d6994942cba4315e3db30669a07cd4aaa78fe Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 2 Aug 2026 23:39:55 -0700 Subject: [PATCH 1/2] [ISSUE #10770] Skip queryAssignment queues without master broker --- .../proxy/grpc/v2/route/RouteActivity.java | 3 ++ .../grpc/v2/route/RouteActivityTest.java | 29 +++++++++++++++++++ 2 files changed, 32 insertions(+) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java index 75f7089c5e0..d5f70728e7a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java @@ -125,6 +125,9 @@ public CompletableFuture queryAssignment(ProxyContext c Map brokerIdMap = brokerMap.get(queueData.getBrokerName()); if (brokerIdMap != null) { Broker broker = brokerIdMap.get(MixAll.MASTER_ID); + if (broker == null) { + continue; + } Permission permission = this.convertToPermission(queueData.getPerm()); if (isFifo && !isLite) { for (int i = 0; i < queueData.getReadQueueNums(); i++) { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java index abbf82452ef..4447a5d0692 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java @@ -173,6 +173,24 @@ public void testQueryAssignmentWithNoReadQueue() throws Throwable { assertEquals(Code.FORBIDDEN, response.getStatus().getCode()); } + @Test + public void testQueryAssignmentWithMissingMasterBroker() throws Throwable { + when(this.messagingProcessor.getTopicRouteDataForProxy(any(), any(), anyString())) + .thenReturn(createProxyTopicRouteDataWithoutMasterBroker()); + + QueryAssignmentResponse response = this.routeActivity.queryAssignment( + createContext(), + QueryAssignmentRequest.newBuilder() + .setEndpoints(grpcEndpoints) + .setTopic(GRPC_TOPIC) + .setGroup(GRPC_GROUP) + .build() + ).get(); + + assertEquals(Code.FORBIDDEN, response.getStatus().getCode()); + assertEquals(0, response.getAssignmentsCount()); + } + @Test public void testQueryAssignment() throws Throwable { when(this.messagingProcessor.getTopicRouteDataForProxy(any(), any(), anyString())) @@ -226,6 +244,17 @@ private static ProxyTopicRouteData createProxyTopicRouteData(int r, int w, int p return proxyTopicRouteData; } + private static ProxyTopicRouteData createProxyTopicRouteDataWithoutMasterBroker() { + ProxyTopicRouteData proxyTopicRouteData = new ProxyTopicRouteData(); + proxyTopicRouteData.getQueueDatas().add(createQueueData(2, 2, PermName.PERM_READ | PermName.PERM_WRITE)); + ProxyTopicRouteData.ProxyBrokerData proxyBrokerData = new ProxyTopicRouteData.ProxyBrokerData(); + proxyBrokerData.setCluster(CLUSTER); + proxyBrokerData.setBrokerName(BROKER_NAME); + proxyBrokerData.getBrokerAddrs().put(1L, addressArrayList); + proxyTopicRouteData.getBrokerDatas().add(proxyBrokerData); + return proxyTopicRouteData; + } + @Test public void testGenPartitionFromQueueData() throws Exception { // test queueData with 8 read queues, 8 write queues, and rw permission, expect 8 rw queues. From 2530042c7992aba18f60b66b2c4f01c2abae941e Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 2 Aug 2026 23:29:10 -0700 Subject: [PATCH 2/2] [ISSUE #10768] Continue queryRoute after missing broker data --- .../proxy/grpc/v2/route/RouteActivity.java | 2 +- .../grpc/v2/route/RouteActivityTest.java | 31 +++++++++++++++++++ 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java index d5f70728e7a..cefcecdc1dc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivity.java @@ -78,7 +78,7 @@ public CompletableFuture queryRoute(ProxyContext ctx, QueryR String brokerName = queueData.getBrokerName(); Map brokerIdMap = brokerMap.get(brokerName); if (brokerIdMap == null) { - break; + continue; } for (Broker broker : brokerIdMap.values()) { messageQueueList.addAll(this.genMessageQueueFromQueueData(queueData, request.getTopic(), topicMessageType, broker)); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java index 4447a5d0692..ca123e8200d 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/route/RouteActivityTest.java @@ -119,6 +119,29 @@ public void testQueryRoute() throws Throwable { } } + @Test + public void testQueryRouteSkipsMissingBrokerAndContinues() throws Throwable { + when(this.messagingProcessor.getTopicRouteDataForProxy(any(), any(), anyString())) + .thenReturn(createProxyTopicRouteDataWithMissingBrokerBeforeValidBroker()); + when(this.messagingProcessor.getMetadataService().getTopicMessageType(any(), anyString())) + .thenReturn(TopicMessageType.NORMAL); + + QueryRouteResponse response = this.routeActivity.queryRoute( + createContext(), + QueryRouteRequest.newBuilder() + .setEndpoints(grpcEndpoints) + .setTopic(Resource.newBuilder().setName(TOPIC).build()) + .build() + ).get(); + + assertEquals(Code.OK, response.getStatus().getCode()); + assertEquals(4, response.getMessageQueuesCount()); + for (MessageQueue messageQueue : response.getMessageQueuesList()) { + assertEquals(BROKER_NAME, messageQueue.getBroker().getName()); + assertEquals(grpcEndpoints, messageQueue.getBroker().getEndpoints()); + } + } + @Test public void testQueryRouteTopicExist() throws Throwable { when(this.messagingProcessor.getTopicRouteDataForProxy(any(), any(), anyString())) @@ -255,6 +278,14 @@ private static ProxyTopicRouteData createProxyTopicRouteDataWithoutMasterBroker( return proxyTopicRouteData; } + private static ProxyTopicRouteData createProxyTopicRouteDataWithMissingBrokerBeforeValidBroker() { + ProxyTopicRouteData proxyTopicRouteData = createProxyTopicRouteData(2, 2, PermName.PERM_READ | PermName.PERM_WRITE); + QueueData missingBrokerQueueData = createQueueData(2, 2, PermName.PERM_READ | PermName.PERM_WRITE); + missingBrokerQueueData.setBrokerName("missingBrokerName"); + proxyTopicRouteData.getQueueDatas().add(0, missingBrokerQueueData); + return proxyTopicRouteData; + } + @Test public void testGenPartitionFromQueueData() throws Exception { // test queueData with 8 read queues, 8 write queues, and rw permission, expect 8 rw queues.