From 284bc433d1a21833f19be2dc66a7f166cdaf86da Mon Sep 17 00:00:00 2001 From: koo-taejin Date: Fri, 17 Oct 2014 11:56:31 +0900 Subject: [PATCH] =?UTF-8?q?#22=20Cluster=EA=B0=84=20Stream=20=EA=B8=B0?= =?UTF-8?q?=EB=8A=A5=20=EA=B3=A0=EC=95=88=20=EB=B0=8F=20=EA=B0=9C=EB=B0=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 이전에 사용하던 StreamChannel 관련 코드 제거 및 신규 StreamChannel로 변경 2. MessageListener에서 StreamChnnel 관련 처리를 분리 3. PinpointSocketFactory의 인터페이스 변경 4. 단순 테스트 코드 수정 --- .../collector/cluster/WebClusterPoint.java | 10 +- .../collector/receiver/tcp/TCPReceiver.java | 7 - .../cluster/ClusterPointRouterTest.java | 7 - .../ZookeeperProfilerClusterStressTest.java | 5 +- .../cluster/zookeeper/ZookeeperTestUtils.java | 7 - .../pinpoint/profiler/DefaultAgent.java | 25 +- .../tools/NetworkAvailabilityChecker.java | 7 +- .../profiler/HeartBeatCheckerTest.java | 15 +- .../profiler/HeartBitCheckerStressTest.java | 15 +- .../sender/TcpDataSenderReconnectTest.java | 15 +- .../profiler/sender/TcpDataSenderTest.java | 17 +- .../pinpoint/profiler/util/MockAgent.java | 10 - .../RequestResponseServerMessageListener.java | 8 - .../pinpoint/rpc/client/PinpointSocket.java | 32 +-- .../rpc/client/PinpointSocketFactory.java | 48 ++-- .../rpc/client/PinpointSocketHandler.java | 115 +++++++-- .../client/ReconnectStateSocketHandler.java | 23 +- .../client/SocketClientPipelineFactory.java | 3 + .../pinpoint/rpc/client/SocketHandler.java | 13 +- .../pinpoint/rpc/client/StreamChannel.java | 225 ---------------- .../rpc/client/StreamChannelManager.java | 81 ------ .../pinpoint/rpc/server/ChannelContext.java | 23 +- .../rpc/server/PinpointServerSocket.java | 45 ++-- .../rpc/server/ServerMessageListener.java | 3 - .../rpc/server/ServerStreamChannel.java | 159 ------------ .../server/ServerStreamChannelManager.java | 65 ----- .../SimpleLoggingServerMessageListener.java | 7 - .../stream/StreamChannelMessageListener.java | 19 -- .../RecordedStreamChannelMessageListener.java | 25 +- .../rpc/client/PinpointSocketFactoryTest.java | 49 ---- .../rpc/server/ControlPacketServerTest.java | 7 - .../pinpoint/rpc/server/EventListnerTest.java | 6 - .../rpc/server/MessageListenerTest.java | 55 ++-- .../rpc/server/TestSeverMessageListener.java | 34 --- .../rpc/stream/StreamChannelManagerTest.java | 244 ++++++++++++++++++ .../web/server/PinpointSocketManager.java | 7 +- .../pinpoint/web/cluster/ClusterTest.java | 4 +- 37 files changed, 511 insertions(+), 929 deletions(-) delete mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannel.java delete mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannelManager.java delete mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannel.java delete mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannelManager.java delete mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannelMessageListener.java create mode 100644 rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/WebClusterPoint.java b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/WebClusterPoint.java index 840b0ca7d..bf7d089aa 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/WebClusterPoint.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/WebClusterPoint.java @@ -22,16 +22,12 @@ public class WebClusterPoint implements ClusterPoint { private final Logger logger = LoggerFactory.getLogger(this.getClass()); private final PinpointSocketFactory factory; - private final MessageListener messageListener; - // InetSocketAddress List로 전달 하는게 좋을거 같은데 이걸 Key로만들기가 쉽지 않네; - private final Map clusterRepository = new HashMap(); public WebClusterPoint(String id, MessageListener messageListener) { - this.messageListener = messageListener; - this.factory = new PinpointSocketFactory(); this.factory.setTimeoutMillis(1000 * 5); + this.factory.setMessageListener(messageListener); Map properties = new HashMap(); properties.put("id", id); @@ -74,7 +70,7 @@ public class WebClusterPoint implements ClusterPoint { PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); logger.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -82,7 +78,7 @@ public class WebClusterPoint implements ClusterPoint { } } logger.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java index d0877875d..34d5ffc91 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java @@ -29,10 +29,8 @@ import com.nhn.pinpoint.common.util.PinpointThreadFactory; import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; import com.nhn.pinpoint.thrift.io.DeserializerFactory; import com.nhn.pinpoint.thrift.io.Header; @@ -136,11 +134,6 @@ public class TCPReceiver { requestResponse(requestPacket, channel); } - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.warn("unsupported streamPacket received {}", streamPacket); - } - @Override public int handleEnableWorker(Map properties) { if (properties == null) { diff --git a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest.java b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest.java index 95bd1a1e4..e48b8f201 100644 --- a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest.java +++ b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest.java @@ -18,11 +18,9 @@ import com.nhn.pinpoint.collector.receiver.tcp.AgentProperties; import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.ChannelContext; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; @RunWith(SpringJUnit4ClassRunner.class) @@ -88,11 +86,6 @@ public class ClusterPointRouterTest { logger.warn("Unsupport request received {} {}", requestPacket, channel); } - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.warn("unsupported streamPacket received {}", streamPacket); - } - @Override public int handleEnableWorker(Map properties) { logger.warn("do handleEnableWorker {}", properties); diff --git a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterStressTest.java b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterStressTest.java index eb397597b..89889f0a4 100644 --- a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterStressTest.java +++ b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterStressTest.java @@ -216,6 +216,7 @@ public class ZookeeperProfilerClusterStressTest { this.factory = new PinpointSocketFactory(); this.factory.setProperties(properties); + this.factory.setMessageListener(messageListener); } private void connect(InetSocketAddress address) { @@ -249,7 +250,7 @@ public class ZookeeperProfilerClusterStressTest { PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); logger.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -257,7 +258,7 @@ public class ZookeeperProfilerClusterStressTest { } } logger.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperTestUtils.java b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperTestUtils.java index af8550fbc..c5958f84f 100644 --- a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperTestUtils.java +++ b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperTestUtils.java @@ -13,9 +13,7 @@ import com.nhn.pinpoint.rpc.client.MessageListener; import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; final class ZookeeperTestUtils { @@ -86,11 +84,6 @@ final class ZookeeperTestUtils { LOGGER.warn("Unsupport request received {} {}", requestPacket, channel); } - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - LOGGER.warn("unsupported streamPacket received {}", streamPacket); - } - @Override public int handleEnableWorker(Map properties) { LOGGER.warn("do handleEnableWorker {}", properties); diff --git a/profiler/src/main/java/com/navercorp/pinpoint/profiler/DefaultAgent.java b/profiler/src/main/java/com/navercorp/pinpoint/profiler/DefaultAgent.java index 7f611ff65..559b0a4d5 100644 --- a/profiler/src/main/java/com/navercorp/pinpoint/profiler/DefaultAgent.java +++ b/profiler/src/main/java/com/navercorp/pinpoint/profiler/DefaultAgent.java @@ -137,9 +137,8 @@ public class DefaultAgent implements Agent { this.tAgentInfo = createTAgentInfo(); - this.factory = createPinpointSocketFactory(); - this.socket = createPinpointSocket(this.profilerConfig.getCollectorServerIp(), this.profilerConfig.getCollectorTcpServerPort(), factory, - this.profilerConfig.isTcpDataSenderCommandAcceptEnable()); + this.factory = createPinpointSocketFactory(this.profilerConfig.isTcpDataSenderCommandAcceptEnable()); + this.socket = createPinpointSocket(this.profilerConfig.getCollectorServerIp(), this.profilerConfig.getCollectorTcpServerPort(), factory); this.tcpDataSender = createTcpDataSender(socket); @@ -285,7 +284,7 @@ public class DefaultAgent implements Agent { return serverMetaDataHolder; } - protected PinpointSocketFactory createPinpointSocketFactory() { + protected PinpointSocketFactory createPinpointSocketFactory(boolean isSupportServerMode) { Map properties = this.agentInformation.toMap(); properties.put(AgentPropertiesType.IP.getName(), serverInfo.getHostip()); @@ -293,28 +292,22 @@ public class DefaultAgent implements Agent { pinpointSocketFactory.setTimeoutMillis(1000 * 5); pinpointSocketFactory.setProperties(properties); + if (isSupportServerMode) { + pinpointSocketFactory.setMessageListener(new CommandDispatcher()); + } + return pinpointSocketFactory; } protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { - return createPinpointSocket(host, port, factory, false); - } - - protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory, boolean useMessageListener) { // 1.2 버전이 Tcp Data Command 허용하는 버전이 아니기 떄문에 true이던 false이던 무조건 SimpleLoggingMessageListener를 이용하게 함 // SimpleLoggingMessageListener.LISTENER 는 서로 통신을 하지 않게 설정되어 있음 (테스트코드는 pinpoint-rpc에 존재) // 1.3 버전으로 할 경우 아래 분기에서 MessageListener 변경 필요 - MessageListener messageListener = null; - if (useMessageListener) { - messageListener = new CommandDispatcher(); - } else { - messageListener = SimpleLoggingMessageListener.LISTENER; - } PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); logger.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -322,7 +315,7 @@ public class DefaultAgent implements Agent { } } logger.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/profiler/src/main/java/com/navercorp/pinpoint/profiler/tools/NetworkAvailabilityChecker.java b/profiler/src/main/java/com/navercorp/pinpoint/profiler/tools/NetworkAvailabilityChecker.java index c66204e10..92c2af271 100644 --- a/profiler/src/main/java/com/navercorp/pinpoint/profiler/tools/NetworkAvailabilityChecker.java +++ b/profiler/src/main/java/com/navercorp/pinpoint/profiler/tools/NetworkAvailabilityChecker.java @@ -93,18 +93,17 @@ public class NetworkAvailabilityChecker implements PinpointTools { PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); pinpointSocketFactory.setTimeoutMillis(1000 * 5); pinpointSocketFactory.setProperties(Collections.emptyMap()); + pinpointSocketFactory.setMessageListener(new CommandDispatcher()); return pinpointSocketFactory; } private static PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { - MessageListener messageListener = new CommandDispatcher(); - PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); LOGGER.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -112,7 +111,7 @@ public class NetworkAvailabilityChecker implements PinpointTools { } } LOGGER.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBeatCheckerTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBeatCheckerTest.java index d60c446fe..952efafbf 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBeatCheckerTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBeatCheckerTest.java @@ -15,15 +15,12 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.profiler.sender.TcpDataSender; import com.nhn.pinpoint.rpc.PinpointSocketException; -import com.nhn.pinpoint.rpc.client.MessageListener; import com.nhn.pinpoint.rpc.client.PinpointSocket; import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; import com.nhn.pinpoint.thrift.dto.TAgentInfo; import com.nhn.pinpoint.thrift.dto.TResult; @@ -210,11 +207,6 @@ public class HeartBeatCheckerTest { } } - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.info("handleStreamPacket:{}", streamPacket); - } - @Override public int handleEnableWorker(Map arg0) { return 0; @@ -225,18 +217,17 @@ public class HeartBeatCheckerTest { PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); pinpointSocketFactory.setTimeoutMillis(1000 * 5); pinpointSocketFactory.setProperties(Collections.EMPTY_MAP); + pinpointSocketFactory.setMessageListener(new CommandDispatcher()); return pinpointSocketFactory; } private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { - MessageListener messageListener = new CommandDispatcher(); - PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); logger.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -244,7 +235,7 @@ public class HeartBeatCheckerTest { } } logger.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBitCheckerStressTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBitCheckerStressTest.java index 82b663bdc..90c85b849 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBitCheckerStressTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/HeartBitCheckerStressTest.java @@ -12,15 +12,12 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.profiler.sender.TcpDataSender; import com.nhn.pinpoint.rpc.PinpointSocketException; -import com.nhn.pinpoint.rpc.client.MessageListener; import com.nhn.pinpoint.rpc.client.PinpointSocket; import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; import com.nhn.pinpoint.thrift.dto.TAgentInfo; import com.nhn.pinpoint.thrift.dto.TResult; @@ -160,11 +157,6 @@ public class HeartBitCheckerStressTest { e.printStackTrace(); } } - - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.info("handleStreamPacket:{}", streamPacket); - } @Override public int handleEnableWorker(Map arg0) { @@ -176,18 +168,17 @@ public class HeartBitCheckerStressTest { PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); pinpointSocketFactory.setTimeoutMillis(1000 * 5); pinpointSocketFactory.setProperties(Collections.EMPTY_MAP); + pinpointSocketFactory.setMessageListener(new CommandDispatcher()); return pinpointSocketFactory; } private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { - MessageListener messageListener = new CommandDispatcher(); - PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); logger.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -195,7 +186,7 @@ public class HeartBitCheckerStressTest { } } logger.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java index 975e36818..b40e20c10 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java @@ -9,15 +9,12 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.rpc.PinpointSocketException; -import com.nhn.pinpoint.rpc.client.MessageListener; import com.nhn.pinpoint.rpc.client.PinpointSocket; import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; import com.nhn.pinpoint.thrift.dto.TApiMetaData; @@ -48,11 +45,6 @@ public class TcpDataSenderReconnectTest { logger.info("handleRequest:{}", requestPacket); } - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.info("handleStreamPacket:{}", streamPacket); - } - @Override public int handleEnableWorker(Map properties) { return 0; @@ -95,18 +87,17 @@ public class TcpDataSenderReconnectTest { PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); pinpointSocketFactory.setTimeoutMillis(1000 * 5); pinpointSocketFactory.setProperties(Collections.EMPTY_MAP); + pinpointSocketFactory.setMessageListener(new CommandDispatcher()); return pinpointSocketFactory; } private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { - MessageListener messageListener = new CommandDispatcher(); - PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); logger.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -114,7 +105,7 @@ public class TcpDataSenderReconnectTest { } } logger.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java index cdd43008d..2e5bc835d 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java @@ -15,15 +15,12 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.rpc.PinpointSocketException; -import com.nhn.pinpoint.rpc.client.MessageListener; import com.nhn.pinpoint.rpc.client.PinpointSocket; import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; import com.nhn.pinpoint.thrift.dto.TApiMetaData; @@ -57,11 +54,6 @@ public class TcpDataSenderTest { public void handleRequest(RequestPacket requestPacket, SocketChannel channel) { logger.info("handleRequest:{}", requestPacket); } - - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.info("handleStreamPacket:{}", streamPacket); - } @Override public int handleEnableWorker(Map arg0) { @@ -83,6 +75,8 @@ public class TcpDataSenderTest { this.sendLatch = new CountDownLatch(2); PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + socketFactory.setMessageListener(new CommandDispatcher()); + PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory); TcpDataSender sender = new TcpDataSender(socket); @@ -113,15 +107,12 @@ public class TcpDataSenderTest { return pinpointSocketFactory; } - private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { - MessageListener messageListener = new CommandDispatcher(); - PinpointSocket socket = null; for (int i = 0; i < 3; i++) { try { - socket = factory.connect(host, port, messageListener); + socket = factory.connect(host, port); logger.info("tcp connect success:{}/{}", host, port); return socket; } catch (PinpointSocketException e) { @@ -129,7 +120,7 @@ public class TcpDataSenderTest { } } logger.warn("change background tcp connect mode {}/{} ", host, port); - socket = factory.scheduledConnect(host, port, messageListener); + socket = factory.scheduledConnect(host, port); return socket; } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/util/MockAgent.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/util/MockAgent.java index d6dd7e286..2703c46bf 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/util/MockAgent.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/util/MockAgent.java @@ -64,21 +64,11 @@ public class MockAgent extends DefaultAgent { return new HoldingSpanStorageFactory(getSpanDataSender()); } - @Override - protected PinpointSocketFactory createPinpointSocketFactory() { - return null; - } - @Override protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { return null; } - @Override - protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory, boolean useMessageListener) { - return null; - } - @Override protected EnhancedDataSender createTcpDataSender(PinpointSocket socket) { return new LoggingDataSender(); diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java index 4fb2d19e1..6f87f4755 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java @@ -8,9 +8,7 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; /** @@ -34,12 +32,6 @@ public class RequestResponseServerMessageListener implements ServerMessageListen channel.sendResponseMessage(requestPacket, requestPacket.getPayload()); } - - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.info("handlerStream {} {}", streamChannel, streamChannel); - } - @Override public int handleEnableWorker(Map properties) { logger.info("handleEnableWorker {}", properties); diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocket.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocket.java index cc9ec2cd0..3c3d2db3b 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocket.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocket.java @@ -10,6 +10,8 @@ import com.nhn.pinpoint.rpc.DefaultFuture; import com.nhn.pinpoint.rpc.Future; import com.nhn.pinpoint.rpc.PinpointSocketException; import com.nhn.pinpoint.rpc.ResponseMessage; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener; import com.nhn.pinpoint.rpc.util.AssertUtils; @@ -21,7 +23,6 @@ import com.nhn.pinpoint.rpc.util.AssertUtils; public class PinpointSocket { private final Logger logger = LoggerFactory.getLogger(this.getClass()); - private final MessageListener messageListener; private volatile SocketHandler socketHandler; @@ -32,28 +33,19 @@ public class PinpointSocket { public PinpointSocket() { this(new ReconnectStateSocketHandler()); } - - public PinpointSocket(MessageListener messageListener) { - this(new ReconnectStateSocketHandler(), messageListener); - } public PinpointSocket(SocketHandler socketHandler) { - this(socketHandler, SimpleLoggingMessageListener.LISTENER); - } - - public PinpointSocket(SocketHandler socketHandler, MessageListener messageListener) { AssertUtils.assertNotNull(socketHandler, "socketHandler"); - AssertUtils.assertNotNull(messageListener, "messageListener"); - - this.messageListener = messageListener; - socketHandler.setMessageListener(this.messageListener); + + if (socketHandler.isSupportServerMode()) { + socketHandler.turnOnServerMode(); + } this.socketHandler = socketHandler; socketHandler.setPinpointSocket(this); } - void reconnectSocketHandler(SocketHandler socketHandler) { AssertUtils.assertNotNull(socketHandler, "socketHandler"); @@ -64,8 +56,11 @@ public class PinpointSocket { } logger.warn("reconnectSocketHandler:{}", socketHandler); - // Pinpoint 소켓 내부 객체가 되기전에 listener를 먼저 등록 - socketHandler.setMessageListener(messageListener); + // Pinpoint 소켓 내부 객체가 되기전에 listener를 먼저 등록 + if (socketHandler.isSupportServerMode()) { + socketHandler.turnOnServerMode(); + } + this.socketHandler = socketHandler; notifyReconnectEvent(); @@ -118,12 +113,11 @@ public class PinpointSocket { return socketHandler.request(bytes); } - - public StreamChannel createStreamChannel() { + public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) { // 실패를 리턴하는 StreamChannel을 던져야 되는데. StreamChannel을 interface로 변경해야 됨. // 일단 그냥 ex를 던지도록 하겠음. ensureOpen(); - return socketHandler.createStreamChannel(); + return socketHandler.createStreamChannel(payload, clientStreamChannelMessageListener); } private Future returnFailureFuture() { diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketFactory.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketFactory.java index f1774ef84..5af2d3bf9 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketFactory.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketFactory.java @@ -29,6 +29,8 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.common.util.PinpointThreadFactory; import com.nhn.pinpoint.rpc.PinpointSocketException; +import com.nhn.pinpoint.rpc.stream.DisabledServerStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.ServerStreamChannelMessageListener; import com.nhn.pinpoint.rpc.util.AssertUtils; import com.nhn.pinpoint.rpc.util.LoggerFactorySetup; import com.nhn.pinpoint.rpc.util.TimerFactory; @@ -61,6 +63,8 @@ public class PinpointSocketFactory { private long enableWorkerPacketDelay = DEFAULT_ENABLE_WORKER_PACKET_DELAY; private long timeoutMillis = DEFAULT_TIMEOUTMILLIS; + private MessageListener messageListener = SimpleLoggingMessageListener.LISTENER; + private ServerStreamChannelMessageListener serverStreamChannelMessageListener = DisabledServerStreamChannelMessageListener.INSTANCE; static { LoggerFactorySetup.setupSlf4jLoggerFactory(); @@ -181,33 +185,21 @@ public class PinpointSocketFactory { } public PinpointSocket connect(String host, int port) throws PinpointSocketException { - return connect(host, port, SimpleLoggingMessageListener.LISTENER); - } - - public PinpointSocket connect(String host, int port, MessageListener messageListener) throws PinpointSocketException { - AssertUtils.assertNotNull(messageListener); - SocketAddress address = new InetSocketAddress(host, port); ChannelFuture connectFuture = bootstrap.connect(address); SocketHandler socketHandler = getSocketHandler(connectFuture, address); - PinpointSocket pinpointSocket = new PinpointSocket(socketHandler, messageListener); + PinpointSocket pinpointSocket = new PinpointSocket(socketHandler); traceSocket(pinpointSocket); return pinpointSocket; } public PinpointSocket reconnect(String host, int port) throws PinpointSocketException { - return reconnect(host, port, SimpleLoggingMessageListener.LISTENER); - } - - public PinpointSocket reconnect(String host, int port, MessageListener messageListener) throws PinpointSocketException { - AssertUtils.assertNotNull(messageListener); - SocketAddress address = new InetSocketAddress(host, port); ChannelFuture connectFuture = bootstrap.connect(address); SocketHandler socketHandler = getSocketHandler(connectFuture, address); - PinpointSocket pinpointSocket = new PinpointSocket(socketHandler, messageListener); + PinpointSocket pinpointSocket = new PinpointSocket(socketHandler); traceSocket(pinpointSocket); return pinpointSocket; } @@ -218,13 +210,7 @@ public class PinpointSocketFactory { } public PinpointSocket scheduledConnect(String host, int port) { - return scheduledConnect(host, port, SimpleLoggingMessageListener.LISTENER); - } - - public PinpointSocket scheduledConnect(String host, int port, MessageListener messageListener) { - AssertUtils.assertNotNull(messageListener); - - PinpointSocket pinpointSocket = new PinpointSocket(new ReconnectStateSocketHandler(), messageListener); + PinpointSocket pinpointSocket = new PinpointSocket(new ReconnectStateSocketHandler()); SocketAddress address = new InetSocketAddress(host, port); reconnect(pinpointSocket, address); return pinpointSocket; @@ -383,4 +369,24 @@ public class PinpointSocketFactory { this.properties = Collections.unmodifiableMap(agentProperties); } + public MessageListener getMessageListener() { + return messageListener; + } + + public void setMessageListener(MessageListener messageListener) { + AssertUtils.assertNotNull(messageListener, "messageListener must not be null"); + + this.messageListener = messageListener; + } + + public ServerStreamChannelMessageListener getServerStreamChannelMessageListener() { + return serverStreamChannelMessageListener; + } + + public void setServerStreamChannelMessageListener(ServerStreamChannelMessageListener serverStreamChannelMessageListener) { + AssertUtils.assertNotNull(messageListener, "messageListener must not be null"); + + this.serverStreamChannelMessageListener = serverStreamChannelMessageListener; + } + } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketHandler.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketHandler.java index da354183e..14d4cb7dd 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketHandler.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/PinpointSocketHandler.java @@ -37,8 +37,13 @@ import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.ResponsePacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; -import com.nhn.pinpoint.rpc.util.AssertUtils; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.DisabledServerStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.ServerStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.StreamChannelManager; import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils; +import com.nhn.pinpoint.rpc.util.IDGenerator; import com.nhn.pinpoint.rpc.util.MapUtils; import com.nhn.pinpoint.rpc.util.TimerFactory; @@ -60,7 +65,6 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke private final State state = new State(); private volatile Channel channel; - private volatile MessageListener messageListener = SimpleLoggingMessageListener.LISTENER; private long timeoutMillis = DEFAULT_TIMEOUTMILLIS; private long pingDelay = DEFAULT_PING_DELAY; @@ -74,8 +78,10 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke private SocketAddress connectSocketAddress; private volatile PinpointSocket pinpointSocket; + private final MessageListener messageListener; + private final ServerStreamChannelMessageListener serverStreamChannelMessageListener; + private final RequestManager requestManager; - private final StreamChannelManager streamChannelManager; private final ChannelFutureListener pingWriteFailFutureListener = new WriteFailFutureListener(this.logger, "ping write fail.", "ping write success."); private final ChannelFutureListener sendWriteFailFutureListener = new WriteFailFutureListener(this.logger, "send() write fail.", "send() write fail."); @@ -95,10 +101,25 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke this.channelTimer = timer; this.pinpointSocketFactory = pinpointSocketFactory; this.requestManager = new RequestManager(timer); - this.streamChannelManager = new StreamChannelManager(); this.pingDelay = pingDelay; this.enableWorkerPacketDelay = enableWorkerPacketDelay; this.timeoutMillis = timeoutMillis; + + MessageListener messageLisener = pinpointSocketFactory.getMessageListener(); + if (messageLisener != null) { + this.messageListener = messageLisener; + } else { + this.messageListener = SimpleLoggingMessageListener.LISTENER; + } + + ServerStreamChannelMessageListener serverStreamChannelMessageListener = pinpointSocketFactory.getServerStreamChannelMessageListener(); + if (serverStreamChannelMessageListener != null) { + this.serverStreamChannelMessageListener = serverStreamChannelMessageListener; + } else { + this.serverStreamChannelMessageListener = DisabledServerStreamChannelMessageListener.INSTANCE; + } + + pinpointSocketFactory.getServerStreamChannelMessageListener(); } public Timer getChannelTimer() { @@ -134,24 +155,25 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke if (!state.changeRun()) { throw new IllegalStateException("invalid open state:" + state.getString()); } + + Channel channel = this.channel; + if (channel != null) { + prepareChannel(channel); + } } - @Override - public void setMessageListener(MessageListener messageListener) { - AssertUtils.assertNotNull(messageListener, "messageListener"); - - logger.info("{} registered Listner({}).", toString(), messageListener); - - if (messageListener != SimpleLoggingMessageListener.LISTENER) { - this.messageListener = messageListener; - - // MessageListener 등록시 EnableWorkerPacket전달 - sendEnableWorkerPacket(); - - RegisterEnableWorkerPacketJob job = new RegisterEnableWorkerPacketJob(enableWorkerPacketRetryCount); - reservationEnableWorkerPacketJob(job); - } - } + private void prepareChannel(Channel channel) { + ServerStreamChannelMessageListener serverStreamChannelMessageListener = this.serverStreamChannelMessageListener; + + StreamChannelManager streamChannelManager = new StreamChannelManager(channel, IDGenerator.createOddIdGenerator(), serverStreamChannelMessageListener); + + SocketHandlerContext context = new SocketHandlerContext(channel, streamChannelManager); + channel.setAttachment(context); + } + + private SocketHandlerContext getChannelContext(Channel channel) { + return (SocketHandlerContext) channel.getAttachment(); + } @Override public void initReconnect() { @@ -373,13 +395,14 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke return messageFuture; } - - - public StreamChannel createStreamChannel() { + + @Override + public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) { ensureOpen(); final Channel channel = this.channel; - return this.streamChannelManager.createStreamChannel(channel); + SocketHandlerContext context = getChannelContext(channel); + return context.getStreamChannelManager().openStreamChannel(payload, clientStreamChannelMessageListener); } @@ -405,7 +428,10 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke case PacketType.APPLICATION_STREAM_CREATE_SUCCESS: case PacketType.APPLICATION_STREAM_CREATE_FAIL: case PacketType.APPLICATION_STREAM_RESPONSE: - this.streamChannelManager.messageReceived((StreamPacket) message, e.getChannel()); + case PacketType.APPLICATION_STREAM_PING: + case PacketType.APPLICATION_STREAM_PONG: + SocketHandlerContext context = getChannelContext(channel); + context.getStreamChannelManager().messageReceived((StreamPacket) message); return; case PacketType.CONTROL_SERVER_CLOSE: messageReceivedServerClosed(e.getChannel()); @@ -550,7 +576,12 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke private void releaseResource() { logger.debug("releaseResource()"); this.requestManager.close(); - this.streamChannelManager.close(); + + if (this.channel != null) { + SocketHandlerContext context = getChannelContext(channel); + context.getStreamChannelManager().close(); + } + this.channelTimer.stop(); } @@ -588,5 +619,37 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke public boolean isConnected() { return this.state.isRun(); } + + @Override + public boolean isSupportServerMode() { + return messageListener != SimpleLoggingMessageListener.LISTENER; + } + + @Override + public void turnOnServerMode() { + // MessageListener 등록시 EnableWorkerPacket전달 + sendEnableWorkerPacket(); + + RegisterEnableWorkerPacketJob job = new RegisterEnableWorkerPacketJob(enableWorkerPacketRetryCount); + reservationEnableWorkerPacketJob(job); + } + class SocketHandlerContext { + private final Channel channel; + private final StreamChannelManager streamChannelManager; + + public SocketHandlerContext(Channel channel, StreamChannelManager streamChannelManager) { + this.channel = channel; + this.streamChannelManager = streamChannelManager; + } + + public Channel getChannel() { + return channel; + } + + public StreamChannelManager getStreamChannelManager() { + return streamChannelManager; + } + } + } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/ReconnectStateSocketHandler.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/ReconnectStateSocketHandler.java index 44ec9a38c..d633b3ecc 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/ReconnectStateSocketHandler.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/ReconnectStateSocketHandler.java @@ -4,6 +4,8 @@ import com.nhn.pinpoint.rpc.DefaultFuture; import com.nhn.pinpoint.rpc.Future; import com.nhn.pinpoint.rpc.PinpointSocketException; import com.nhn.pinpoint.rpc.ResponseMessage; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener; import java.net.SocketAddress; @@ -22,10 +24,6 @@ public class ReconnectStateSocketHandler implements SocketHandler { public void open() { throw new IllegalStateException(); } - - @Override - public void setMessageListener(MessageListener messageListener) { - } @Override public void initReconnect() { @@ -70,10 +68,10 @@ public class ReconnectStateSocketHandler implements SocketHandler { } @Override - public StreamChannel createStreamChannel() { - throw new UnsupportedOperationException(); + public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) { + throw new UnsupportedOperationException(); } - + @Override public void sendPing() { } @@ -82,4 +80,15 @@ public class ReconnectStateSocketHandler implements SocketHandler { public boolean isConnected() { return false; } + + @Override + public boolean isSupportServerMode() { + return false; + } + + @Override + public void turnOnServerMode() { + throw new UnsupportedOperationException(); + } + } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketClientPipelineFactory.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketClientPipelineFactory.java index 43049fd47..6cde8c164 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketClientPipelineFactory.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketClientPipelineFactory.java @@ -31,12 +31,15 @@ public class SocketClientPipelineFactory implements ChannelPipelineFactory { ChannelPipeline pipeline = Channels.pipeline(); pipeline.addLast("encoder", new PacketEncoder()); pipeline.addLast("decoder", new PacketDecoder()); + long pingDelay = pinpointSocketFactory.getPingDelay(); long enableWorkerPacketDelay = pinpointSocketFactory.getEnableWorkerPacketDelay(); long timeoutMillis = pinpointSocketFactory.getTimeoutMillis(); + PinpointSocketHandler pinpointSocketHandler = new PinpointSocketHandler(pinpointSocketFactory, pingDelay, enableWorkerPacketDelay, timeoutMillis); pipeline.addLast("writeTimeout", new WriteTimeoutHandler(pinpointSocketHandler.getChannelTimer(), 3000, TimeUnit.MILLISECONDS)); pipeline.addLast("socketHandler", pinpointSocketHandler); + return pipeline; } } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketHandler.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketHandler.java index 50f01e3f0..f999a29fd 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketHandler.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/SocketHandler.java @@ -1,9 +1,11 @@ package com.nhn.pinpoint.rpc.client; +import java.net.SocketAddress; + import com.nhn.pinpoint.rpc.Future; import com.nhn.pinpoint.rpc.ResponseMessage; - -import java.net.SocketAddress; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener; /** * @author emeroad @@ -29,11 +31,14 @@ public interface SocketHandler { Future request(byte[] bytes); - StreamChannel createStreamChannel(); + ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener); void sendPing(); boolean isConnected(); - void setMessageListener(MessageListener messageListener); + boolean isSupportServerMode(); + + void turnOnServerMode(); + } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannel.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannel.java deleted file mode 100644 index 136913f87..000000000 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannel.java +++ /dev/null @@ -1,225 +0,0 @@ -package com.nhn.pinpoint.rpc.client; - -import java.util.concurrent.atomic.AtomicInteger; - -import org.jboss.netty.channel.Channel; -import org.jboss.netty.channel.ChannelFuture; -import org.jboss.netty.channel.ChannelFutureListener; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.nhn.pinpoint.rpc.DefaultFuture; -import com.nhn.pinpoint.rpc.FailureEventHandler; -import com.nhn.pinpoint.rpc.Future; -import com.nhn.pinpoint.rpc.StreamCreateResponse; -import com.nhn.pinpoint.rpc.packet.PacketType; -import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; -import com.nhn.pinpoint.rpc.stream.StreamChannelMessageListener; - -/** - * @author emeroad - */ -public class StreamChannel { - - private final Logger logger = LoggerFactory.getLogger(this.getClass()); - - private static final int NONE = 0; - // OPEN 호출 - private static final int OPEN = 1; - // OPEN 결과 대기 - private static final int OPEN_AWAIT = 2; - // 동작중 - private static final int RUN = 3; - // 닫힘 - private static final int CLOSED = 4; - - private final AtomicInteger state = new AtomicInteger(NONE); - - private final int channelId; - - private StreamChannelManager streamChannelManager; - - private StreamChannelMessageListener streamChannelMessageListener; - - private DefaultFuture openLatch; - private Channel channel; - - public StreamChannel(int channelId) { - this.channelId = channelId; - } - - public int getChannelId() { - return channelId; - } - - public void setChannel(Channel channel) { - this.channel = channel; - } - - public Future open(byte[] bytes) { - if (!state.compareAndSet(NONE, OPEN)) { - throw new IllegalStateException("invalid state:" + state.get()); - } - StreamCreatePacket streamCreatePacket = new StreamCreatePacket(channelId, bytes); - - this.openLatch = new DefaultFuture(); - openLatch.setFailureEventHandler(new FailureEventHandler() { - @Override - public boolean fireFailure() { - streamChannelManager.closeChannel(channelId); - return false; - } - }); - ChannelFuture channelFuture = this.channel.write(streamCreatePacket); - channelFuture.addListener(new ChannelFutureListener() { - @Override - public void operationComplete(ChannelFuture future) throws Exception { - if (!future.isSuccess()) { - future.setFailure(future.getCause()); - } - } - }); - - - if (!state.compareAndSet(OPEN, OPEN_AWAIT)) { - throw new IllegalStateException("invalid state"); - } - return openLatch; - } - - - public boolean receiveStreamPacket(StreamPacket packet) { - final short packetType = packet.getPacketType(); - switch (packetType) { - case PacketType.APPLICATION_STREAM_CREATE_SUCCESS: - logger.debug("APPLICATION_STREAM_CREATE_SUCCESS {}", channel); - StreamCreateResponse success = new StreamCreateResponse(true); - success.setMessage(packet.getPayload()); - return openChannel(RUN, success); - - case PacketType.APPLICATION_STREAM_CREATE_FAIL: - logger.debug("APPLICATION_STREAM_CREATE_FAIL {}", channel); - StreamCreateResponse failResult = new StreamCreateResponse(false); - failResult.setMessage(packet.getPayload()); - return openChannel(CLOSED, failResult); - - case PacketType.APPLICATION_STREAM_RESPONSE: { - logger.debug("APPLICATION_STREAM_RESPONSE {}", channel); - - StreamResponsePacket streamResponsePacket = (StreamResponsePacket) packet; - StreamChannelMessageListener streamChannelMessageListener = this.streamChannelMessageListener; - if (streamChannelMessageListener != null) { - streamChannelMessageListener.handleStreamData(this, streamResponsePacket); - } - return true; - } - case PacketType.APPLICATION_STREAM_CLOSE: { - logger.debug("APPLICATION_STREAM_CLOSE {}", channel); - - this.closeInternal(); - - StreamClosePacket streamClosePacket = (StreamClosePacket) packet; - StreamChannelMessageListener streamChannelMessageListener = this.streamChannelMessageListener; - if (streamChannelMessageListener != null) { - streamChannelMessageListener.handleStreamClose(this, streamClosePacket); - } - - return true; - } - } - return false; - } - - private boolean openChannel(int channelState, StreamCreateResponse streamCreateResponse) { - if (state.compareAndSet(OPEN_AWAIT, channelState)) { - notifyOpenResult(streamCreateResponse); - return true; - } else { - logger.info("invalid stream channel state:{}", state.get()); - return false; - } - } - - - - - private boolean notifyOpenResult(StreamCreateResponse failResult) { - DefaultFuture openLatch = this.openLatch; - if (openLatch != null) { - return openLatch.setResult(failResult); - } - return false; - } - - - - public boolean close() { - return close0(true); - } - - - - boolean closeInternal() { - return close0(false); - } - - private boolean close0(boolean safeClose) { - if (!state.compareAndSet(RUN, CLOSED)) { - return false; - } - - if (safeClose) { - StreamClosePacket closePacket = new StreamClosePacket(this.channelId, StreamClosePacket.SUCCESS); - this.channel.write(closePacket); - - StreamChannelManager streamChannelManager = this.streamChannelManager; - if (streamChannelManager != null) { - streamChannelManager.closeChannel(channelId); - this.streamChannelManager = null; - } - } - return true; - } - - public void setStreamChannelManager(StreamChannelManager streamChannelManager) { - this.streamChannelManager = streamChannelManager; - } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - - StreamChannel that = (StreamChannel) o; - - if (channelId != that.channelId) return false; - if (channel != null ? !channel.equals(that.channel) : that.channel != null) return false; - - return true; - } - - @Override - public int hashCode() { - int result = channelId; - result = 31 * result + (channel != null ? channel.hashCode() : 0); - return result; - } - - public void setStreamChannelMessageListener(StreamChannelMessageListener streamChannelMessageListener) { - this.streamChannelMessageListener = streamChannelMessageListener; - } - - @Override - public String toString() { - final StringBuilder sb = new StringBuilder(); - sb.append("StreamChannel"); - sb.append("{channelId=").append(channelId); - sb.append(", channel=").append(channel); - sb.append('}'); - return sb.toString(); - } -} - diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannelManager.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannelManager.java deleted file mode 100644 index a940fd512..000000000 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/StreamChannelManager.java +++ /dev/null @@ -1,81 +0,0 @@ -package com.nhn.pinpoint.rpc.client; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; -import java.util.concurrent.atomic.AtomicInteger; - -import org.jboss.netty.channel.Channel; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.nhn.pinpoint.rpc.PinpointSocketException; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; - -/** - * @author emeroad - */ -public class StreamChannelManager { - - private Logger logger = LoggerFactory.getLogger(this.getClass()); - - private final AtomicInteger idAllocator = new AtomicInteger(0); - - private final ConcurrentMap channelMap = new ConcurrentHashMap(); - - public StreamChannel createStreamChannel(Channel channel) { - final int channelId = allocateChannelId(); - StreamChannel streamChannel = new StreamChannel(channelId); - streamChannel.setChannel(channel); - - StreamChannel old = channelMap.put(channelId, streamChannel); - if (old != null) { - throw new PinpointSocketException("already channelId exist:" + channelId + " streamChannel:" + old); - } - // handle을 붙여서 리턴. - streamChannel.setStreamChannelManager(this); - - return streamChannel; - } - - private int allocateChannelId() { - return idAllocator.get(); - } - - - public StreamChannel findStreamChannel(int channelId) { - return this.channelMap.get(channelId); - } - - public boolean closeChannel(int channelId) { - StreamChannel remove = this.channelMap.remove(channelId); - return remove != null; - } - - public void close() { - logger.debug("close()"); - final ConcurrentMap channelMap = this.channelMap; - - int forceCloseChannel = 0; - for (Map.Entry entry : channelMap.entrySet()) { - if(entry.getValue().closeInternal()) { - forceCloseChannel++; - } - } - channelMap.clear(); - if(forceCloseChannel > 0) { - logger.info("streamChannelManager forceCloseChannel {}", forceCloseChannel); - } - } - - - public boolean messageReceived(StreamPacket streamPacket, Channel channel) { - final int channelId = streamPacket.getStreamChannelId(); - final StreamChannel streamChannel = findStreamChannel(channelId); - if (streamChannel == null) { - logger.warn("streamChannel not found. channelId:{} ", channelId, channel); - return false; - } - return streamChannel.receiveStreamPacket(streamPacket); - } -} 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 1a3b85188..b88efbddf 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 @@ -6,11 +6,16 @@ import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.StreamChannelContext; +import com.nhn.pinpoint.rpc.stream.StreamChannelManager; + public class ChannelContext { private final Logger logger = LoggerFactory.getLogger(this.getClass()); - private final ServerStreamChannelManager streamChannelManager; + private final StreamChannelManager streamChannelManager; private final SocketChannel socketChannel; @@ -20,11 +25,11 @@ public class ChannelContext { private volatile Map channelProperties = Collections.emptyMap(); - public ChannelContext(SocketChannel socketChannel, ServerStreamChannelManager streamChannelManager) { + public ChannelContext(SocketChannel socketChannel, StreamChannelManager streamChannelManager) { this(socketChannel, streamChannelManager, DoNothingChannelStateEventListener.INSTANCE); } - public ChannelContext(SocketChannel socketChannel, ServerStreamChannelManager streamChannelManager, SocketChannelStateChangeEventListener stateChangeEventListener) { + public ChannelContext(SocketChannel socketChannel, StreamChannelManager streamChannelManager, SocketChannelStateChangeEventListener stateChangeEventListener) { this.socketChannel = socketChannel; this.streamChannelManager = streamChannelManager; @@ -33,16 +38,16 @@ public class ChannelContext { this.state = new PinpointServerSocketState(); } - public ServerStreamChannel getStreamChannel(int channelId) { + public StreamChannelContext getStreamChannel(int channelId) { return streamChannelManager.findStreamChannel(channelId); } - public ServerStreamChannel createStreamChannel(int channelId) { - return streamChannelManager.createStreamChannel(channelId); + public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) { + return streamChannelManager.openStreamChannel(payload, clientStreamChannelMessageListener); } public void closeAllStreamChannel() { - streamChannelManager.closeInternal(); + streamChannelManager.close(); } public SocketChannel getSocketChannel() { @@ -112,5 +117,9 @@ public class ChannelContext { this.channelProperties = Collections.unmodifiableMap(properties); return true; } + + public StreamChannelManager getStreamChannelManager() { + return streamChannelManager; + } } 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 4a6622975..458317038 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 @@ -47,11 +47,14 @@ import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.ResponsePacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.ServerClosePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket; import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; +import com.nhn.pinpoint.rpc.stream.DisabledServerStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.ServerStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.StreamChannelManager; +import com.nhn.pinpoint.rpc.util.AssertUtils; import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils; import com.nhn.pinpoint.rpc.util.CpuUtils; +import com.nhn.pinpoint.rpc.util.IDGenerator; import com.nhn.pinpoint.rpc.util.LoggerFactorySetup; import com.nhn.pinpoint.rpc.util.TimerFactory; @@ -76,6 +79,8 @@ public class PinpointServerSocket extends SimpleChannelHandler { private final Timer requestManagerTimer; private ServerMessageListener messageListener = SimpleLoggingServerMessageListener.LISTENER; + private ServerStreamChannelMessageListener serverStreamChannelMessageListener = DisabledServerStreamChannelMessageListener.INSTANCE; + private WriteFailFutureListener traceSendAckWriteFailFutureListener = new WriteFailFutureListener(logger, "TraceSendAckPacket send fail.", "TraceSendAckPacket send() success."); private InetAddress[] ignoreAddressList; @@ -126,6 +131,12 @@ public class PinpointServerSocket extends SimpleChannelHandler { } this.messageListener = messageListener; } + + public void setServerStreamChannelMessageListener(ServerStreamChannelMessageListener serverStreamChannelMessageListener) { + AssertUtils.assertNotNull(serverStreamChannelMessageListener, "serverStreamChannelMessageListener must not be null"); + + this.serverStreamChannelMessageListener = serverStreamChannelMessageListener; + } private void setOptions(ServerBootstrap bootstrap) { // read write timeout이 있어야 되나? nio라서 없어도 되던가? @@ -214,6 +225,8 @@ public class PinpointServerSocket extends SimpleChannelHandler { case PacketType.APPLICATION_STREAM_CREATE_SUCCESS: case PacketType.APPLICATION_STREAM_CREATE_FAIL: case PacketType.APPLICATION_STREAM_RESPONSE: + case PacketType.APPLICATION_STREAM_PING: + case PacketType.APPLICATION_STREAM_PONG: handleStreamPacket((StreamPacket) message, channel); return; case PacketType.CONTROL_ENABLE_WORKER: @@ -258,31 +271,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { private void handleStreamPacket(StreamPacket packet, Channel channel) { ChannelContext context = getChannelContext(channel); - if (packet instanceof StreamCreatePacket) { - logger.debug("StreamCreate {}, streamId:{}", channel, packet.getStreamChannelId()); - try { - ServerStreamChannel streamChannel = context.createStreamChannel(packet.getStreamChannelId()); - boolean success = streamChannel.receiveChannelCreate((StreamCreatePacket) packet); - if (success) { - messageListener.handleStream(packet, streamChannel); - } - - } catch (PinpointSocketException e) { - logger.warn("channel create fail. channel:{} Caused:{}", channel, e); - } - } else if (packet instanceof StreamClosePacket) { - logger.debug("StreamDestroy {}, streamId:{}", channel, packet.getStreamChannelId()); - ServerStreamChannel streamChannel = context.getStreamChannel(packet.getStreamChannelId()); - // null이 나올수 있음. - boolean close = streamChannel.close(); - if (close) { - messageListener.handleStream(packet, streamChannel); - } else { - logger.warn("invalid streamClosePacket. already close. channel:{} Caused:{}", channel); - } - } else { - logger.warn("invalid streamPacket. channel:{}", channel); - } + context.getStreamChannelManager().messageReceived(packet); } private Map decodeSocketProperties(ControlEnableWorkerPacket message) { @@ -439,7 +428,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { private void prepareChannel(Channel channel) { SocketChannel socketChannel = new SocketChannel(channel, DEFAULT_TIMEOUTMILLIS, requestManagerTimer); - ServerStreamChannelManager streamChannelManager = new ServerStreamChannelManager(channel); + StreamChannelManager streamChannelManager = new StreamChannelManager(channel, IDGenerator.createEvenIdGenerator(), serverStreamChannelMessageListener); ChannelContext channelContext = new ChannelContext(socketChannel, streamChannelManager, channelStateChangeEventListener); diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java index 27d8ca3de..3ad87234e 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java @@ -4,7 +4,6 @@ import java.util.Map; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; /** * @author emeroad @@ -15,8 +14,6 @@ public interface ServerMessageListener { // 외부 노출 Channel은 별도의 Tcp Channel로 감싸는걸로 변경할 것. void handleRequest(RequestPacket requestPacket, SocketChannel channel); - void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel); - int handleEnableWorker(Map properties); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannel.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannel.java deleted file mode 100644 index 16112b348..000000000 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannel.java +++ /dev/null @@ -1,159 +0,0 @@ -package com.nhn.pinpoint.rpc.server; - -import java.util.concurrent.atomic.AtomicInteger; - -import org.jboss.netty.channel.Channel; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreateFailPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreateSuccessPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket; - -/** - * @author emeroad - */ -public class ServerStreamChannel { - - private final Logger logger = LoggerFactory.getLogger(this.getClass()); - - private static final int NONE = 0; - // OPEN이 도착함. - private static final int OPEN_ARRIVED = 1; - // create success 던짐. 동작중 - private static final int RUN = 2; - // 닫힘 - private static final int CLOSED = 2; - - private final AtomicInteger state = new AtomicInteger(NONE); - - private final int channelId; - - private ServerStreamChannelManager serverStreamChannelManager; - - private Channel channel; - - public ServerStreamChannel(int channelId) { - this.channelId = channelId; - } - - public int getChannelId() { - return channelId; - } - - public void setChannel(Channel channel) { - this.channel = channel; - } - - -// public boolean receiveStreamPacket(StreamPacket packet) { -// final short packetType = packet.getPacketType(); -// switch (packetType) { -// case PacketType.APPLICATION_STREAM_CREATE: -// logger.info("APPLICATION_STREAM_CREATE_SUCCESS"); -// return receiveChannelCreate((StreamCreatePacket) packet); -// } -// return false; -// } - - public boolean receiveChannelCreate(StreamCreatePacket streamCreateResponse) { - if (state.compareAndSet(NONE, OPEN_ARRIVED)) { - return true; - } else { - logger.info("invalid state:{}", state.get()); - return false; - } - } - - public boolean sendOpenResult(boolean success, byte[] bytes) { - if(success ) { - if(!state.compareAndSet(OPEN_ARRIVED, RUN)) { - return false; - } - StreamCreateSuccessPacket streamCreateSuccessPacket = new StreamCreateSuccessPacket(channelId); - this.channel.write(streamCreateSuccessPacket); - return true; - } else { - if(!state.compareAndSet(OPEN_ARRIVED, CLOSED)) { - return false; - } - StreamCreateFailPacket streamCreateFailPacket = new StreamCreateFailPacket(channelId, StreamCreateFailPacket.UNKNWON_ERROR); - this.channel.write(streamCreateFailPacket); - return true; - } - } - - public boolean sendStreamMessage(byte[] bytes) { - if (state.get() != RUN) { - return false; - } - StreamResponsePacket response = new StreamResponsePacket(channelId, bytes); - this.channel.write(response); - return true; - } - - - public boolean close() { - return close0(true); - } - - boolean closeInternal() { - return close0(false); - } - - private boolean close0(boolean safeClose) { - if (!state.compareAndSet(RUN, CLOSED)) { - return false; - } - - if (safeClose) { - StreamClosePacket streamClosePacket = new StreamClosePacket(channelId, StreamClosePacket.SUCCESS); - this.channel.write(streamClosePacket); - - ServerStreamChannelManager serverStreamChannelManager = this.serverStreamChannelManager; - if (serverStreamChannelManager != null) { - serverStreamChannelManager.closeChannel(channelId); - this.serverStreamChannelManager = null; - } - } - return true; - } - - public void setServerStreamChannelManager(ServerStreamChannelManager serverStreamChannelManager) { - this.serverStreamChannelManager = serverStreamChannelManager; - } - - @Override - public boolean equals(Object o) { - if (this == o) return true; - if (o == null || getClass() != o.getClass()) return false; - - ServerStreamChannel that = (ServerStreamChannel) o; - - if (channelId != that.channelId) return false; - if (channel != null ? !channel.equals(that.channel) : that.channel != null) return false; - - return true; - } - - @Override - public int hashCode() { - int result = channelId; - result = 31 * result + (channel != null ? channel.hashCode() : 0); - return result; - } - - @Override - public String toString() { - final StringBuilder sb = new StringBuilder(); - sb.append("ServerStreamChannel"); - sb.append("{state=").append(state); - sb.append(", channelId=").append(channelId); - sb.append(", channel=").append(channel); - sb.append('}'); - return sb.toString(); - } -} - diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannelManager.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannelManager.java deleted file mode 100644 index cef0ef09f..000000000 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerStreamChannelManager.java +++ /dev/null @@ -1,65 +0,0 @@ -package com.nhn.pinpoint.rpc.server; - -import com.nhn.pinpoint.rpc.PinpointSocketException; -import org.jboss.netty.channel.Channel; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - -import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.ConcurrentMap; - -/** - * @author emeroad - */ -public class ServerStreamChannelManager { - - private final Logger logger = LoggerFactory.getLogger(this.getClass()); - - private final Channel channel; - private final ConcurrentMap channelMap = new ConcurrentHashMap(); - - public ServerStreamChannelManager(Channel channel) { - if (channel == null) { - throw new NullPointerException("channel"); - } - this.channel = channel; - } - - public ServerStreamChannel createStreamChannel(int channelId) { - ServerStreamChannel streamChannel = new ServerStreamChannel(channelId); - streamChannel.setChannel(channel); - - ServerStreamChannel old = channelMap.put(channelId, streamChannel); - if (old != null) { - throw new PinpointSocketException("already channelId exist:" + channelId + " streamChannel:" + old); - } - // handle을 붙여서 리턴. - streamChannel.setServerStreamChannelManager(this); - - return streamChannel; - } - - - public ServerStreamChannel findStreamChannel(int channelId) { - return this.channelMap.get(channelId); - } - - public boolean closeChannel(int channelId) { - ServerStreamChannel remove = this.channelMap.remove(channelId); - return remove != null; - } - - - public void closeInternal() { - final boolean debugEnabled = logger.isDebugEnabled(); - for (Map.Entry streamChannel : this.channelMap.entrySet()) { - streamChannel.getValue().closeInternal(); - if (debugEnabled) { - logger.debug("ServerStreamChannel.closeInternal() id:{}, {}", streamChannel.getKey(), channel); - } - } - - - } -} diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java index b4200de83..ea23e87e0 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java @@ -8,7 +8,6 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; /** * @author emeroad @@ -29,12 +28,6 @@ public class SimpleLoggingServerMessageListener implements ServerMessageListener logger.info("handlerRequest {} {}", requestPacket, channel); } - - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.info("handlerStream {} {}", streamChannel, streamChannel); - } - @Override public int handleEnableWorker(Map properties) { logger.info("handleEnableWorker {}", properties); diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannelMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannelMessageListener.java deleted file mode 100644 index 6ab5ca757..000000000 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/stream/StreamChannelMessageListener.java +++ /dev/null @@ -1,19 +0,0 @@ -package com.nhn.pinpoint.rpc.stream; - -import com.nhn.pinpoint.rpc.client.StreamChannel; -import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket; - -/** - * @author koo.taejin - */ -public interface StreamChannelMessageListener { - - short handleStreamCreate(StreamChannel streamChannel, StreamCreatePacket packet); - - void handleStreamData(StreamChannel streamChannel, StreamResponsePacket packet); - - void handleStreamClose(StreamChannel streamChannel, StreamClosePacket packet); - -} 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 93c45c498..32a2732a4 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/RecordedStreamChannelMessageListener.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/RecordedStreamChannelMessageListener.java @@ -8,17 +8,17 @@ import java.util.concurrent.CountDownLatch; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.nhn.pinpoint.rpc.client.StreamChannel; import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket; import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket; -import com.nhn.pinpoint.rpc.stream.StreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext; +import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.stream.StreamChannelContext; /** * @author emeroad * @author koo.taejin */ -public class RecordedStreamChannelMessageListener implements StreamChannelMessageListener { +public class RecordedStreamChannelMessageListener implements ClientStreamChannelMessageListener { private final Logger logger = LoggerFactory.getLogger(this.getClass()); @@ -29,24 +29,17 @@ public class RecordedStreamChannelMessageListener implements StreamChannelMessag public RecordedStreamChannelMessageListener(int receiveMessageCount) { this.latch = new CountDownLatch(receiveMessageCount); } - + @Override - public short handleStreamCreate(StreamChannel streamChannel, StreamCreatePacket packet) { - // TODO Auto-generated method stub - return 0; - } - - @Override - public void handleStreamData(StreamChannel streamChannel, StreamResponsePacket packet) { - logger.info("handleStreamData {}, {}", streamChannel, packet); + public void handleStreamData(ClientStreamChannelContext streamChannelContext, StreamResponsePacket packet) { + logger.info("handleStreamData {}, {}", streamChannelContext, packet); receivedMessageList.add(packet.getPayload()); latch.countDown(); } - @Override - public void handleStreamClose(StreamChannel streamChannel, StreamClosePacket packet) { - logger.info("handleClose {}, {}", streamChannel, packet); + public void handleStreamClose(StreamChannelContext 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/client/PinpointSocketFactoryTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/client/PinpointSocketFactoryTest.java index 07e7b0f70..c32eb381b 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/client/PinpointSocketFactoryTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/client/PinpointSocketFactoryTest.java @@ -180,55 +180,6 @@ public class PinpointSocketFactoryTest { } - - - - @Test - public void stream() throws IOException, InterruptedException { - PinpointServerSocket ss = new PinpointServerSocket(); - - TestSeverMessageListener testSeverMessageListener = new TestSeverMessageListener(); - ss.setMessageListener(testSeverMessageListener); - ss.bind("localhost", 10234); - PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); - try { - PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234); - - - StreamChannel streamChannel = socket.createStreamChannel(); - byte[] openBytes = TestByteUtils.createRandomByte(30); - - // 현재 서버에서 3번 보내게 되어 있음. - RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4); - streamChannel.setStreamChannelMessageListener(clientListener); - - Future open = streamChannel.open(openBytes); - open.await(); - StreamCreateResponse response = open.getResult(); - Assert.assertTrue(response.isSuccess()); - // stream 메시지를 대기함. - clientListener.getLatch().await(); - List receivedMessage = clientListener.getReceivedMessage(); - List sendMessage = testSeverMessageListener.getSendMessage(); - - // 한개는 close 패킷임. - Assert.assertEquals(receivedMessage.size(), sendMessage.size()); - for(int i =0; i channelContextList = ss.getDuplexCommunicationChannelContext(); @@ -99,17 +105,18 @@ public class MessageListenerTest { PinpointServerSocket ss = new PinpointServerSocket(); ss.bind("127.0.0.1", 10234); - PinpointSocketFactory socketFactory = createPinpointSocketFactory(); - + PinpointSocketFactory socketFactory1 = createPinpointSocketFactory(); + EchoMessageListener echoMessageListener1 = new EchoMessageListener(); + socketFactory1.setMessageListener(echoMessageListener1); + + PinpointSocketFactory socketFactory2 = createPinpointSocketFactory(); + EchoMessageListener echoMessageListener2 = new EchoMessageListener(); + socketFactory2.setMessageListener(echoMessageListener2); + try { - - EchoMessageListener echoMessageListener1 = new EchoMessageListener(); - EchoMessageListener echoMessageListener2 = new EchoMessageListener(); - - // 리스터를 등록한 것만 RegisterAgent 로 나옴 - PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, echoMessageListener1); - PinpointSocket socket2 = socketFactory.connect("127.0.0.1", 10234, echoMessageListener2); + PinpointSocket socket = socketFactory1.connect("127.0.0.1", 10234); + PinpointSocket socket2 = socketFactory2.connect("127.0.0.1", 10234); Thread.sleep(500); @@ -131,7 +138,9 @@ public class MessageListenerTest { socket.close(); socket2.close(); } finally { - socketFactory.release(); + socketFactory1.release(); + socketFactory2.release(); + ss.close(); } } @@ -143,12 +152,12 @@ public class MessageListenerTest { Map params = getParams(); PinpointSocketFactory socketFactory = createPinpointSocketFactory(params); + socketFactory.setMessageListener(new EchoMessageListener()); try { - EchoMessageListener echoMessageListener1 = new EchoMessageListener(); // 리스터를 등록한 것만 RegisterAgent 로 나옴 - PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, echoMessageListener1); + PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234); Thread.sleep(500); @@ -171,10 +180,11 @@ public class MessageListenerTest { ss.bind("127.0.0.1", 10234); PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + socketFactory.setMessageListener(SimpleLoggingMessageListener.LISTENER); try { // Listener가 없을때 디폴트로 등록하는 SimpleLoggingMessageListener.LISTENER인 경우 상호 연결이 불가능함 - PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, SimpleLoggingMessageListener.LISTENER); + PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234); Thread.sleep(500); @@ -201,11 +211,12 @@ public class MessageListenerTest { PinpointSocketFactory socketFactory = createPinpointSocketFactory(); socketFactory.setEnableWorkerPacketDelay(500); + socketFactory.setMessageListener(new EchoMessageListener()); try { // Listener가 없을때 디폴트로 등록하는 SimpleLoggingMessageListener.LISTENER인 경우 상호 연결이 불가능함 - PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, new EchoMessageListener()); + PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234); Thread.sleep(5000); List channelContextList = ss.getDuplexCommunicationChannelContext(); diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java index b4856ec31..cc9c5e905 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java @@ -7,13 +7,9 @@ import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.nhn.pinpoint.rpc.TestByteUtils; import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket; -import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; /** * @author emeroad @@ -37,42 +33,12 @@ public class TestSeverMessageListener implements ServerMessageListener { channel.sendResponseMessage(requestPacket, requestPacket.getPayload()); } - - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.debug("streamPacket:{} channel:{}", streamPacket, streamChannel); - if (streamPacket instanceof StreamCreatePacket) { - byte[] payload = streamPacket.getPayload(); - this.open = payload; - streamChannel.sendOpenResult(true, payload); - sendStreamMessage(streamChannel); - sendStreamMessage(streamChannel); - sendStreamMessage(streamChannel); - - sendClose(streamChannel); - } else if(streamPacket instanceof StreamClosePacket) { - // 채널 종료해야 함. - } - - } - @Override public int handleEnableWorker(Map properties) { logger.debug("handleEnableWorker properties:{} channel:{}", properties); return ControlEnableWorkerConfirmPacket.SUCCESS; } - private void sendClose(ServerStreamChannel streamChannel) { - sendMessageList.add(new byte[0]); - streamChannel.close(); - } - - private void sendStreamMessage(ServerStreamChannel streamChannel) { - byte[] randomByte = TestByteUtils.createRandomByte(10); - streamChannel.sendStreamMessage(randomByte); - sendMessageList.add(randomByte); - } - public byte[] getOpen() { return open; } 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 new file mode 100644 index 000000000..d1448c140 --- /dev/null +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java @@ -0,0 +1,244 @@ +package com.nhn.pinpoint.rpc.stream; + +import java.io.IOException; +import java.util.List; +import java.util.concurrent.CopyOnWriteArrayList; + +import junit.framework.Assert; + +import org.junit.Test; + +import com.nhn.pinpoint.rpc.PinpointSocketException; +import com.nhn.pinpoint.rpc.RecordedStreamChannelMessageListener; +import com.nhn.pinpoint.rpc.TestByteUtils; +import com.nhn.pinpoint.rpc.client.MessageListener; +import com.nhn.pinpoint.rpc.client.PinpointSocket; +import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; +import com.nhn.pinpoint.rpc.client.SimpleLoggingMessageListener; +import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket; +import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket; +import com.nhn.pinpoint.rpc.server.ChannelContext; +import com.nhn.pinpoint.rpc.server.PinpointServerSocket; +import com.nhn.pinpoint.rpc.server.ServerMessageListener; +import com.nhn.pinpoint.rpc.server.TestSeverMessageListener; + +public class StreamChannelManagerTest { + + // Client to Server Stream + @Test + public void stream1() throws IOException, InterruptedException { + SimpleStreamBO bo = new SimpleStreamBO(); + + PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), new ServerListener(bo)); + ss.bind("localhost", 10234); + + PinpointSocketFactory pinpointSocketFactory = createSocketFactory(); + try { + PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234); + + RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4); + + ClientStreamChannelContext clientContext = socket.createStreamChannel(new byte[0], clientListener); + + int sendCount = 4; + + for (int i = 0; i < 4; i++) { + sendRandomBytes(bo); + } + + Thread.sleep(100); + + Assert.assertEquals(sendCount, clientListener.getReceivedMessage().size()); + + clientContext.getStreamChannel().close(); + socket.close(); + } finally { + pinpointSocketFactory.release(); + ss.close(); + } + } + + @Test(expected = PinpointSocketException.class) + public void stream2() throws IOException, InterruptedException { + PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), null); + ss.bind("localhost", 10234); + + PinpointSocketFactory pinpointSocketFactory = createSocketFactory(); + try { + PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234); + + RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4); + + ClientStreamChannelContext clientContext = socket.createStreamChannel(new byte[0], clientListener); + + Thread.sleep(100); + + clientContext.getStreamChannel().close(); + socket.close(); + } finally { + pinpointSocketFactory.release(); + ss.close(); + } + } + + @Test(expected = PinpointSocketException.class) + public void stream3() throws IOException, InterruptedException { + SimpleStreamBO bo = new SimpleStreamBO(); + + PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), new ServerListener(bo)); + ss.bind("localhost", 10234); + + PinpointSocketFactory pinpointSocketFactory = createSocketFactory(); + + PinpointSocket socket = null; + try { + socket = pinpointSocketFactory.connect("127.0.0.1", 10234); + + RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4); + + ClientStreamChannelContext clientContext = socket.createStreamChannel(new byte[0], clientListener); + Thread.sleep(100); + + clientContext.getStreamChannel().close(); + + Thread.sleep(100); + + sendRandomBytes(bo); + + } finally { + if (socket != null) { + socket.close(); + } + + pinpointSocketFactory.release(); + ss.close(); + } + } + + // ServerSocket to Client Stream + @Test + public void stream4() throws IOException, InterruptedException { + PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), null); + ss.bind("localhost", 10234); + + SimpleStreamBO bo = new SimpleStreamBO(); + + PinpointSocketFactory pinpointSocketFactory = createSocketFactory(new TestListener(), new ServerListener(bo)); + + try { + PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234); + + Thread.sleep(100); + + List contextList = ss.getDuplexCommunicationChannelContext(); + Assert.assertEquals(1, contextList.size()); + + ChannelContext context = contextList.get(0); + + RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4); + + ClientStreamChannelContext clientContext = context.createStreamChannel(new byte[0], clientListener); + + int sendCount = 4; + + for (int i = 0; i < 4; i++) { + sendRandomBytes(bo); + } + + Thread.sleep(100); + + Assert.assertEquals(sendCount, clientListener.getReceivedMessage().size()); + + clientContext.getStreamChannel().close(); + socket.close(); + } finally { + pinpointSocketFactory.release(); + ss.close(); + } + } + + private PinpointServerSocket createServerSocket(ServerMessageListener severMessageListener, + ServerStreamChannelMessageListener serverStreamChannelMessageListener) { + PinpointServerSocket serverSocket = new PinpointServerSocket(); + + if (severMessageListener != null) { + serverSocket.setMessageListener(severMessageListener); + } + + if (serverStreamChannelMessageListener != null) { + serverSocket.setServerStreamChannelMessageListener(serverStreamChannelMessageListener); + } + + return serverSocket; + } + + private PinpointSocketFactory createSocketFactory() { + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + return pinpointSocketFactory; + } + + private PinpointSocketFactory createSocketFactory(MessageListener messageListener, ServerStreamChannelMessageListener serverStreamChannelMessageListener) { + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + pinpointSocketFactory.setMessageListener(messageListener); + pinpointSocketFactory.setServerStreamChannelMessageListener(serverStreamChannelMessageListener); + + return pinpointSocketFactory; + } + + class TestListener extends SimpleLoggingMessageListener { + + } + + private void sendRandomBytes(SimpleStreamBO bo) { + byte[] openBytes = TestByteUtils.createRandomByte(30); + bo.sendResponse(openBytes); + } + + class ServerListener implements ServerStreamChannelMessageListener { + + private final SimpleStreamBO bo; + + public ServerListener(SimpleStreamBO bo) { + this.bo = bo; + } + + @Override + public short handleStreamCreate(ServerStreamChannelContext streamChannelContext, StreamCreatePacket packet) { + bo.addServerStreamChannelContext(streamChannelContext); + return 0; + } + + @Override + public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) { + + } + + } + + class SimpleStreamBO { + + private final List serverStreamChannelContextList; + + public SimpleStreamBO() { + serverStreamChannelContextList = new CopyOnWriteArrayList(); + } + + public void addServerStreamChannelContext(ServerStreamChannelContext context) { + serverStreamChannelContextList.add(context); + } + + public void removeServerStreamChannelContext(ServerStreamChannelContext context) { + serverStreamChannelContextList.remove(context); + } + + void sendResponse(byte[] data) { + + for (ServerStreamChannelContext context : serverStreamChannelContextList) { + context.getStreamChannel().sendData(data); + } + + } + + } + +} diff --git a/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java b/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java index 0a72019fb..8c1aad263 100644 --- a/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java +++ b/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java @@ -22,8 +22,8 @@ import com.nhn.pinpoint.rpc.packet.stream.StreamPacket; import com.nhn.pinpoint.rpc.server.ChannelContext; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; -import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; +import com.nhn.pinpoint.rpc.stream.ServerStreamChannel; import com.nhn.pinpoint.web.cluster.ClusterManager; import com.nhn.pinpoint.web.cluster.zookeeper.ZookeeperClusterManager; import com.nhn.pinpoint.web.config.WebConfig; @@ -168,11 +168,6 @@ public class PinpointSocketManager { logger.warn("Unsupport request received {} {}", requestPacket, channel); } - @Override - public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.warn("unsupported streamPacket received {}", streamPacket); - } - @Override public int handleEnableWorker(Map properties) { logger.warn("do handleEnableWorker {}", properties); diff --git a/web/src/test/java/com/navercorp/pinpoint/web/cluster/ClusterTest.java b/web/src/test/java/com/navercorp/pinpoint/web/cluster/ClusterTest.java index 20a62efd5..44cd47609 100644 --- a/web/src/test/java/com/navercorp/pinpoint/web/cluster/ClusterTest.java +++ b/web/src/test/java/com/navercorp/pinpoint/web/cluster/ClusterTest.java @@ -123,7 +123,9 @@ public class ClusterTest { Assert.assertEquals(0, socketManager.getCollectorChannelContext().size()); factory = new PinpointSocketFactory(); - socket = factory.connect(DEFAULT_IP, DEFAULT_ACCEPTOR_PORT, new SimpleListener()); + factory.setMessageListener(new SimpleListener()); + + socket = factory.connect(DEFAULT_IP, DEFAULT_ACCEPTOR_PORT); Thread.sleep(1000);