From f45cf33972c2d96d5bb284b8a27e74a2b4186bee Mon Sep 17 00:00:00 2001 From: zjncs <18910855655@163.com> Date: Fri, 11 Sep 2026 13:59:26 +0800 Subject: [PATCH] fix(proxy): return UNRECOGNIZED_CLIENT_TYPE when client settings is null in receiveMessage A gRPC client that calls ReceiveMessage before any settings were cached (or after its settings were cleaned up) made getClientSettings return null, and the following settings.getClientType() threw a NullPointerException which the catch block turned into an opaque INTERNAL_SERVER_ERROR response. ClientActivity.heartbeat already guards this case with UNRECOGNIZED_CLIENT_TYPE; apply the same guard in receiveMessage so the client gets a meaningful code. --- .../v2/consumer/ReceiveMessageActivity.java | 4 +++ .../consumer/ReceiveMessageActivityTest.java | 25 +++++++++++++++++++ 2 files changed, 29 insertions(+) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java index 39f2995d6fc..d5cdc93f2f6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java @@ -64,6 +64,10 @@ public void receiveMessage(ProxyContext ctx, ReceiveMessageRequest request, try { Settings settings = this.grpcClientSettingsManager.getClientSettings(ctx); + if (settings == null) { + writer.writeAndComplete(ctx, Code.UNRECOGNIZED_CLIENT_TYPE, "cannot find client settings for this client"); + return; + } ctx.setClientType(settings.getClientType().name()); Subscription subscription = settings.getSubscription(); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java index 5341c259761..b4a14e35fd2 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivityTest.java @@ -422,6 +422,31 @@ public void testReceiveMessage() { assertEquals(Code.MESSAGE_NOT_FOUND, getResponseCodeFromReceiveMessageResponseList(responseArgumentCaptor.getAllValues())); } + @Test + public void testReceiveMessageWithoutClientSettings() { + StreamObserver receiveStreamObserver = mock(ServerCallStreamObserver.class); + ArgumentCaptor responseArgumentCaptor = ArgumentCaptor.forClass(ReceiveMessageResponse.class); + doNothing().when(receiveStreamObserver).onNext(responseArgumentCaptor.capture()); + + when(this.grpcClientSettingsManager.getClientSettings(any())).thenReturn(null); + + this.receiveMessageActivity.receiveMessage( + createContext(), + ReceiveMessageRequest.newBuilder() + .setGroup(Resource.newBuilder().setName(CONSUMER_GROUP).build()) + .setMessageQueue(MessageQueue.newBuilder().setTopic(Resource.newBuilder().setName(TOPIC).build()).build()) + .setAutoRenew(true) + .setFilterExpression(FilterExpression.newBuilder() + .setType(FilterType.TAG) + .setExpression("*") + .build()) + .build(), + receiveStreamObserver + ); + + assertEquals(Code.UNRECOGNIZED_CLIENT_TYPE, getResponseCodeFromReceiveMessageResponseList(responseArgumentCaptor.getAllValues())); + } + private Code getResponseCodeFromReceiveMessageResponseList(List responseList) { for (ReceiveMessageResponse response : responseList) { if (response.hasStatus()) {