mirror of
https://gitflic.ru/project/openide/openide.git
synced 2026-09-27 10:03:11 +07:00
convert to kotlin BinaryRequestHandlerTest
This commit is contained in:
@@ -0,0 +1,5 @@
|
||||
<root>
|
||||
<item name='io.netty.bootstrap.AbstractBootstrap B handler(io.netty.channel.ChannelHandler) 0'>
|
||||
<annotation name='org.jetbrains.annotations.NotNull'/>
|
||||
</item>
|
||||
</root>
|
||||
@@ -0,0 +1,11 @@
|
||||
<root>
|
||||
<item name='io.netty.channel.ChannelInitializer void channelRegistered(io.netty.channel.ChannelHandlerContext) 0'>
|
||||
<annotation name='org.jetbrains.annotations.NotNull'/>
|
||||
</item>
|
||||
<item name='io.netty.channel.ChannelInitializer void initChannel(C) 0'>
|
||||
<annotation name='org.jetbrains.annotations.NotNull'/>
|
||||
</item>
|
||||
<item name='io.netty.channel.ChannelPipeline io.netty.channel.ChannelPipeline addLast(io.netty.channel.ChannelHandler...) 0'>
|
||||
<annotation name='org.jetbrains.annotations.NotNull'/>
|
||||
</item>
|
||||
</root>
|
||||
@@ -0,0 +1,150 @@
|
||||
package org.jetbrains.ide
|
||||
|
||||
import io.netty.channel.ChannelInitializer
|
||||
import io.netty.channel.Channel
|
||||
import org.jetbrains.io.Decoder
|
||||
import io.netty.channel.ChannelHandlerContext
|
||||
import io.netty.buffer.ByteBuf
|
||||
import com.intellij.util.Consumer
|
||||
import java.util.UUID
|
||||
import io.netty.channel.ChannelHandler
|
||||
import com.intellij.openapi.util.AsyncResult
|
||||
import io.netty.util.CharsetUtil
|
||||
import org.jetbrains.io.ChannelExceptionHandler
|
||||
import org.jetbrains.io.NettyUtil
|
||||
import com.intellij.util.net.NetUtils
|
||||
import io.netty.buffer.Unpooled
|
||||
import junit.framework.TestCase
|
||||
import org.junit.rules.RuleChain
|
||||
import org.junit.Rule
|
||||
import org.junit.Test
|
||||
|
||||
public class BinaryRequestHandlerTest {
|
||||
private val fixtureManager = FixtureRule()
|
||||
|
||||
private val _chain = RuleChain
|
||||
.outerRule(fixtureManager)
|
||||
|
||||
Rule
|
||||
public fun getChain(): RuleChain = _chain
|
||||
|
||||
Test
|
||||
public fun test() {
|
||||
val text = "Hello!"
|
||||
val result = AsyncResult<String>()
|
||||
|
||||
val bootstrap = NettyUtil.oioClientBootstrap().handler(object : ChannelInitializer<Channel>() {
|
||||
override fun initChannel(channel: Channel) {
|
||||
channel.pipeline().addLast(object : Decoder() {
|
||||
override fun messageReceived(context: ChannelHandlerContext, message: ByteBuf) {
|
||||
val requiredLength = 4 + text.length()
|
||||
val buffer = getBufferIfSufficient(message, requiredLength, context)
|
||||
if (buffer == null) {
|
||||
message.release()
|
||||
}
|
||||
else {
|
||||
val response = buffer.toString(buffer.readerIndex(), requiredLength, CharsetUtil.UTF_8)
|
||||
buffer.skipBytes(requiredLength)
|
||||
buffer.release()
|
||||
result.setDone(response)
|
||||
}
|
||||
}
|
||||
}, ChannelExceptionHandler.getInstance())
|
||||
}
|
||||
})
|
||||
|
||||
val port = BuiltInServerManager.getInstance().waitForStart().getPort()
|
||||
val channel = bootstrap.connect(NetUtils.getLoopbackAddress(), port).syncUninterruptibly().channel()
|
||||
val buffer = channel.alloc().buffer()
|
||||
buffer.writeByte('C'.toInt())
|
||||
buffer.writeByte('H'.toInt())
|
||||
buffer.writeLong(MyBinaryRequestHandler.ID.getMostSignificantBits())
|
||||
buffer.writeLong(MyBinaryRequestHandler.ID.getLeastSignificantBits())
|
||||
|
||||
val message = Unpooled.copiedBuffer(text, CharsetUtil.UTF_8)
|
||||
buffer.writeShort(message.readableBytes())
|
||||
|
||||
channel.write(buffer)
|
||||
channel.writeAndFlush(message).syncUninterruptibly()
|
||||
|
||||
try {
|
||||
result.doWhenRejected(object : Consumer<String> {
|
||||
override fun consume(error: String) {
|
||||
TestCase.fail(error)
|
||||
}
|
||||
})
|
||||
|
||||
TestCase.assertEquals("got-" + text, result.getResultSync(5000))
|
||||
}
|
||||
finally {
|
||||
channel.close()
|
||||
}
|
||||
}
|
||||
|
||||
class MyBinaryRequestHandler : BinaryRequestHandler() {
|
||||
class object {
|
||||
val ID = UUID.fromString("E5068DD6-1DB7-437C-A3FC-3CA53B6E1AC9")
|
||||
}
|
||||
|
||||
override fun getId(): UUID {
|
||||
return ID
|
||||
}
|
||||
|
||||
override fun getInboundHandler(): ChannelHandler {
|
||||
return MyDecoder()
|
||||
}
|
||||
|
||||
private class MyDecoder : Decoder() {
|
||||
private var state = State.HEADER
|
||||
private var contentLength = -1
|
||||
|
||||
private enum class State {
|
||||
HEADER
|
||||
CONTENT
|
||||
}
|
||||
|
||||
override fun messageReceived(context: ChannelHandlerContext, message: ByteBuf) {
|
||||
while (true) {
|
||||
when (state) {
|
||||
State.HEADER -> {
|
||||
run {
|
||||
val buffer = getBufferIfSufficient(message, 2, context)
|
||||
if (buffer == null) {
|
||||
message.release()
|
||||
return
|
||||
}
|
||||
|
||||
contentLength = buffer.readUnsignedShort()
|
||||
state = State.CONTENT
|
||||
}
|
||||
run {
|
||||
val buffer = getBufferIfSufficient(message, contentLength, context)
|
||||
if (buffer == null) {
|
||||
message.release()
|
||||
return
|
||||
}
|
||||
|
||||
val messageText = buffer.toString(buffer.readerIndex(), contentLength, CharsetUtil.UTF_8)
|
||||
buffer.skipBytes(contentLength)
|
||||
state = State.HEADER
|
||||
context.writeAndFlush(Unpooled.copiedBuffer("got-" + messageText, CharsetUtil.UTF_8))
|
||||
}
|
||||
}
|
||||
|
||||
State.CONTENT -> {
|
||||
val buffer = getBufferIfSufficient(message, contentLength, context)
|
||||
if (buffer == null) {
|
||||
message.release()
|
||||
return
|
||||
}
|
||||
val messageText = buffer.toString(buffer.readerIndex(), contentLength, CharsetUtil.UTF_8)
|
||||
buffer.skipBytes(contentLength)
|
||||
state = State.HEADER
|
||||
context.writeAndFlush(Unpooled.copiedBuffer("got-" + messageText, CharsetUtil.UTF_8))
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,140 +0,0 @@
|
||||
package org.jetbrains.ide;
|
||||
|
||||
import com.intellij.openapi.util.AsyncResult;
|
||||
import com.intellij.testFramework.LightPlatformTestCase;
|
||||
import com.intellij.testFramework.PlatformTestCase;
|
||||
import com.intellij.util.Consumer;
|
||||
import com.intellij.util.net.NetUtils;
|
||||
import io.netty.bootstrap.Bootstrap;
|
||||
import io.netty.buffer.ByteBuf;
|
||||
import io.netty.buffer.Unpooled;
|
||||
import io.netty.channel.Channel;
|
||||
import io.netty.channel.ChannelHandler;
|
||||
import io.netty.channel.ChannelHandlerContext;
|
||||
import io.netty.channel.ChannelInitializer;
|
||||
import io.netty.util.CharsetUtil;
|
||||
import junit.framework.TestCase;
|
||||
import org.jetbrains.annotations.NotNull;
|
||||
import org.jetbrains.io.ChannelExceptionHandler;
|
||||
import org.jetbrains.io.Decoder;
|
||||
import org.jetbrains.io.NettyUtil;
|
||||
|
||||
import java.util.UUID;
|
||||
|
||||
public class BinaryRequestHandlerTest extends LightPlatformTestCase {
|
||||
@SuppressWarnings("JUnitTestCaseWithNonTrivialConstructors")
|
||||
public BinaryRequestHandlerTest() {
|
||||
PlatformTestCase.initPlatformLangPrefix();
|
||||
}
|
||||
|
||||
public void test() throws InterruptedException {
|
||||
final String text = "Hello!";
|
||||
final AsyncResult<String> result = new AsyncResult<String>();
|
||||
|
||||
Bootstrap bootstrap = NettyUtil.oioClientBootstrap().handler(new ChannelInitializer() {
|
||||
@Override
|
||||
protected void initChannel(Channel channel) throws Exception {
|
||||
channel.pipeline().addLast(new Decoder() {
|
||||
@Override
|
||||
protected void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf message) throws Exception {
|
||||
int requiredLength = 4 + text.length();
|
||||
ByteBuf buffer = getBufferIfSufficient(message, requiredLength, context);
|
||||
if (buffer == null) {
|
||||
message.release();
|
||||
}
|
||||
else {
|
||||
String response = buffer.toString(buffer.readerIndex(), requiredLength, CharsetUtil.UTF_8);
|
||||
buffer.skipBytes(requiredLength);
|
||||
buffer.release();
|
||||
result.setDone(response);
|
||||
}
|
||||
}
|
||||
}, ChannelExceptionHandler.getInstance());
|
||||
}
|
||||
});
|
||||
|
||||
int port = BuiltInServerManager.getInstance().waitForStart().getPort();
|
||||
Channel channel = bootstrap.connect(NetUtils.getLoopbackAddress(), port).syncUninterruptibly().channel();
|
||||
ByteBuf buffer = channel.alloc().buffer();
|
||||
buffer.writeByte('C');
|
||||
buffer.writeByte('H');
|
||||
buffer.writeLong(MyBinaryRequestHandler.ID.getMostSignificantBits());
|
||||
buffer.writeLong(MyBinaryRequestHandler.ID.getLeastSignificantBits());
|
||||
|
||||
ByteBuf message = Unpooled.copiedBuffer(text, CharsetUtil.UTF_8);
|
||||
buffer.writeShort(message.readableBytes());
|
||||
|
||||
channel.write(buffer);
|
||||
channel.writeAndFlush(message).syncUninterruptibly();
|
||||
|
||||
try {
|
||||
result.doWhenRejected(new Consumer<String>() {
|
||||
@Override
|
||||
public void consume(String error) {
|
||||
TestCase.fail(error);
|
||||
}
|
||||
});
|
||||
|
||||
TestCase.assertEquals("got-" + text, result.getResultSync(5000));
|
||||
}
|
||||
finally {
|
||||
channel.close();
|
||||
}
|
||||
}
|
||||
|
||||
static class MyBinaryRequestHandler extends BinaryRequestHandler {
|
||||
private static final UUID ID = UUID.fromString("E5068DD6-1DB7-437C-A3FC-3CA53B6E1AC9");
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public UUID getId() {
|
||||
return ID;
|
||||
}
|
||||
|
||||
@NotNull
|
||||
@Override
|
||||
public ChannelHandler getInboundHandler() {
|
||||
return new MyDecoder();
|
||||
}
|
||||
|
||||
private static class MyDecoder extends Decoder {
|
||||
private State state = State.HEADER;
|
||||
private int contentLength = -1;
|
||||
|
||||
private enum State {
|
||||
HEADER, CONTENT
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void messageReceived(@NotNull ChannelHandlerContext context, @NotNull ByteBuf message) throws Exception {
|
||||
while (true) {
|
||||
switch (state) {
|
||||
case HEADER: {
|
||||
ByteBuf buffer = getBufferIfSufficient(message, 2, context);
|
||||
if (buffer == null) {
|
||||
message.release();
|
||||
return;
|
||||
}
|
||||
|
||||
contentLength = buffer.readUnsignedShort();
|
||||
state = State.CONTENT;
|
||||
}
|
||||
|
||||
case CONTENT: {
|
||||
ByteBuf buffer = getBufferIfSufficient(message, contentLength, context);
|
||||
if (buffer == null) {
|
||||
message.release();
|
||||
return;
|
||||
}
|
||||
|
||||
String messageText = buffer.toString(buffer.readerIndex(), contentLength, CharsetUtil.UTF_8);
|
||||
buffer.skipBytes(contentLength);
|
||||
state = State.HEADER;
|
||||
context.writeAndFlush(Unpooled.copiedBuffer("got-" + messageText, CharsetUtil.UTF_8));
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user