From c99a669400bd0c39752edd6f6225d5d241ee0377 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sat, 15 Aug 2026 08:00:45 -0700 Subject: [PATCH] fix(proxy): validate missing remoting ext fields Signed-off-by: liuhy --- .../activity/AbstractRemotingActivity.java | 4 ++-- .../AbstractRemotingActivityTest.java | 24 ++++++++++++++++++- 2 files changed, 25 insertions(+), 3 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivity.java index 2d09c394299..1418891a796 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivity.java @@ -65,13 +65,13 @@ protected RemotingCommand request(ChannelHandlerContext ctx, RemotingCommand req ProxyContext context, long timeoutMillis) throws Exception { String brokerName; if (request.getCode() == RequestCode.SEND_MESSAGE_V2 || request.getCode() == RequestCode.SEND_BATCH_MESSAGE) { - if (request.getExtFields().get(BROKER_NAME_FIELD_FOR_SEND_MESSAGE_V2) == null) { + if (request.getExtFields() == null || request.getExtFields().get(BROKER_NAME_FIELD_FOR_SEND_MESSAGE_V2) == null) { return RemotingCommand.buildErrorResponse(ResponseCode.VERSION_NOT_SUPPORTED, "Request doesn't have field bname"); } brokerName = request.getExtFields().get(BROKER_NAME_FIELD_FOR_SEND_MESSAGE_V2); } else { - if (request.getExtFields().get(BROKER_NAME_FIELD) == null) { + if (request.getExtFields() == null || request.getExtFields().get(BROKER_NAME_FIELD) == null) { return RemotingCommand.buildErrorResponse(ResponseCode.VERSION_NOT_SUPPORTED, "Request doesn't have field bname"); } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivityTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivityTest.java index 11dd6bc40c4..5bf958ce9e1 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivityTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/remoting/activity/AbstractRemotingActivityTest.java @@ -120,6 +120,28 @@ public void testRequestInvalid() throws Exception { verify(ctx, never()).writeAndFlush(any()); } + @Test + public void testRequestWithoutExtFieldsIsInvalid() throws Exception { + RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.PULL_MESSAGE, null); + request.setExtFields(null); + + RemotingCommand remotingCommand = remotingActivity.request(ctx, request, null, 10000); + + assertThat(remotingCommand.getCode()).isEqualTo(ResponseCode.VERSION_NOT_SUPPORTED); + verify(ctx, never()).writeAndFlush(any()); + } + + @Test + public void testSendMessageV2WithoutExtFieldsIsInvalid() throws Exception { + RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.SEND_MESSAGE_V2, null); + request.setExtFields(null); + + RemotingCommand remotingCommand = remotingActivity.request(ctx, request, null, 10000); + + assertThat(remotingCommand.getCode()).isEqualTo(ResponseCode.VERSION_NOT_SUPPORTED); + verify(ctx, never()).writeAndFlush(any()); + } + @Test public void testRequestProxyException() throws Exception { ArgumentCaptor captor = ArgumentCaptor.forClass(RemotingCommand.class); @@ -199,4 +221,4 @@ public void testRequestDefaultException() throws Exception { verify(ctx, times(1)).writeAndFlush(captor.capture()); assertThat(captor.getValue().getCode()).isEqualTo(ResponseCode.SYSTEM_ERROR); } -} \ No newline at end of file +}