[ISSUE #3949] Add HeaderInterceptor

This commit is contained in:
zhouxiang
2022-07-13 11:29:07 +08:00
parent 80b2f95d17
commit 4059688c76
2 changed files with 59 additions and 0 deletions
@@ -31,6 +31,7 @@ import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import org.apache.rocketmq.common.thread.ThreadPoolMonitor;
import org.apache.rocketmq.proxy.configuration.ConfigurationManager;
import org.apache.rocketmq.proxy.grpc.interceptor.HeaderInterceptor;
import org.apache.rocketmq.proxy.grpc.service.GrpcService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -91,6 +92,7 @@ public class GrpcServer {
.channelType(NioServerSocketChannel.class)
.addService(messagingProcessor)
.executor(this.executor)
.intercept(new HeaderInterceptor())
.build();
log.info(
@@ -0,0 +1,57 @@
/*
* 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 com.google.common.net.HostAndPort;
import io.grpc.Grpc;
import io.grpc.Metadata;
import io.grpc.ServerCall;
import io.grpc.ServerCallHandler;
import io.grpc.ServerInterceptor;
import java.net.InetSocketAddress;
import java.net.SocketAddress;
import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants;
public class HeaderInterceptor implements ServerInterceptor {
@Override
public <ReqT, RespT> ServerCall.Listener<ReqT> interceptCall(ServerCall<ReqT, RespT> call, Metadata headers,
ServerCallHandler<ReqT, RespT> next) {
SocketAddress remoteSocketAddress = call.getAttributes()
.get(Grpc.TRANSPORT_ATTR_REMOTE_ADDR);
String remoteAddress = parseSocketAddress(remoteSocketAddress);
headers.put(InterceptorConstants.REMOTE_ADDRESS, remoteAddress);
SocketAddress localSocketAddress = call.getAttributes()
.get(Grpc.TRANSPORT_ATTR_LOCAL_ADDR);
String localAddress = parseSocketAddress(localSocketAddress);
headers.put(InterceptorConstants.LOCAL_ADDRESS, localAddress);
return next.startCall(call, headers);
}
private String parseSocketAddress(SocketAddress socketAddress) {
if (socketAddress instanceof InetSocketAddress) {
InetSocketAddress inetSocketAddress = (InetSocketAddress) socketAddress;
return HostAndPort.fromParts(inetSocketAddress.getAddress()
.getHostAddress(), inetSocketAddress.getPort())
.toString();
}
return "";
}
}