From 4a67242cd9d6c2525e00722ff27245425f17c0d1 Mon Sep 17 00:00:00 2001 From: Vladimir Krivosheev Date: Fri, 10 Jun 2016 15:13:22 +0200 Subject: [PATCH] WEB-21991 New Node.js debug protocol incompatibilities --- .../server/BuildMessageDispatcher.java | 8 +-- .../jps/javac/ExternalJavacManager.java | 8 +-- .../netty/io/netty/channel/annotations.xml | 6 ++ .../ide/BuiltInServerManagerImpl.java | 8 +++ .../org/jetbrains/io/jsonRpc/ClientManager.kt | 8 +-- .../socket/RpcBinaryRequestHandler.java | 6 +- .../io/webSocket/MessageChannelHandler.java | 66 ++++++++++-------- .../webSocket/WebSocketHandshakeHandler.java | 4 +- .../io/webSocket/WebSocketProtocolHandler.kt | 67 +++++++++++++++++++ .../io/DelegatingHttpRequestHandler.kt | 6 +- .../src/org/jetbrains/io/NettyUtil.java | 2 +- .../src/org/jetbrains/io/netty.kt | 2 +- .../debugger-ui/src/RemoteVmConnection.kt | 15 +++-- 13 files changed, 151 insertions(+), 55 deletions(-) create mode 100644 platform/built-in-server/src/org/jetbrains/io/webSocket/WebSocketProtocolHandler.kt diff --git a/java/compiler/impl/src/com/intellij/compiler/server/BuildMessageDispatcher.java b/java/compiler/impl/src/com/intellij/compiler/server/BuildMessageDispatcher.java index 096a113cf7f9..2624cde1af17 100644 --- a/java/compiler/impl/src/com/intellij/compiler/server/BuildMessageDispatcher.java +++ b/java/compiler/impl/src/com/intellij/compiler/server/BuildMessageDispatcher.java @@ -1,5 +1,5 @@ /* - * Copyright 2000-2014 JetBrains s.r.o. + * Copyright 2000-2016 JetBrains s.r.o. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -123,7 +123,7 @@ class BuildMessageDispatcher extends SimpleChannelInboundHandlerAdapter { @Override public void channelUnregistered(ChannelHandlerContext ctx) throws Exception { - JavacProcessDescriptor descriptor = ctx.attr(SESSION_DESCRIPTOR).get(); + JavacProcessDescriptor descriptor = ctx.channel().attr(SESSION_DESCRIPTOR).get(); if (descriptor != null) { descriptor.setDone(); } @@ -335,7 +335,7 @@ public class ExternalJavacManager { @Override public void channelRead0(final ChannelHandlerContext context, JavacRemoteProto.Message message) throws Exception { - JavacProcessDescriptor descriptor = context.attr(SESSION_DESCRIPTOR).get(); + JavacProcessDescriptor descriptor = context.channel().attr(SESSION_DESCRIPTOR).get(); UUID sessionId; if (descriptor == null) { @@ -345,7 +345,7 @@ public class ExternalJavacManager { descriptor = myMessageHandlers.get(sessionId); if (descriptor != null) { descriptor.channel = context.channel(); - context.attr(SESSION_DESCRIPTOR).set(descriptor); + context.channel().attr(SESSION_DESCRIPTOR).set(descriptor); } } else { diff --git a/lib/annotations/netty/io/netty/channel/annotations.xml b/lib/annotations/netty/io/netty/channel/annotations.xml index 13f5c446cf4a..121521d60f9d 100644 --- a/lib/annotations/netty/io/netty/channel/annotations.xml +++ b/lib/annotations/netty/io/netty/channel/annotations.xml @@ -5,6 +5,12 @@ + + + + + + diff --git a/platform/built-in-server/src/org/jetbrains/ide/BuiltInServerManagerImpl.java b/platform/built-in-server/src/org/jetbrains/ide/BuiltInServerManagerImpl.java index 6030acfc5e50..a080eb512141 100644 --- a/platform/built-in-server/src/org/jetbrains/ide/BuiltInServerManagerImpl.java +++ b/platform/built-in-server/src/org/jetbrains/ide/BuiltInServerManagerImpl.java @@ -14,6 +14,7 @@ import com.intellij.openapi.util.text.StringUtil; import com.intellij.util.Url; import com.intellij.util.UrlImpl; import com.intellij.util.net.NetUtils; +import io.netty.channel.EventLoopGroup; import io.netty.channel.oio.OioEventLoopGroup; import org.jetbrains.annotations.NonNls; import org.jetbrains.annotations.NotNull; @@ -123,6 +124,13 @@ public class BuiltInServerManagerImpl extends BuiltInServerManager { return server; } + @NotNull + public EventLoopGroup getEventLoopGroup() { + waitForStart(); + assert server != null; + return server.getEventLoopGroup(); + } + @Override public boolean isOnBuiltInWebServer(@Nullable Url url) { return url != null && !StringUtil.isEmpty(url.getAuthority()) && isOnBuiltInWebServerByAuthority(url.getAuthority()); diff --git a/platform/built-in-server/src/org/jetbrains/io/jsonRpc/ClientManager.kt b/platform/built-in-server/src/org/jetbrains/io/jsonRpc/ClientManager.kt index 4eb035b48410..10561b9cbec8 100644 --- a/platform/built-in-server/src/org/jetbrains/io/jsonRpc/ClientManager.kt +++ b/platform/built-in-server/src/org/jetbrains/io/jsonRpc/ClientManager.kt @@ -5,7 +5,7 @@ import com.intellij.openapi.util.SimpleTimer import gnu.trove.THashSet import gnu.trove.TObjectProcedure import io.netty.buffer.ByteBuf -import io.netty.channel.ChannelHandlerContext +import io.netty.channel.Channel import io.netty.util.AttributeKey import org.jetbrains.concurrency.Promise import org.jetbrains.io.webSocket.WebSocketServerOptions @@ -64,7 +64,7 @@ class ClientManager(private val listener: ClientListener?, val exceptionHandler: }) } - fun disconnectClient(context: ChannelHandlerContext, client: Client, closeChannel: Boolean): Boolean { + fun disconnectClient(channel: Channel, client: Client, closeChannel: Boolean): Boolean { synchronized (clients) { if (!clients.remove(client)) { return false @@ -72,10 +72,10 @@ class ClientManager(private val listener: ClientListener?, val exceptionHandler: } try { - context.attr(CLIENT).remove() + channel.attr(CLIENT).remove() if (closeChannel) { - context.channel().close() + channel.close() } client.rejectAsyncResults(exceptionHandler) diff --git a/platform/built-in-server/src/org/jetbrains/io/jsonRpc/socket/RpcBinaryRequestHandler.java b/platform/built-in-server/src/org/jetbrains/io/jsonRpc/socket/RpcBinaryRequestHandler.java index 8f0b44c7cfcd..5dfc75111434 100644 --- a/platform/built-in-server/src/org/jetbrains/io/jsonRpc/socket/RpcBinaryRequestHandler.java +++ b/platform/built-in-server/src/org/jetbrains/io/jsonRpc/socket/RpcBinaryRequestHandler.java @@ -58,7 +58,7 @@ public class RpcBinaryRequestHandler extends BinaryRequestHandler implements Exc @Override public ChannelHandler getInboundHandler(@NotNull ChannelHandlerContext context) { SocketClient client = new SocketClient(context.channel()); - context.attr(ClientManagerKt.getCLIENT()).set(client); + context.channel().attr(ClientManagerKt.getCLIENT()).set(client); clientManager.getValue().addClient(client); connected(client, null); return new MyDecoder(client); @@ -127,10 +127,10 @@ public class RpcBinaryRequestHandler extends BinaryRequestHandler implements Exc @Override public void channelInactive(ChannelHandlerContext context) throws Exception { - Client client = context.attr(ClientManagerKt.getCLIENT()).get(); + Client client = context.channel().attr(ClientManagerKt.getCLIENT()).get(); // if null, so, has already been explicitly removed if (client != null) { - clientManager.getValue().disconnectClient(context, client, false); + clientManager.getValue().disconnectClient(context.channel(), client, false); } } } diff --git a/platform/built-in-server/src/org/jetbrains/io/webSocket/MessageChannelHandler.java b/platform/built-in-server/src/org/jetbrains/io/webSocket/MessageChannelHandler.java index 1cb37b07287a..8a94ed4807b3 100644 --- a/platform/built-in-server/src/org/jetbrains/io/webSocket/MessageChannelHandler.java +++ b/platform/built-in-server/src/org/jetbrains/io/webSocket/MessageChannelHandler.java @@ -1,18 +1,19 @@ package org.jetbrains.io.webSocket; +import io.netty.channel.Channel; import io.netty.channel.ChannelHandler; import io.netty.channel.ChannelHandlerContext; -import io.netty.handler.codec.http.websocketx.*; +import io.netty.handler.codec.http.websocketx.CloseWebSocketFrame; +import io.netty.handler.codec.http.websocketx.TextWebSocketFrame; import org.jetbrains.annotations.NotNull; import org.jetbrains.io.ChannelBufferToString; -import org.jetbrains.io.SimpleChannelInboundHandlerAdapter; import org.jetbrains.io.jsonRpc.Client; import org.jetbrains.io.jsonRpc.ClientManager; import org.jetbrains.io.jsonRpc.ClientManagerKt; import org.jetbrains.io.jsonRpc.MessageServer; @ChannelHandler.Sharable -final class MessageChannelHandler extends SimpleChannelInboundHandlerAdapter { +final class MessageChannelHandler extends WebSocketProtocolHandler { private final ClientManager clientManager; private final MessageServer messageServer; @@ -22,51 +23,62 @@ final class MessageChannelHandler extends SimpleChannelInboundHandlerAdapter ReferenceCountUtil.release(message) + is PingWebSocketFrame -> context.channel().writeAndFlush(PongWebSocketFrame(message.content())) + is CloseWebSocketFrame -> closeFrameReceived(context.channel(), message) + is TextWebSocketFrame -> { + try { + textFrameReceived(context.channel(), message) + } + finally { + // client should release buffer as soon as possible, so, message could be released already + if (message.refCnt() > 0) { + message.release() + } + } + } + else -> throw UnsupportedOperationException("${message.javaClass.name} frame types not supported") + } + } + + abstract protected fun textFrameReceived(channel: Channel, message: TextWebSocketFrame) + + protected open fun closeFrameReceived(channel: Channel, message: CloseWebSocketFrame) { + channel.close() + } + + @Suppress("OverridingDeprecatedMember") + override fun exceptionCaught(context: ChannelHandlerContext, cause: Throwable) { + NettyUtil.logAndClose(cause, LOG, context.channel()) + } +} + +open class WebSocketProtocolHandshakeHandler(private val handshaker: WebSocketClientHandshaker) : ChannelInboundHandlerAdapter() { + override final fun channelRead(context: ChannelHandlerContext, message: Any) { + val channel = context.channel() + if (!handshaker.isHandshakeComplete) { + handshaker.finishHandshake(channel, message as FullHttpResponse) + context.pipeline().remove(this) + completed() + return + } + + if (message is FullHttpResponse) { + throw IllegalStateException("Unexpected FullHttpResponse (getStatus=${message.status()}, content=${message.content().toString(CharsetUtil.UTF_8)})") + } + + context.fireChannelRead(message) + } + + open protected fun completed() { + } +} \ No newline at end of file diff --git a/platform/platform-impl/src/org/jetbrains/io/DelegatingHttpRequestHandler.kt b/platform/platform-impl/src/org/jetbrains/io/DelegatingHttpRequestHandler.kt index 806d3e6a351b..ee17b364b0d7 100644 --- a/platform/platform-impl/src/org/jetbrains/io/DelegatingHttpRequestHandler.kt +++ b/platform/platform-impl/src/org/jetbrains/io/DelegatingHttpRequestHandler.kt @@ -41,7 +41,7 @@ internal class DelegatingHttpRequestHandler : DelegatingHttpRequestHandlerBase() return isSupported(request) && !request.isWriteFromBrowserWithoutOrigin() && isAccessible(request) && process(urlDecoder, request, context) } - val prevHandlerAttribute = context.attr(PREV_HANDLER) + val prevHandlerAttribute = context.channel().attr(PREV_HANDLER) val connectedHandler = prevHandlerAttribute.get() if (connectedHandler != null) { if (connectedHandler.checkAndProcess()) { @@ -79,11 +79,13 @@ internal class DelegatingHttpRequestHandler : DelegatingHttpRequestHandlerBase() return false } + @Suppress("OverridingDeprecatedMember") override fun exceptionCaught(context: ChannelHandlerContext, cause: Throwable) { try { - context.attr(PREV_HANDLER).remove() + context.channel().attr(PREV_HANDLER).remove() } finally { + @Suppress("DEPRECATION") super.exceptionCaught(context, cause) } } diff --git a/platform/platform-impl/src/org/jetbrains/io/NettyUtil.java b/platform/platform-impl/src/org/jetbrains/io/NettyUtil.java index 33bd2693a253..7b779dfbc95a 100644 --- a/platform/platform-impl/src/org/jetbrains/io/NettyUtil.java +++ b/platform/platform-impl/src/org/jetbrains/io/NettyUtil.java @@ -85,7 +85,7 @@ public final class NettyUtil { int maxAttemptCount, @NotNull Condition stopCondition) throws Throwable { int attemptCount = 0; - if (bootstrap.group() instanceof NioEventLoopGroup) { + if (bootstrap.config().group() instanceof NioEventLoopGroup) { return connectNio(bootstrap, remoteAddress, promise, maxAttemptCount, stopCondition, attemptCount); } diff --git a/platform/platform-impl/src/org/jetbrains/io/netty.kt b/platform/platform-impl/src/org/jetbrains/io/netty.kt index 044302fdb578..653910718ef2 100644 --- a/platform/platform-impl/src/org/jetbrains/io/netty.kt +++ b/platform/platform-impl/src/org/jetbrains/io/netty.kt @@ -68,7 +68,7 @@ fun oioClientBootstrap(): Bootstrap { } inline fun ChannelFuture.addChannelListener(crossinline listener: (future: ChannelFuture) -> Unit) { - addListener(GenericFutureListener { listener(it) }) + addListener(GenericFutureListener { listener(it) }) } // if NIO, so, it is shared and we must not shutdown it diff --git a/platform/script-debugger/debugger-ui/src/RemoteVmConnection.kt b/platform/script-debugger/debugger-ui/src/RemoteVmConnection.kt index cfd13cf1a3b7..43f91b442d6e 100644 --- a/platform/script-debugger/debugger-ui/src/RemoteVmConnection.kt +++ b/platform/script-debugger/debugger-ui/src/RemoteVmConnection.kt @@ -23,13 +23,14 @@ import com.intellij.ui.ColoredListCellRenderer import com.intellij.ui.components.JBList import com.intellij.util.io.socketConnection.ConnectionStatus import io.netty.bootstrap.Bootstrap +import io.netty.channel.ChannelFuture +import io.netty.util.concurrent.GenericFutureListener import org.jetbrains.concurrency.AsyncPromise import org.jetbrains.concurrency.Promise import org.jetbrains.concurrency.rejectedPromise import org.jetbrains.concurrency.resolvedPromise import org.jetbrains.debugger.Vm import org.jetbrains.io.NettyUtil -import org.jetbrains.io.addChannelListener import org.jetbrains.io.connect import org.jetbrains.rpc.LOG import java.net.ConnectException @@ -43,12 +44,16 @@ abstract class RemoteVmConnection : VmConnection() { private val connectCancelHandler = AtomicReference<() -> Unit>() + protected val channelCloseListener = GenericFutureListener { + close("Process disconnected unexpectedly", ConnectionStatus.DISCONNECTED) + } + abstract fun createBootstrap(address: InetSocketAddress, vmResult: AsyncPromise): Bootstrap @JvmOverloads fun open(address: InetSocketAddress, stopCondition: Condition? = null): Promise { port = address.port - isLocalAddress = address.getAddress().isAnyLocalAddress() || address.getAddress().isLoopbackAddress() + isLocalAddress = address.address.isAnyLocalAddress || address.address.isLoopbackAddress setState(ConnectionStatus.WAITING_FOR_CONNECTION, "Connecting to ${address.hostName}:${port}") val result = AsyncPromise() val future = ApplicationManager.getApplication().executeOnPooledThread { @@ -77,11 +82,7 @@ abstract class RemoteVmConnection : VmConnection() { createBootstrap(address, result) .connect(address, connectionPromise, maxAttemptCount = if (stopCondition == null) NettyUtil.DEFAULT_CONNECT_ATTEMPT_COUNT else -1, stopCondition = stopCondition) - ?.let { - it.closeFuture().addChannelListener { - close("Process disconnected unexpectedly", ConnectionStatus.DISCONNECTED) - } - } + ?.let { it.closeFuture().addListener(channelCloseListener) } } connectCancelHandler.set {