From a47e489734be444cccb7f97fc6b2edd29d86d895 Mon Sep 17 00:00:00 2001 From: liuhy Date: Mon, 3 Aug 2026 08:13:22 -0700 Subject: [PATCH 1/2] fix(proxy): tolerate unresolved grpc peer addresses --- .../grpc/interceptor/HeaderInterceptor.java | 9 +- .../interceptor/HeaderInterceptorTest.java | 92 +++++++++++++++++++ 2 files changed, 99 insertions(+), 2 deletions(-) create mode 100644 proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java index e3e78841559..d4f1b6906af 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptor.java @@ -72,9 +72,14 @@ public ServerCall.Listener interceptCall( private String parseSocketAddress(SocketAddress socketAddress) { if (socketAddress instanceof InetSocketAddress) { InetSocketAddress inetSocketAddress = (InetSocketAddress) socketAddress; + String host = inetSocketAddress.getAddress() == null + ? inetSocketAddress.getHostString() + : inetSocketAddress.getAddress().getHostAddress(); + if (StringUtils.isBlank(host)) { + return ""; + } return HostAndPort.fromParts( - inetSocketAddress.getAddress() - .getHostAddress(), + host, inetSocketAddress.getPort() ).toString(); } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java new file mode 100644 index 00000000000..f819ead2a7d --- /dev/null +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java @@ -0,0 +1,92 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 org.apache.rocketmq.proxy.grpc.interceptor; + +import io.grpc.Attributes; +import io.grpc.Grpc; +import io.grpc.Metadata; +import io.grpc.MethodDescriptor; +import io.grpc.ServerCall; +import io.grpc.ServerCallHandler; +import io.grpc.Status; +import java.net.InetSocketAddress; +import org.apache.rocketmq.common.constant.GrpcConstants; +import org.junit.Test; + +import static org.junit.Assert.assertEquals; + +public class HeaderInterceptorTest { + + @Test + public void testInterceptCallWithUnresolvedSocketAddress() { + HeaderInterceptor interceptor = new HeaderInterceptor(); + Metadata headers = new Metadata(); + Attributes attributes = Attributes.newBuilder() + .set(Grpc.TRANSPORT_ATTR_REMOTE_ADDR, InetSocketAddress.createUnresolved("remote.example.com", 8081)) + .set(Grpc.TRANSPORT_ATTR_LOCAL_ADDR, InetSocketAddress.createUnresolved("local.example.com", 8080)) + .build(); + + ServerCall call = new TestServerCall(attributes); + ServerCallHandler handler = (serverCall, metadata) -> new ServerCall.Listener() { + }; + + interceptor.interceptCall(call, headers, handler); + + assertEquals("remote.example.com:8081", headers.get(GrpcConstants.REMOTE_ADDRESS)); + assertEquals("local.example.com:8080", headers.get(GrpcConstants.LOCAL_ADDRESS)); + } + + private static class TestServerCall extends ServerCall { + private final Attributes attributes; + + private TestServerCall(Attributes attributes) { + this.attributes = attributes; + } + + @Override + public void request(int numMessages) { + } + + @Override + public void sendHeaders(Metadata headers) { + } + + @Override + public void sendMessage(String message) { + } + + @Override + public void close(Status status, Metadata trailers) { + } + + @Override + public boolean isCancelled() { + return false; + } + + @Override + public Attributes getAttributes() { + return attributes; + } + + @Override + public MethodDescriptor getMethodDescriptor() { + return null; + } + } +} From 3de75063ef35e509a81eb786270ea809e54c38cc Mon Sep 17 00:00:00 2001 From: liuhy Date: Tue, 4 Aug 2026 03:51:27 -0700 Subject: [PATCH 2/2] test(proxy): harden unresolved address test double --- .../rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java index f819ead2a7d..6ba8d377d96 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/interceptor/HeaderInterceptorTest.java @@ -86,7 +86,7 @@ public Attributes getAttributes() { @Override public MethodDescriptor getMethodDescriptor() { - return null; + throw new UnsupportedOperationException("Method descriptor is not used by this test double"); } } }