diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/stream/BasicStreamPacket.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/stream/BasicStreamPacket.java index 0e89aa386..83a15121c 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/stream/BasicStreamPacket.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/stream/BasicStreamPacket.java @@ -29,7 +29,21 @@ public abstract class BasicStreamPacket implements StreamPacket { public static final short PACKET_UNSUPPORT = 142; public static final short UNKNWON_ERROR = 200; + + public static final short ROUTE_TYPE_ERROR = 330; + public static final short ROUTE_TYPE_SERVER_UNSUPPORT = 331; + public static final short ROUTE_TYPE_CLIENT = 336; + public static final short ROUTE_TYPE_UNKOWN = 339; + + public static final short ROUTE_PACKET_ERROR = 340; + public static final short ROUTE_PACKET_UNKNOWN = 341; + public static final short ROUTE_PACKET_UNSUPPORT = 342; + + public static final short ROUTE_NOT_FOUND = 350; + + public static final short ROUTE_CONNECTION_ERROR = 360; + private static final byte[] EMPTY_PAYLOAD = new byte[0]; diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ChannelContext.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ChannelContext.java index b88efbddf..76a519a11 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ChannelContext.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ChannelContext.java @@ -2,6 +2,7 @@ package com.nhn.pinpoint.rpc.server; import java.util.Collections; import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -23,7 +24,7 @@ public class ChannelContext { private final SocketChannelStateChangeEventListener stateChangeEventListener; - private volatile Map channelProperties = Collections.emptyMap(); + private final AtomicReference> properties = new AtomicReference>(); public ChannelContext(SocketChannel socketChannel, StreamChannelManager streamChannelManager) { this(socketChannel, streamChannelManager, DoNothingChannelStateEventListener.INSTANCE); @@ -101,21 +102,16 @@ public class ChannelContext { } public Map getChannelProperties() { - return channelProperties; + Map properties = this.properties.get(); + return properties == null ? Collections.emptyMap() : properties; } - public boolean setChannelProperties(Map properties) { - if (properties == null) { + public boolean setChannelProperties(Map value) { + if (value == null) { return false; } - if (this.channelProperties != Collections.emptyMap()) { - logger.warn("Already Register ChannelProperties.({}).", this.channelProperties); - return false; - } - - this.channelProperties = Collections.unmodifiableMap(properties); - return true; + return this.properties.compareAndSet(null, Collections.unmodifiableMap(value)); } public StreamChannelManager getStreamChannelManager() { diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerSocket.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerSocket.java index f6bd2ab9d..c1498a28e 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerSocket.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerSocket.java @@ -312,7 +312,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { private void sendHandshakeResponseMessage(int requestId, HandshakeResponseCode handShakeResponseCode, Channel channel) { try { - logger.info("write HandshakeResponsePakcet channel:{}, HandshakeResponseCode:{}.", channel, handShakeResponseCode); + logger.info("write HandshakeResponsePakcet. channel:{}, HandshakeResponseCode:{}.", channel, handShakeResponseCode); Map result = new HashMap(); result.put(ControlHandshakeResponsePacket.CODE, handShakeResponseCode.getCode()); diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ClientStreamChannelMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ClientStreamChannelMessageListener.java index abf5f6d9e..103d0ba4c 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ClientStreamChannelMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ClientStreamChannelMessageListener.java @@ -10,6 +10,6 @@ public interface ClientStreamChannelMessageListener { void handleStreamData(ClientStreamChannelContext streamChannelContext, StreamResponsePacket packet); - void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet); + void handleStreamClose(ClientStreamChannelContext streamChannelContext, StreamClosePacket packet); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/DisabledServerStreamChannelMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/DisabledServerStreamChannelMessageListener.java index 2ecc5dfe9..785481d3c 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/DisabledServerStreamChannelMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/DisabledServerStreamChannelMessageListener.java @@ -20,7 +20,7 @@ public class DisabledServerStreamChannelMessageListener implements ServerStreamC } @Override - public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) { + public void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet) { logger.info("{} handleStreamClose unsupported operation. StreamChannel:{}, Packet:{}", this.getClass().getSimpleName(), streamChannelContext, packet); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/LoggingStreamChannelMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/LoggingStreamChannelMessageListener.java index 45397a3b8..ab6f51029 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/LoggingStreamChannelMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/LoggingStreamChannelMessageListener.java @@ -26,7 +26,7 @@ public class LoggingStreamChannelMessageListener { } @Override - public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) { + public void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet) { LOGGER.info("handleStreamClose StreamChannel:{}, Packet:{}", streamChannelContext, packet); } @@ -40,7 +40,7 @@ public class LoggingStreamChannelMessageListener { } @Override - public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) { + public void handleStreamClose(ClientStreamChannelContext streamChannelContext, StreamClosePacket packet) { LOGGER.info("handleStreamClose StreamChannel:{}, Packet:{}", streamChannelContext, packet); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ServerStreamChannelMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ServerStreamChannelMessageListener.java index c18f1f808..baa146fd7 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ServerStreamChannelMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/ServerStreamChannelMessageListener.java @@ -10,6 +10,6 @@ public interface ServerStreamChannelMessageListener { short handleStreamCreate(ServerStreamChannelContext streamChannelContext, StreamCreatePacket packet); - void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet); + void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannel.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannel.java index f3567c9a4..0ebe5f342 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannel.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannel.java @@ -100,11 +100,11 @@ public abstract class StreamChannel { this.streamChannelManager.clearResourceAndSendClose(getStreamId(), BasicStreamPacket.CHANNEL_CLOSE); } - protected Channel getChannel() { + public Channel getChannel() { return channel; } - protected int getStreamId() { + public int getStreamId() { return streamChannelId; } diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/RecordedStreamChannelMessageListener.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/RecordedStreamChannelMessageListener.java index 32a2732a4..dc015810f 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/RecordedStreamChannelMessageListener.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/RecordedStreamChannelMessageListener.java @@ -12,7 +12,6 @@ import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket; import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext; import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener; -import com.nhn.pinpoint.rpc.stream.StreamChannelContext; /** * @author emeroad @@ -38,7 +37,7 @@ public class RecordedStreamChannelMessageListener implements ClientStreamChannel } @Override - public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) { + public void handleStreamClose(ClientStreamChannelContext streamChannelContext, StreamClosePacket packet) { logger.info("handleClose {}, {}", streamChannelContext, packet); receivedMessageList.add(packet.getPayload()); latch.countDown(); diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java index 02adb1693..11a0393f7 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java @@ -308,8 +308,8 @@ public class StreamChannelManagerTest { } @Override - public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) { - bo.removeServerStreamChannelContext((ServerStreamChannelContext) streamChannelContext); + public void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet) { + bo.removeServerStreamChannelContext(streamChannelContext); } }