PY-18029 Use DirectedMessageHandler both in TNettyClientTransport and TNettyServerTransport

This commit is contained in:
Alexander Koshevoy
2018-08-22 23:16:40 +03:00
parent 2d772fcdc7
commit 6d50c5978b
4 changed files with 8 additions and 19 deletions
@@ -1,6 +1,5 @@
package com.jetbrains.python.console.thrift.client
package com.jetbrains.python.console.thrift
import com.jetbrains.python.console.thrift.DirectedMessage
import com.jetbrains.python.console.thrift.DirectedMessage.MessageDirection.REQUEST
import com.jetbrains.python.console.thrift.DirectedMessage.MessageDirection.RESPONSE
import io.netty.buffer.ByteBuf
@@ -42,7 +42,7 @@ To implement this idea let each message be flagged whether the message is a requ
On Java side the message with direction flag is `com.jetbrains.python.console.thrift.DirectedMessage`. The incoming message is parsed into
`DirectedMessage` using `com.jetbrains.python.console.thrift.DirectedMessageCodec`. The message content is dispatched then via
`com.jetbrains.python.console.thrift.client.DirectedMessageHandler` either to request or response stream and will be handled accordingly by
`com.jetbrains.python.console.thrift.DirectedMessageHandler` either to request or response stream and will be handled accordingly by
*server-side Thrift service* or *client-side Thrift service*.
*Netty* handlers are asynchronous so that requests for *server-side service* and responses for *client-side service* are processed
@@ -2,6 +2,7 @@ package com.jetbrains.python.console.thrift.client
import com.jetbrains.python.console.thrift.DirectedMessage
import com.jetbrains.python.console.thrift.DirectedMessageCodec
import com.jetbrains.python.console.thrift.DirectedMessageHandler
import com.jetbrains.python.console.thrift.TCumulativeTransport
import io.netty.bootstrap.Bootstrap
import io.netty.buffer.ByteBuf
@@ -4,12 +4,13 @@ import com.intellij.openapi.diagnostic.Logger
import com.intellij.util.ConcurrencyUtil
import com.jetbrains.python.console.thrift.DirectedMessage
import com.jetbrains.python.console.thrift.DirectedMessageCodec
import com.jetbrains.python.console.thrift.DirectedMessageHandler
import com.jetbrains.python.console.thrift.TCumulativeTransport
import io.netty.bootstrap.ServerBootstrap
import io.netty.channel.ChannelHandlerContext
import io.netty.channel.ChannelInboundHandlerAdapter
import io.netty.channel.ChannelInitializer
import io.netty.channel.ChannelOption
import io.netty.channel.SimpleChannelInboundHandler
import io.netty.channel.nio.NioEventLoopGroup
import io.netty.channel.socket.SocketChannel
import io.netty.channel.socket.nio.NioServerSocketChannel
@@ -109,26 +110,14 @@ class TNettyServerTransport(port: Int) : TServerTransport() {
val thriftTransport = TNettyTransport(ch)
val reverseTransport = TNettyClientTransport(ch)
ch.pipeline().addLast(object : SimpleChannelInboundHandler<DirectedMessage>() {
override fun channelRead0(ctx: ChannelHandlerContext, msg: DirectedMessage) {
when (msg.direction) {
DirectedMessage.MessageDirection.REQUEST -> {
thriftTransport.outputStream
}
DirectedMessage.MessageDirection.RESPONSE -> {
reverseTransport.outputStream
}
}.let {
it.write(msg.content)
it.flush()
}
}
ch.pipeline().addLast(DirectedMessageHandler(reverseTransport.outputStream, thriftTransport.outputStream))
ch.pipeline().addLast(object: ChannelInboundHandlerAdapter() {
override fun channelInactive(ctx: ChannelHandlerContext) {
thriftTransport.close()
reverseTransport.close()
ctx.fireChannelInactive()
super.channelInactive(ctx)
}
})