From 59e6ffccf7497e6d95fdbea8cc8dc4e38b7bc8e1 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Wed, 16 Mar 2022 10:45:59 +0200 Subject: [PATCH 1/9] [Proxy] Refactor proxy code to simplify and clarify it --- .../proxy/server/DirectProxyHandler.java | 101 ++++++++---------- .../proxy/server/ParserProxyHandler.java | 8 +- 2 files changed, 53 insertions(+), 56 deletions(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index edf4f18411366..02901521f2150 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -28,7 +28,6 @@ import io.netty.channel.ChannelFuture; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; -import io.netty.channel.ChannelId; import io.netty.channel.ChannelInitializer; import io.netty.channel.ChannelOption; import io.netty.channel.socket.SocketChannel; @@ -43,11 +42,7 @@ import io.netty.util.concurrent.Future; import io.netty.util.concurrent.FutureListener; import java.net.InetSocketAddress; -import java.net.URI; -import java.net.URISyntaxException; import java.util.Arrays; -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.TimeUnit; import java.util.function.Supplier; import javax.net.ssl.SSLSession; @@ -71,11 +66,11 @@ public class DirectProxyHandler { @Getter private final Channel inboundChannel; + private final ProxyConnection proxyConnection; @Getter Channel outboundChannel; @Getter private final Rate inboundChannelRequestsRate; - protected static Map inboundOutboundChannelMap = new ConcurrentHashMap<>(); private final String originalPrincipal; private final AuthData clientAuthData; private final String clientAuthMethod; @@ -86,12 +81,13 @@ public class DirectProxyHandler { private final ProxyService service; private final Runnable onHandshakeCompleteAction; - public DirectProxyHandler(ProxyService service, ProxyConnection proxyConnection, String targetBrokerUrl, + public DirectProxyHandler(ProxyService service, ProxyConnection proxyConnection, String brokerHostAndPort, InetSocketAddress targetBrokerAddress, int protocolVersion, Supplier sslHandlerSupplier) { this.service = service; this.authentication = proxyConnection.getClientAuthentication(); this.inboundChannel = proxyConnection.ctx().channel(); + this.proxyConnection = proxyConnection; this.inboundChannelRequestsRate = new Rate(); this.originalPrincipal = proxyConnection.clientAuthRole; this.clientAuthData = proxyConnection.clientAuthData; @@ -110,6 +106,16 @@ public DirectProxyHandler(ProxyService service, ProxyConnection proxyConnection, b.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, brokerProxyConnectTimeoutMs); } b.group(inboundChannel.eventLoop()).channel(inboundChannel.getClass()).option(ChannelOption.AUTO_READ, false); + + String remoteHost; + try { + remoteHost = parseHost(brokerHostAndPort); + } catch (IllegalArgumentException e) { + log.warn("[{}] Failed to parse broker host '{}'", inboundChannel, brokerHostAndPort, e); + inboundChannel.close(); + return; + } + b.handler(new ChannelInitializer() { @Override protected void initChannel(SocketChannel ch) { @@ -123,55 +129,39 @@ protected void initChannel(SocketChannel ch) { } ch.pipeline().addLast("frameDecoder", new LengthFieldBasedFrameDecoder( Commands.DEFAULT_MAX_MESSAGE_SIZE + Commands.MESSAGE_SIZE_FRAME_PADDING, 0, 4, 0, 4)); - ch.pipeline().addLast("proxyOutboundHandler", new ProxyBackendHandler(config, protocolVersion)); + ch.pipeline().addLast("proxyOutboundHandler", + new ProxyBackendHandler(config, protocolVersion, remoteHost)); } }); - URI targetBroker; - try { - // targetBrokerUrl is coming in the "hostname:6650" form, so we need - // to extract host and port - targetBroker = new URI("pulsar://" + targetBrokerUrl); - } catch (URISyntaxException e) { - log.warn("[{}] Failed to parse broker url '{}'", inboundChannel, targetBrokerUrl, e); - inboundChannel.close(); - return; - } - ChannelFuture f = b.connect(targetBrokerAddress); outboundChannel = f.channel(); f.addListener(future -> { if (!future.isSuccess()) { // Close the connection if the connection attempt has failed. log.warn("[{}] Establishing connection to {} ({}) failed. Closing inbound channel.", inboundChannel, - targetBrokerAddress, targetBrokerUrl, future.cause()); + targetBrokerAddress, brokerHostAndPort, future.cause()); inboundChannel.close(); return; } - final ProxyBackendHandler cnx = (ProxyBackendHandler) outboundChannel.pipeline() - .get("proxyOutboundHandler"); - cnx.setRemoteHostName(targetBroker.getHost()); - - // if enable full parsing feature - if (service.getProxyLogLevel() == 2) { - //Set a map between inbound and outbound, - //so can find inbound by outbound or find outbound by inbound - inboundOutboundChannelMap.put(outboundChannel.id(), inboundChannel.id()); - } + }); + } - if (!config.isHaProxyProtocolEnabled()) { - return; - } + private String parseHost(String brokerPortAndHost) { + int pos = brokerPortAndHost.indexOf(':'); + if (pos > 0) { + return brokerPortAndHost.substring(0, pos); + } else { + throw new IllegalArgumentException("Illegal broker host:port '" + brokerPortAndHost + "'"); + } + } - if (proxyConnection.hasHAProxyMessage()) { - outboundChannel.writeAndFlush(encodeProxyProtocolMessage(proxyConnection.getHAProxyMessage())); - } else { - if (!(inboundChannel.remoteAddress() instanceof InetSocketAddress)) { - return; - } - if (!(outboundChannel.localAddress() instanceof InetSocketAddress)) { - return; - } + private void writeHAProxyMessage() { + if (proxyConnection.hasHAProxyMessage()) { + outboundChannel.writeAndFlush(encodeProxyProtocolMessage(proxyConnection.getHAProxyMessage())); + } else { + if (inboundChannel.remoteAddress() instanceof InetSocketAddress + && outboundChannel.localAddress() instanceof InetSocketAddress) { InetSocketAddress clientAddress = (InetSocketAddress) inboundChannel.remoteAddress(); String sourceAddress = clientAddress.getAddress().getHostAddress(); int sourcePort = clientAddress.getPort(); @@ -179,11 +169,12 @@ protected void initChannel(SocketChannel ch) { String destinationAddress = proxyAddress.getAddress().getHostAddress(); int destinationPort = proxyAddress.getPort(); HAProxyMessage msg = new HAProxyMessage(HAProxyProtocolVersion.V1, HAProxyCommand.PROXY, - HAProxyProxiedProtocol.TCP4, sourceAddress, destinationAddress, sourcePort, destinationPort); + HAProxyProxiedProtocol.TCP4, sourceAddress, destinationAddress, sourcePort, + destinationPort); outboundChannel.writeAndFlush(encodeProxyProtocolMessage(msg)); msg.release(); } - }); + } } private ByteBuf encodeProxyProtocolMessage(HAProxyMessage msg) { @@ -220,19 +211,25 @@ enum BackendState { public class ProxyBackendHandler extends PulsarDecoder implements FutureListener { private BackendState state = BackendState.Init; - private String remoteHostName; + private final String remoteHostName; protected ChannelHandlerContext ctx; private final ProxyConfiguration config; private final int protocolVersion; - public ProxyBackendHandler(ProxyConfiguration config, int protocolVersion) { + public ProxyBackendHandler(ProxyConfiguration config, int protocolVersion, String remoteHostName) { this.config = config; this.protocolVersion = protocolVersion; + this.remoteHostName = remoteHostName; } @Override public void channelActive(ChannelHandlerContext ctx) throws Exception { this.ctx = ctx; + + if (config.isHaProxyProtocolEnabled()) { + writeHAProxyMessage(); + } + // Send the Connect command to broker authenticationDataProvider = authentication.getAuthData(remoteHostName); AuthData authData = authenticationDataProvider.authenticate(AuthData.INIT_AUTH_DATA); @@ -389,20 +386,20 @@ private void startDirectProxying(CommandConnected connected) { inboundChannel.pipeline().addBefore("handler", "inboundParser", new ParserProxyHandler(service, inboundChannel, ParserProxyHandler.FRONTEND_CONN, - connected.getMaxMessageSize())); + connected.getMaxMessageSize(), outboundChannel.id())); outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", new ParserProxyHandler(service, outboundChannel, ParserProxyHandler.BACKEND_CONN, - connected.getMaxMessageSize())); + connected.getMaxMessageSize(), inboundChannel.id())); } else { inboundChannel.pipeline().addBefore("handler", "inboundParser", new ParserProxyHandler(service, inboundChannel, ParserProxyHandler.FRONTEND_CONN, - Commands.DEFAULT_MAX_MESSAGE_SIZE)); + Commands.DEFAULT_MAX_MESSAGE_SIZE, outboundChannel.id())); outboundChannel.pipeline().addBefore("proxyOutboundHandler", "outboundParser", new ParserProxyHandler(service, outboundChannel, ParserProxyHandler.BACKEND_CONN, - Commands.DEFAULT_MAX_MESSAGE_SIZE)); + Commands.DEFAULT_MAX_MESSAGE_SIZE, inboundChannel.id())); } } } @@ -418,10 +415,6 @@ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) { ctx.close(); } - public void setRemoteHostName(String remoteHostName) { - this.remoteHostName = remoteHostName; - } - private boolean verifyTlsHostName(String hostname, ChannelHandlerContext ctx) { ChannelHandler sslHandler = ctx.channel().pipeline().get("tls"); diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java index fea1a401af61d..41a9b594f9363 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ParserProxyHandler.java @@ -25,6 +25,7 @@ import io.netty.buffer.Unpooled; import io.netty.channel.Channel; import io.netty.channel.ChannelHandlerContext; +import io.netty.channel.ChannelId; import io.netty.channel.ChannelInboundHandlerAdapter; import java.nio.charset.StandardCharsets; import java.util.ArrayList; @@ -53,6 +54,7 @@ public class ParserProxyHandler extends ChannelInboundHandlerAdapter { private final String connType; private final int maxMessageSize; + private final ChannelId peerChannelId; private final ProxyService service; @@ -66,11 +68,13 @@ public class ParserProxyHandler extends ChannelInboundHandlerAdapter { */ private static final Map consumerHashMap = new ConcurrentHashMap<>(); - public ParserProxyHandler(ProxyService service, Channel channel, String type, int maxMessageSize) { + public ParserProxyHandler(ProxyService service, Channel channel, String type, int maxMessageSize, + ChannelId peerChannelId) { this.service = service; this.channel = channel; this.connType = type; this.maxMessageSize = maxMessageSize; + this.peerChannelId = peerChannelId; } private void logging(Channel conn, BaseCommand.Type cmdtype, String info, List messages) { @@ -154,7 +158,7 @@ public void channelRead(ChannelHandlerContext ctx, Object msg) { break; } topicName = TopicName.get(ParserProxyHandler.consumerHashMap.get(cmd.getMessage().getConsumerId() - + "," + DirectProxyHandler.inboundOutboundChannelMap.get(ctx.channel().id()))); + + "," + peerChannelId)); msgBytes = new MutableLong(0); MessageParser.parseMessage(topicName, -1L, -1L, buffer, (message) -> { From 09a80eefb4b49cfbde80130f26e0d2c07bcbbbe9 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 17 Mar 2022 07:17:59 +0200 Subject: [PATCH 2/9] Make parsing of broker host and port IPv6 compatible by using lastIndexOf(':') --- .../proxy/server/BrokerProxyValidator.java | 2 +- .../proxy/server/DirectProxyHandler.java | 2 +- .../server/BrokerProxyValidatorTest.java | 20 +++++++++++++++++++ 3 files changed, 22 insertions(+), 2 deletions(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/BrokerProxyValidator.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/BrokerProxyValidator.java index debe1f7fcac87..b0529c2a777e1 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/BrokerProxyValidator.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/BrokerProxyValidator.java @@ -113,7 +113,7 @@ private static List parseCommaSeparatedConfigValue(String configValue) { } public CompletableFuture resolveAndCheckTargetAddress(String hostAndPort) { - int pos = hostAndPort.indexOf(':'); + int pos = hostAndPort.lastIndexOf(':'); String host = hostAndPort.substring(0, pos); int port = Integer.parseInt(hostAndPort.substring(pos + 1)); if (!isPortAllowed(port)) { diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index 02901521f2150..08c457c36739c 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -148,7 +148,7 @@ protected void initChannel(SocketChannel ch) { } private String parseHost(String brokerPortAndHost) { - int pos = brokerPortAndHost.indexOf(':'); + int pos = brokerPortAndHost.lastIndexOf(':'); if (pos > 0) { return brokerPortAndHost.substring(0, pos); } else { diff --git a/pulsar-proxy/src/test/java/org/apache/pulsar/proxy/server/BrokerProxyValidatorTest.java b/pulsar-proxy/src/test/java/org/apache/pulsar/proxy/server/BrokerProxyValidatorTest.java index 8e457554cf5ad..fba3c36e26616 100644 --- a/pulsar-proxy/src/test/java/org/apache/pulsar/proxy/server/BrokerProxyValidatorTest.java +++ b/pulsar-proxy/src/test/java/org/apache/pulsar/proxy/server/BrokerProxyValidatorTest.java @@ -90,6 +90,26 @@ public void shouldAllowAllWithWildcard() throws Exception { brokerProxyValidator.resolveAndCheckTargetAddress("myhost.mydomain:6650").get(); } + @Test + public void shouldAllowIPv6Address() throws Exception { + BrokerProxyValidator brokerProxyValidator = new BrokerProxyValidator( + createMockedAddressResolver("fd4d:801b:73fa:abcd:0000:0000:0000:0001"), + "*" + , "fd4d:801b:73fa:abcd::/64" + , "6650"); + brokerProxyValidator.resolveAndCheckTargetAddress("myhost.mydomain:6650").get(); + } + + @Test + public void shouldAllowIPv6AddressNumeric() throws Exception { + BrokerProxyValidator brokerProxyValidator = new BrokerProxyValidator( + createMockedAddressResolver("fd4d:801b:73fa:abcd:0000:0000:0000:0001"), + "*" + , "fd4d:801b:73fa:abcd::/64" + , "6650"); + brokerProxyValidator.resolveAndCheckTargetAddress("fd4d:801b:73fa:abcd:0000:0000:0000:0001:6650").get(); + } + private AddressResolver createMockedAddressResolver(String ipAddressResult) { AddressResolver inetSocketAddressResolver = mock(AddressResolver.class); when(inetSocketAddressResolver.resolve(any())).then(invocationOnMock -> { From 11346c5a6f7ae43542b36fedbc29f471aaa66c67 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 17 Mar 2022 07:33:36 +0200 Subject: [PATCH 3/9] Use thenAcceptAsync so that exceptions are handled --- .../java/org/apache/pulsar/proxy/server/ProxyConnection.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index 382e4751cd7c5..e0b21242c63b7 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -271,13 +271,13 @@ private synchronized void completeConnect(AuthData clientData) throws PulsarClie } brokerProxyValidator.resolveAndCheckTargetAddress(proxyToBrokerUrl) - .thenAccept(address -> ctx().executor().submit(() -> { + .thenAcceptAsync(address -> { // Client already knows which broker to connect. Let's open a // connection there and just pass bytes in both directions state = State.ProxyConnectionToBroker; directProxyHandler = new DirectProxyHandler(service, this, proxyToBrokerUrl, address, protocolVersionToAdvertise, sslHandlerSupplier); - })) + }, ctx.executor()) .exceptionally(throwable -> { if (throwable instanceof TargetAddressDeniedException || throwable.getCause() instanceof TargetAddressDeniedException) { From 72974ad7087fdcd86398d5dd081235cc8afcf5a2 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 17 Mar 2022 12:32:34 +0200 Subject: [PATCH 4/9] Enable auto read in proxy and remove ctx.read() calls --- .../pulsar/proxy/server/DirectProxyHandler.java | 15 ++++----------- .../pulsar/proxy/server/ProxyConnection.java | 4 +--- 2 files changed, 5 insertions(+), 14 deletions(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index 08c457c36739c..b61c4488295dd 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -105,7 +105,8 @@ public DirectProxyHandler(ProxyService service, ProxyConnection proxyConnection, if (brokerProxyConnectTimeoutMs > 0) { b.option(ChannelOption.CONNECT_TIMEOUT_MILLIS, brokerProxyConnectTimeoutMs); } - b.group(inboundChannel.eventLoop()).channel(inboundChannel.getClass()).option(ChannelOption.AUTO_READ, false); + b.group(inboundChannel.eventLoop()) + .channel(inboundChannel.getClass()); String remoteHost; try { @@ -237,7 +238,6 @@ public void channelActive(ChannelHandlerContext ctx) throws Exception { command = Commands.newConnect(authentication.getAuthMethodName(), authData, protocolVersion, "Pulsar proxy", null /* target broker */, originalPrincipal, clientAuthData, clientAuthMethod); outboundChannel.writeAndFlush(command); - outboundChannel.read(); } @Override @@ -298,7 +298,6 @@ protected void handleAuthChallenge(CommandAuthChallenge authChallenge) { } outboundChannel.writeAndFlush(request); - outboundChannel.read(); } catch (Exception e) { log.error("Error mutual verify", e); } @@ -308,9 +307,7 @@ protected void handleAuthChallenge(CommandAuthChallenge authChallenge) { public void operationComplete(Future future) { // This is invoked when the write operation on the paired connection // is completed - if (future.isSuccess()) { - outboundChannel.read(); - } else { + if (!future.isSuccess()) { log.warn("[{}] [{}] Failed to write on proxy connection. Closing both connections.", inboundChannel, outboundChannel, future.cause()); inboundChannel.close(); @@ -347,11 +344,7 @@ protected void handleConnected(CommandConnected connected) { connected.hasMaxMessageSize() ? connected.getMaxMessageSize() : Commands.INVALID_MAX_MESSAGE_SIZE; inboundChannel.writeAndFlush(Commands.newConnected(connected.getProtocolVersion(), maxMessageSize)) .addListener(future -> { - if (future.isSuccess()) { - // Start reading from both connections - inboundChannel.read(); - outboundChannel.read(); - } else { + if (!future.isSuccess()) { log.warn("[{}] [{}] Failed to write to inbound connection. Closing both connections.", inboundChannel, outboundChannel, future.cause()); diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index e0b21242c63b7..04a3ab5ec83ae 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -221,9 +221,7 @@ public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exce public void operationComplete(Future future) { // This is invoked when the write operation on the paired connection is // completed - if (future.isSuccess()) { - ctx.read(); - } else { + if (!future.isSuccess()) { LOG.warn("[{}] Error in writing to inbound channel. Closing", remoteAddress, future.cause()); directProxyHandler.outboundChannel.close(); } From ca6e8df92c8490729fd2de50d917e6c90a0f37fd Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 17 Mar 2022 12:36:53 +0200 Subject: [PATCH 5/9] Fix comment --- .../java/org/apache/pulsar/proxy/server/ProxyConnection.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index 04a3ab5ec83ae..cd973ee81dbf0 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -222,7 +222,7 @@ public void operationComplete(Future future) { // This is invoked when the write operation on the paired connection is // completed if (!future.isSuccess()) { - LOG.warn("[{}] Error in writing to inbound channel. Closing", remoteAddress, future.cause()); + LOG.warn("[{}] Error in writing to outbound channel. Closing", remoteAddress, future.cause()); directProxyHandler.outboundChannel.close(); } } From 4a240613a2742553069fffb976a88caab6c04574 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 17 Mar 2022 12:44:15 +0200 Subject: [PATCH 6/9] Refactor future listeners --- .../proxy/server/DirectProxyHandler.java | 24 +++++++------------ .../pulsar/proxy/server/ProxyConnection.java | 23 +++++++----------- 2 files changed, 18 insertions(+), 29 deletions(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index b61c4488295dd..52b515e617ff4 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -39,8 +39,6 @@ import io.netty.handler.ssl.SslHandler; import io.netty.handler.timeout.ReadTimeoutHandler; import io.netty.util.CharsetUtil; -import io.netty.util.concurrent.Future; -import io.netty.util.concurrent.FutureListener; import java.net.InetSocketAddress; import java.util.Arrays; import java.util.concurrent.TimeUnit; @@ -209,7 +207,7 @@ enum BackendState { Init, HandshakeCompleted } - public class ProxyBackendHandler extends PulsarDecoder implements FutureListener { + public class ProxyBackendHandler extends PulsarDecoder { private BackendState state = BackendState.Init; private final String remoteHostName; @@ -258,7 +256,14 @@ public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exce if (msg instanceof ByteBuf) { ProxyService.BYTES_COUNTER.inc(((ByteBuf) msg).readableBytes()); } - inboundChannel.writeAndFlush(msg).addListener(this); + inboundChannel.writeAndFlush(msg) + .addListener(future -> { + if (!future.isSuccess()) { + log.warn("[{}] [{}] Failed to write on proxy connection. Closing both connections.", + inboundChannel, outboundChannel, future.cause()); + inboundChannel.close(); + } + }); break; default: @@ -303,17 +308,6 @@ protected void handleAuthChallenge(CommandAuthChallenge authChallenge) { } } - @Override - public void operationComplete(Future future) { - // This is invoked when the write operation on the paired connection - // is completed - if (!future.isSuccess()) { - log.warn("[{}] [{}] Failed to write on proxy connection. Closing both connections.", inboundChannel, - outboundChannel, future.cause()); - inboundChannel.close(); - } - } - @Override protected void messageReceived() { // no-op diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index cd973ee81dbf0..09e47fa39f984 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -24,8 +24,6 @@ import io.netty.channel.ChannelHandlerContext; import io.netty.handler.codec.haproxy.HAProxyMessage; import io.netty.handler.ssl.SslHandler; -import io.netty.util.concurrent.Future; -import io.netty.util.concurrent.FutureListener; import java.net.SocketAddress; import java.util.Collections; import java.util.List; @@ -65,7 +63,7 @@ * Handles incoming discovery request from client and sends appropriate response back to client. * */ -public class ProxyConnection extends PulsarHandler implements FutureListener { +public class ProxyConnection extends PulsarHandler { private static final Logger LOG = LoggerFactory.getLogger(ProxyConnection.class); // ConnectionPool is used by the proxy to issue lookup requests private ConnectionPool connectionPool; @@ -209,7 +207,14 @@ public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exce directProxyHandler.getInboundChannelRequestsRate().recordEvent(bytes); ProxyService.BYTES_COUNTER.inc(bytes); } - directProxyHandler.outboundChannel.writeAndFlush(msg).addListener(this); + directProxyHandler.outboundChannel.writeAndFlush(msg) + .addListener(future -> { + if (!future.isSuccess()) { + LOG.warn("[{}] Error in writing to outbound channel. Closing", remoteAddress, + future.cause()); + directProxyHandler.outboundChannel.close(); + } + }); break; default: @@ -217,16 +222,6 @@ public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exce } } - @Override - public void operationComplete(Future future) { - // This is invoked when the write operation on the paired connection is - // completed - if (!future.isSuccess()) { - LOG.warn("[{}] Error in writing to outbound channel. Closing", remoteAddress, future.cause()); - directProxyHandler.outboundChannel.close(); - } - } - private synchronized void completeConnect(AuthData clientData) throws PulsarClientException { if (service.getConfiguration().isAuthenticationEnabled()) { if (service.getConfiguration().isForwardAuthorizationCredentials()) { From a04aa5f4372f397b1c5feb895d0f64db7a169dcf Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 17 Mar 2022 19:01:59 +0200 Subject: [PATCH 7/9] Handle backpressure properly by switching auto read off when writability changes - change auto read of the proxy-broker connection based on the writability of the client-proxy connection - change auto read of the client-proxy connection based on the writability of the proxy-broker connection --- .../pulsar/proxy/server/DirectProxyHandler.java | 11 +++++++++++ .../apache/pulsar/proxy/server/ProxyConnection.java | 11 +++++++++++ 2 files changed, 22 insertions(+) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index 52b515e617ff4..89a61e6197c50 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -176,6 +176,8 @@ private void writeHAProxyMessage() { } } + + private ByteBuf encodeProxyProtocolMessage(HAProxyMessage msg) { // Max length of v1 version proxy protocol message is 108 ByteBuf out = Unpooled.buffer(108); @@ -238,6 +240,15 @@ public void channelActive(ChannelHandlerContext ctx) throws Exception { outboundChannel.writeAndFlush(command); } + @Override + public void channelWritabilityChanged(ChannelHandlerContext ctx) throws Exception { + // handle backpressure + // stop/resume reading input from connection between the client and the proxy + // when the writability of the connection between the proxy and the broker changes + inboundChannel.config().setAutoRead(ctx.channel().isWritable()); + super.channelWritabilityChanged(ctx); + } + @Override public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exception { switch (state) { diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index 09e47fa39f984..dd4acfe131d82 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -184,6 +184,17 @@ public void exceptionCaught(ChannelHandlerContext ctx, Throwable cause) throws E ctx.close(); } + @Override + public void channelWritabilityChanged(ChannelHandlerContext ctx) throws Exception { + if (directProxyHandler != null && directProxyHandler.outboundChannel != null) { + // handle backpressure + // stop/resume reading input from connection between the proxy and the broker + // when the writability of the connection between the client and the proxy changes + directProxyHandler.outboundChannel.config().setAutoRead(ctx.channel().isWritable()); + } + super.channelWritabilityChanged(ctx); + } + @Override public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exception { if (msg instanceof HAProxyMessage) { From 8a087550c9ecca348f93f7c63dd7f6882a3ada33 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Thu, 17 Mar 2022 19:23:13 +0200 Subject: [PATCH 8/9] Address review feedback: make utility method static --- .../java/org/apache/pulsar/proxy/server/DirectProxyHandler.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index 89a61e6197c50..6a021c4cdd4d6 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -146,7 +146,7 @@ protected void initChannel(SocketChannel ch) { }); } - private String parseHost(String brokerPortAndHost) { + private static String parseHost(String brokerPortAndHost) { int pos = brokerPortAndHost.lastIndexOf(':'); if (pos > 0) { return brokerPortAndHost.substring(0, pos); From 5dfdf12689e039993950848ee7ec7aedd4c4b0e8 Mon Sep 17 00:00:00 2001 From: Lari Hotari Date: Fri, 18 Mar 2022 08:12:57 +0200 Subject: [PATCH 9/9] Consistently handle write errors - delegate exception handling to exceptionCaught method --- .../proxy/server/DirectProxyHandler.java | 30 +++++------- .../pulsar/proxy/server/ProxyConnection.java | 46 +++++-------------- 2 files changed, 23 insertions(+), 53 deletions(-) diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java index 6a021c4cdd4d6..0cbceb191b18d 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/DirectProxyHandler.java @@ -26,6 +26,7 @@ import io.netty.buffer.Unpooled; import io.netty.channel.Channel; import io.netty.channel.ChannelFuture; +import io.netty.channel.ChannelFutureListener; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.channel.ChannelInitializer; @@ -157,7 +158,8 @@ private static String parseHost(String brokerPortAndHost) { private void writeHAProxyMessage() { if (proxyConnection.hasHAProxyMessage()) { - outboundChannel.writeAndFlush(encodeProxyProtocolMessage(proxyConnection.getHAProxyMessage())); + outboundChannel.writeAndFlush(encodeProxyProtocolMessage(proxyConnection.getHAProxyMessage())) + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } else { if (inboundChannel.remoteAddress() instanceof InetSocketAddress && outboundChannel.localAddress() instanceof InetSocketAddress) { @@ -170,7 +172,8 @@ private void writeHAProxyMessage() { HAProxyMessage msg = new HAProxyMessage(HAProxyProtocolVersion.V1, HAProxyCommand.PROXY, HAProxyProxiedProtocol.TCP4, sourceAddress, destinationAddress, sourcePort, destinationPort); - outboundChannel.writeAndFlush(encodeProxyProtocolMessage(msg)); + outboundChannel.writeAndFlush(encodeProxyProtocolMessage(msg)) + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); msg.release(); } } @@ -237,7 +240,8 @@ public void channelActive(ChannelHandlerContext ctx) throws Exception { ByteBuf command; command = Commands.newConnect(authentication.getAuthMethodName(), authData, protocolVersion, "Pulsar proxy", null /* target broker */, originalPrincipal, clientAuthData, clientAuthMethod); - outboundChannel.writeAndFlush(command); + outboundChannel.writeAndFlush(command) + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } @Override @@ -268,13 +272,7 @@ public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exce ProxyService.BYTES_COUNTER.inc(((ByteBuf) msg).readableBytes()); } inboundChannel.writeAndFlush(msg) - .addListener(future -> { - if (!future.isSuccess()) { - log.warn("[{}] [{}] Failed to write on proxy connection. Closing both connections.", - inboundChannel, outboundChannel, future.cause()); - inboundChannel.close(); - } - }); + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); break; default: @@ -313,7 +311,8 @@ protected void handleAuthChallenge(CommandAuthChallenge authChallenge) { log.debug("{} Mutual auth {}", ctx.channel(), authentication.getAuthMethodName()); } - outboundChannel.writeAndFlush(request); + outboundChannel.writeAndFlush(request) + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } catch (Exception e) { log.error("Error mutual verify", e); } @@ -348,14 +347,7 @@ protected void handleConnected(CommandConnected connected) { int maxMessageSize = connected.hasMaxMessageSize() ? connected.getMaxMessageSize() : Commands.INVALID_MAX_MESSAGE_SIZE; inboundChannel.writeAndFlush(Commands.newConnected(connected.getProtocolVersion(), maxMessageSize)) - .addListener(future -> { - if (!future.isSuccess()) { - log.warn("[{}] [{}] Failed to write to inbound connection. Closing both connections.", - inboundChannel, - outboundChannel, future.cause()); - inboundChannel.close(); - } - }); + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } private void startDirectProxying(CommandConnected connected) { diff --git a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java index dd4acfe131d82..58203eee51c28 100644 --- a/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java +++ b/pulsar-proxy/src/main/java/org/apache/pulsar/proxy/server/ProxyConnection.java @@ -20,6 +20,7 @@ import static com.google.common.base.Preconditions.checkArgument; import io.netty.buffer.ByteBuf; +import io.netty.channel.ChannelFutureListener; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; import io.netty.handler.codec.haproxy.HAProxyMessage; @@ -219,13 +220,7 @@ public void channelRead(final ChannelHandlerContext ctx, Object msg) throws Exce ProxyService.BYTES_COUNTER.inc(bytes); } directProxyHandler.outboundChannel.writeAndFlush(msg) - .addListener(future -> { - if (!future.isSuccess()) { - LOG.warn("[{}] Error in writing to outbound channel. Closing", remoteAddress, - future.cause()); - directProxyHandler.outboundChannel.close(); - } - }); + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); break; default: @@ -270,7 +265,7 @@ private synchronized void completeConnect(AuthData clientData) throws PulsarClie ctx() .writeAndFlush( Commands.newError(-1, ServerError.ServiceNotReady, "Target broker isn't available.")) - .addListener(future -> ctx().close()); + .addListener(ChannelFutureListener.CLOSE); return; } @@ -300,7 +295,7 @@ private synchronized void completeConnect(AuthData clientData) throws PulsarClie .writeAndFlush( Commands.newError(-1, ServerError.ServiceNotReady, "Target broker cannot be validated.")) - .addListener(future -> ctx().close()); + .addListener(ChannelFutureListener.CLOSE); return null; }); } else { @@ -309,7 +304,8 @@ private synchronized void completeConnect(AuthData clientData) throws PulsarClie // partitions metadata lookups state = State.ProxyLookupRequests; lookupProxyHandler = new LookupProxyHandler(service, this); - ctx.writeAndFlush(Commands.newConnected(protocolVersionToAdvertise)); + ctx.writeAndFlush(Commands.newConnected(protocolVersionToAdvertise)) + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); } } @@ -328,7 +324,8 @@ private void doAuthentication(AuthData clientData) throws Exception { } // auth not complete, continue auth with client side. - ctx.writeAndFlush(Commands.newAuthChallenge(authMethod, brokerData, protocolVersionToAdvertise)); + ctx.writeAndFlush(Commands.newAuthChallenge(authMethod, brokerData, protocolVersionToAdvertise)) + .addListener(ChannelFutureListener.FIRE_EXCEPTION_ON_FAILURE); if (LOG.isDebugEnabled()) { LOG.debug("[{}] Authentication in progress client by method {}.", remoteAddress, authMethod); @@ -406,8 +403,8 @@ remoteAddress, protocolVersionToAdvertise, getRemoteEndpointProtocolVersion(), doAuthentication(clientData); } catch (Exception e) { LOG.warn("[{}] Unable to authenticate: ", remoteAddress, e); - ctx.writeAndFlush(Commands.newError(-1, ServerError.AuthenticationError, "Failed to authenticate")); - close(); + ctx.writeAndFlush(Commands.newError(-1, ServerError.AuthenticationError, "Failed to authenticate")) + .addListener(ChannelFutureListener.CLOSE); } } @@ -428,8 +425,8 @@ protected void handleAuthResponse(CommandAuthResponse authResponse) { } catch (Exception e) { String msg = "Unable to handleAuthResponse"; LOG.warn("[{}] {} ", remoteAddress, msg, e); - ctx.writeAndFlush(Commands.newError(-1, ServerError.AuthenticationError, msg)); - close(); + ctx.writeAndFlush(Commands.newError(-1, ServerError.AuthenticationError, msg)) + .addListener(ChannelFutureListener.CLOSE); } } @@ -462,25 +459,6 @@ protected void handleLookup(CommandLookupTopic lookup) { lookupProxyHandler.handleLookup(lookup); } - private synchronized void close() { - if (state != State.Closed) { - state = State.Closed; - if (directProxyHandler != null && directProxyHandler.outboundChannel != null) { - directProxyHandler.outboundChannel.close(); - directProxyHandler = null; - } - if (connectionPool != null) { - try { - connectionPool.close(); - connectionPool = null; - } catch (Exception e) { - LOG.error("Error closing connection pool", e); - } - } - ctx.close(); - } - } - ClientConfigurationData createClientConfiguration() { ClientConfigurationData clientConf = new ClientConfigurationData(); clientConf.setServiceUrl(service.getServiceUrl());