From 609f4687c77f6eb53a5b7f6b2e358f5c99ce1040 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 2 Aug 2026 23:39:55 -0700 Subject: [PATCH] [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.