From c1e8dd4eae9dd893a03c452a91c16ba20745086b Mon Sep 17 00:00:00 2001 From: koo-taejin Date: Tue, 25 Nov 2014 17:18:23 +0900 Subject: [PATCH] #81. always added agent's identity to channel changed pinpointsocket packet name. (enableWorker -> handShake, enableWorkerConfirm -> handShakeResponse) changed pinpointsocket handshake rule. (when support server mode -> always) changed pinpointsocket handshake default retry count. (3 -> Integer.MAX_VALUE) added subCode property to HandShakeResponse. --- .../cluster/ChannelContextClusterPoint.java | 10 +-- .../zookeeper/ZookeeperLatestJobWorker.java | 14 ++-- .../ZookeeperProfilerClusterManager.java | 6 +- ...e.java => AgentHandShakePropertyType.java} | 12 ++- .../collector/receiver/tcp/TCPReceiver.java | 21 +++-- .../cluster/ClusterPointRouterTest.java | 13 ++-- .../cluster/ClusterPointRouterTest2.java | 8 +- .../cluster/zookeeper/ZookeeperTestUtils.java | 7 +- .../resources/applicationContext-test.xml | 8 ++ ...e.java => AgentHandShakePropertyType.java} | 12 ++- .../pinpoint/profiler/AgentInformation.java | 16 ++-- .../pinpoint/profiler/DefaultAgent.java | 14 ++-- .../profiler/AgentInfoSenderTest.java | 6 +- .../sender/TcpDataSenderReconnectTest.java | 6 +- .../profiler/sender/TcpDataSenderTest.java | 6 +- .../RequestResponseServerMessageListener.java | 9 ++- .../pinpoint/rpc/client/PinpointSocket.java | 12 +-- .../rpc/client/PinpointSocketHandler.java | 76 ++++++++++--------- .../client/ReconnectStateSocketHandler.java | 4 +- .../pinpoint/rpc/client/SocketHandler.java | 2 +- .../navercorp/pinpoint/rpc/client/State.java | 27 ++++++- .../pinpoint/rpc/codec/PacketDecoder.java | 12 +-- ...acket.java => ControlHandShakePacket.java} | 16 ++-- ...va => ControlHandShakeResponsePacket.java} | 25 +++--- .../rpc/packet/HandShakeResponseCode.java | 71 +++++++++++++++++ .../rpc/packet/HandShakeResponseType.java | 40 ++++++++++ .../pinpoint/rpc/packet/PacketType.java | 4 +- .../rpc/server/PinpointServerSocket.java | 57 ++++++++------ .../rpc/server/ServerMessageListener.java | 3 +- .../SimpleLoggingServerMessageListener.java | 7 +- .../navercorp/pinpoint/rpc/util/MapUtils.java | 23 +++++- ...e.java => AgentHandShakePropertyType.java} | 12 ++- .../rpc/server/ControlPacketServerTest.java | 26 ++++--- .../pinpoint/rpc/server/EventListnerTest.java | 20 ++--- .../rpc/server/MessageListenerTest.java | 40 +++++++--- .../rpc/server/TestSeverMessageListener.java | 9 ++- ...CommandHeaderTBaseDeserializerFactory.java | 2 - .../CommandHeaderTBaseSerializerFactory.java | 3 - .../web/server/PinpointSocketManager.java | 11 ++- .../pinpoint/web/cluster/ClusterTest.java | 1 - 40 files changed, 447 insertions(+), 224 deletions(-) rename collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/{AgentPropertiesType.java => AgentHandShakePropertyType.java} (80%) rename profiler/src/main/java/com/navercorp/pinpoint/profiler/{AgentPropertiesType.java => AgentHandShakePropertyType.java} (79%) rename rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/{ControlEnableWorkerPacket.java => ControlHandShakePacket.java} (68%) rename rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/{ControlEnableWorkerConfirmPacket.java => ControlHandShakeResponsePacket.java} (58%) create mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseCode.java create mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseType.java rename rpc/src/test/java/com/navercorp/pinpoint/rpc/server/{AgentPropertiesType.java => AgentHandShakePropertyType.java} (75%) diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ChannelContextClusterPoint.java b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ChannelContextClusterPoint.java index 2846ecd2c..c0d763a6b 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ChannelContextClusterPoint.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/ChannelContextClusterPoint.java @@ -4,7 +4,7 @@ import java.util.Map; import org.apache.commons.lang.StringUtils; -import com.nhn.pinpoint.collector.receiver.tcp.AgentPropertiesType; +import com.nhn.pinpoint.collector.receiver.tcp.AgentHandShakePropertyType; import com.nhn.pinpoint.rpc.Future; import com.nhn.pinpoint.rpc.server.ChannelContext; import com.nhn.pinpoint.rpc.server.SocketChannel; @@ -30,16 +30,16 @@ public class ChannelContextClusterPoint implements TargetClusterPoint { AssertUtils.assertNotNull(socketChannel, "SocketChannel may not be null."); Map properties = channelContext.getChannelProperties(); - this.version = MapUtils.getString(properties, AgentPropertiesType.VERSION.getName()); + this.version = MapUtils.getString(properties, AgentHandShakePropertyType.VERSION.getName()); AssertUtils.assertTrue(!StringUtils.isBlank(version), "Version may not be null or empty."); - this.applicationName = MapUtils.getString(properties, AgentPropertiesType.APPLICATION_NAME.getName()); + this.applicationName = MapUtils.getString(properties, AgentHandShakePropertyType.APPLICATION_NAME.getName()); AssertUtils.assertTrue(!StringUtils.isBlank(applicationName), "ApplicationName may not be null or empty."); - this.agentId = MapUtils.getString(properties, AgentPropertiesType.AGENT_ID.getName()); + this.agentId = MapUtils.getString(properties, AgentHandShakePropertyType.AGENT_ID.getName()); AssertUtils.assertTrue(!StringUtils.isBlank(agentId), "AgentId may not be null or empty."); - this.startTimeStamp = MapUtils.getLong(properties, AgentPropertiesType.START_TIMESTAMP.getName()); + this.startTimeStamp = MapUtils.getLong(properties, AgentHandShakePropertyType.START_TIMESTAMP.getName()); AssertUtils.assertTrue(startTimeStamp > 0, "StartTimeStamp is must greater than zero."); } diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperLatestJobWorker.java b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperLatestJobWorker.java index 703bd82b7..618f302fa 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperLatestJobWorker.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperLatestJobWorker.java @@ -22,7 +22,7 @@ import com.nhn.pinpoint.collector.cluster.zookeeper.exception.TimeoutException; import com.nhn.pinpoint.collector.cluster.zookeeper.job.DeleteJob; import com.nhn.pinpoint.collector.cluster.zookeeper.job.Job; import com.nhn.pinpoint.collector.cluster.zookeeper.job.UpdateJob; -import com.nhn.pinpoint.collector.receiver.tcp.AgentPropertiesType; +import com.nhn.pinpoint.collector.receiver.tcp.AgentHandShakePropertyType; import com.nhn.pinpoint.common.util.PinpointThreadFactory; import com.nhn.pinpoint.rpc.server.ChannelContext; import com.nhn.pinpoint.rpc.server.PinpointServerSocketStateCode; @@ -359,9 +359,9 @@ public class ZookeeperLatestJobWorker implements Runnable { private boolean checkRequiredProperties(ChannelContext channelContext) { Map agentProperties = channelContext.getChannelProperties(); - final String applicationName = MapUtils.getString(agentProperties, AgentPropertiesType.APPLICATION_NAME.getName()); - final String agentId = MapUtils.getString(agentProperties, AgentPropertiesType.AGENT_ID.getName()); - final Long startTimeStampe = MapUtils.getLong(agentProperties, AgentPropertiesType.START_TIMESTAMP.getName()); + final String applicationName = MapUtils.getString(agentProperties, AgentHandShakePropertyType.APPLICATION_NAME.getName()); + final String agentId = MapUtils.getString(agentProperties, AgentHandShakePropertyType.AGENT_ID.getName()); + final Long startTimeStampe = MapUtils.getLong(agentProperties, AgentHandShakePropertyType.START_TIMESTAMP.getName()); if (StringUtils.isBlank(applicationName) || StringUtils.isBlank(agentId) || startTimeStampe == null || startTimeStampe <= 0) { logger.warn("ApplicationName({}) and AgnetId({}) and startTimeStampe({}) may not be null.", applicationName, agentId); @@ -375,9 +375,9 @@ public class ZookeeperLatestJobWorker implements Runnable { StringBuilder profilerContents = new StringBuilder(); Map agentProperties = channelContext.getChannelProperties(); - final String applicationName = MapUtils.getString(agentProperties, AgentPropertiesType.APPLICATION_NAME.getName()); - final String agentId = MapUtils.getString(agentProperties, AgentPropertiesType.AGENT_ID.getName()); - final Long startTimeStampe = MapUtils.getLong(agentProperties, AgentPropertiesType.START_TIMESTAMP.getName()); + final String applicationName = MapUtils.getString(agentProperties, AgentHandShakePropertyType.APPLICATION_NAME.getName()); + final String agentId = MapUtils.getString(agentProperties, AgentHandShakePropertyType.AGENT_ID.getName()); + final Long startTimeStampe = MapUtils.getLong(agentProperties, AgentHandShakePropertyType.START_TIMESTAMP.getName()); if (StringUtils.isBlank(applicationName) || StringUtils.isBlank(agentId) || startTimeStampe == null || startTimeStampe <= 0) { logger.warn("ApplicationName({}) and AgnetId({}) and startTimeStampe({}) may not be null.", applicationName, agentId); diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java index 9cf9f92bf..ede77401c 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java @@ -16,7 +16,7 @@ import com.nhn.pinpoint.collector.cluster.WorkerState; import com.nhn.pinpoint.collector.cluster.WorkerStateContext; import com.nhn.pinpoint.collector.cluster.zookeeper.job.DeleteJob; import com.nhn.pinpoint.collector.cluster.zookeeper.job.UpdateJob; -import com.nhn.pinpoint.collector.receiver.tcp.AgentPropertiesType; +import com.nhn.pinpoint.collector.receiver.tcp.AgentHandShakePropertyType; import com.nhn.pinpoint.rpc.server.ChannelContext; import com.nhn.pinpoint.rpc.server.PinpointServerSocketStateCode; import com.nhn.pinpoint.rpc.server.SocketChannelStateChangeEventListener; @@ -158,8 +158,8 @@ public class ZookeeperProfilerClusterManager implements SocketChannelStateChange } private boolean skipAgent(Map agentProperties) { - String applicationName = MapUtils.getString(agentProperties, AgentPropertiesType.APPLICATION_NAME.getName()); - String agentId = MapUtils.getString(agentProperties, AgentPropertiesType.AGENT_ID.getName()); + String applicationName = MapUtils.getString(agentProperties, AgentHandShakePropertyType.APPLICATION_NAME.getName()); + String agentId = MapUtils.getString(agentProperties, AgentHandShakePropertyType.AGENT_ID.getName()); if (StringUtils.isBlank(applicationName) || StringUtils.isBlank(agentId)) { return true; diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/AgentPropertiesType.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/AgentHandShakePropertyType.java similarity index 80% rename from collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/AgentPropertiesType.java rename to collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/AgentHandShakePropertyType.java index 5b5410e3c..da68c9e4c 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/AgentPropertiesType.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/AgentHandShakePropertyType.java @@ -4,12 +4,14 @@ import java.util.Map; import com.nhn.pinpoint.rpc.util.ClassUtils; -public enum AgentPropertiesType { +public enum AgentHandShakePropertyType { // 해당 객체는 profiler, collector 양쪽에 함꼐 있음 // 변경시 함께 변경 필요 // map으로 처리하기 때문에 이전 파라미터 제거 대신 추가할 경우 확장성에는 문제가 없음 + SUPPORT_SERVER("supportServer", Boolean.class), + HOSTNAME("hostName", String.class), IP("ip", String.class), AGENT_ID("agentId", String.class), @@ -23,7 +25,7 @@ public enum AgentPropertiesType { private final String name; private final Class clazzType; - private AgentPropertiesType(String name, Class clazzType) { + private AgentHandShakePropertyType(String name, Class clazzType) { this.name = name; this.clazzType = clazzType; } @@ -37,9 +39,13 @@ public enum AgentPropertiesType { } public static boolean hasAllType(Map properties) { - for (AgentPropertiesType type : AgentPropertiesType.values()) { + for (AgentHandShakePropertyType type : AgentHandShakePropertyType.values()) { Object value = properties.get(type.getName()); + if (type == SUPPORT_SERVER) { + continue; + } + if (value == null) { return false; } 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 af9c3e847..567988881 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 @@ -26,12 +26,14 @@ import com.nhn.pinpoint.collector.receiver.DispatchHandler; import com.nhn.pinpoint.collector.util.PacketUtils; import com.nhn.pinpoint.common.util.ExecutorFactory; import com.nhn.pinpoint.common.util.PinpointThreadFactory; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; import com.nhn.pinpoint.rpc.server.SocketChannel; +import com.nhn.pinpoint.rpc.util.MapUtils; import com.nhn.pinpoint.thrift.io.DeserializerFactory; import com.nhn.pinpoint.thrift.io.Header; import com.nhn.pinpoint.thrift.io.HeaderTBaseDeserializer; @@ -136,17 +138,22 @@ public class TCPReceiver { } @Override - public int handleEnableWorker(Map properties) { + public HandShakeResponseCode handleHandShake(Map properties) { if (properties == null) { - return ControlEnableWorkerConfirmPacket.ILLEGAL_PROTOCOL; + return HandShakeResponseType.ProtocolError.PROTOCOL_ERROR; } - boolean hasAllType = AgentPropertiesType.hasAllType(properties); + boolean hasAllType = AgentHandShakePropertyType.hasAllType(properties); if (!hasAllType) { - return ControlEnableWorkerConfirmPacket.INVALID_PROPERTIES; + return HandShakeResponseType.PropertyError.PROPERTY_ERROR; } - - return ControlEnableWorkerConfirmPacket.SUCCESS; + + boolean supportServer = MapUtils.getBoolean(properties, AgentHandShakePropertyType.SUPPORT_SERVER.getName(), true); + if (supportServer) { + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; + } else { + return HandShakeResponseType.Success.SIMPLEX_COMMUNICATION; + } } }); this.pinpointServerSocket.bind(bindAddress, port); 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 0b4e5007f..105267106 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 @@ -1,5 +1,7 @@ package com.nhn.pinpoint.collector.cluster; +import static org.mockito.Mockito.mock; + import java.net.InetSocketAddress; import java.util.ArrayList; import java.util.HashMap; @@ -7,7 +9,6 @@ import java.util.List; import java.util.Map; import junit.framework.Assert; -import static org.mockito.Mockito.*; import org.junit.Test; import org.junit.runner.RunWith; @@ -19,8 +20,8 @@ import org.springframework.test.context.junit4.SpringJUnit4ClassRunner; import com.nhn.pinpoint.collector.receiver.tcp.AgentProperties; import com.nhn.pinpoint.collector.util.CollectorUtils; -import com.nhn.pinpoint.rpc.PinpointSocketException; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.ChannelContext; @@ -108,9 +109,9 @@ public class ClusterPointRouterTest { } @Override - public int handleEnableWorker(Map properties) { - logger.warn("do handleEnableWorker {}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + public HandShakeResponseCode handleHandShake(Map properties) { + logger.warn("do HandShake {}", properties); + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } diff --git a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest2.java b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest2.java index b3f6cee0e..6e35b9f15 100644 --- a/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest2.java +++ b/collector/src/test/java/com/navercorp/pinpoint/collector/cluster/ClusterPointRouterTest2.java @@ -24,9 +24,9 @@ import com.nhn.pinpoint.collector.receiver.tcp.AgentProperties; import com.nhn.pinpoint.collector.util.CollectorUtils; 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.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.ChannelContext; @@ -153,9 +153,9 @@ public class ClusterPointRouterTest2 { } @Override - public int handleEnableWorker(Map properties) { + public HandShakeResponseCode handleHandShake(Map properties) { logger.warn("do handleEnableWorker {}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } 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 c5958f84f..fb95e3066 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 @@ -10,7 +10,8 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.collector.receiver.tcp.AgentProperties; import com.nhn.pinpoint.rpc.client.MessageListener; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; @@ -85,9 +86,9 @@ final class ZookeeperTestUtils { } @Override - public int handleEnableWorker(Map properties) { + public HandShakeResponseCode handleHandShake(Map properties) { LOGGER.warn("do handleEnableWorker {}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } diff --git a/collector/src/test/resources/applicationContext-test.xml b/collector/src/test/resources/applicationContext-test.xml index 676233c4d..0fbe7848b 100644 --- a/collector/src/test/resources/applicationContext-test.xml +++ b/collector/src/test/resources/applicationContext-test.xml @@ -48,6 +48,14 @@ + + + + + + + + diff --git a/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentPropertiesType.java b/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentHandShakePropertyType.java similarity index 79% rename from profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentPropertiesType.java rename to profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentHandShakePropertyType.java index c158ca019..3952228bf 100644 --- a/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentPropertiesType.java +++ b/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentHandShakePropertyType.java @@ -7,12 +7,14 @@ import com.nhn.pinpoint.rpc.util.ClassUtils; /** * @author koo.taejin */ -public enum AgentPropertiesType { +public enum AgentHandShakePropertyType { // 해당 객체는 profiler, collector 양쪽에 함꼐 있음 // 변경시 함께 변경 필요 // map으로 처리하기 때문에 이전 파라미터 제거 대신 추가할 경우 확장성에는 문제가 없음 + SUPPORT_SERVER("supportServer", Boolean.class), + HOSTNAME("hostName", String.class), IP("ip", String.class), AGENT_ID("agentId", String.class), @@ -26,7 +28,7 @@ public enum AgentPropertiesType { private final String name; private final Class clazzType; - private AgentPropertiesType(String name, Class clazzType) { + private AgentHandShakePropertyType(String name, Class clazzType) { this.name = name; this.clazzType = clazzType; } @@ -40,7 +42,11 @@ public enum AgentPropertiesType { } public static boolean hasAllType(Map properties) { - for (AgentPropertiesType type : AgentPropertiesType.values()) { + for (AgentHandShakePropertyType type : AgentHandShakePropertyType.values()) { + if (type == SUPPORT_SERVER) { + continue; + } + Object value = properties.get(type.getName()); if (value == null) { diff --git a/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentInformation.java b/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentInformation.java index 2c3cba953..9d8e16390 100644 --- a/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentInformation.java +++ b/profiler/src/main/java/com/navercorp/pinpoint/profiler/AgentInformation.java @@ -86,14 +86,14 @@ public class AgentInformation { public Map toMap() { Map map = new HashMap(); - map.put(AgentPropertiesType.AGENT_ID.getName(), this.agentId); - map.put(AgentPropertiesType.APPLICATION_NAME.getName(), this.applicationName); - map.put(AgentPropertiesType.HOSTNAME.getName(), this.machineName); - map.put(AgentPropertiesType.IP.getName(), this.hostIp); - map.put(AgentPropertiesType.PID.getName(), this.pid); - map.put(AgentPropertiesType.SERVICE_TYPE.getName(), this.serverType); - map.put(AgentPropertiesType.START_TIMESTAMP.getName(), this.startTime); - map.put(AgentPropertiesType.VERSION.getName(), this.version); + map.put(AgentHandShakePropertyType.AGENT_ID.getName(), this.agentId); + map.put(AgentHandShakePropertyType.APPLICATION_NAME.getName(), this.applicationName); + map.put(AgentHandShakePropertyType.HOSTNAME.getName(), this.machineName); + map.put(AgentHandShakePropertyType.IP.getName(), this.hostIp); + map.put(AgentHandShakePropertyType.PID.getName(), this.pid); + map.put(AgentHandShakePropertyType.SERVICE_TYPE.getName(), this.serverType); + map.put(AgentHandShakePropertyType.START_TIMESTAMP.getName(), this.startTime); + map.put(AgentHandShakePropertyType.VERSION.getName(), this.version); return map; } 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 8e11b9704..5cf8889dd 100644 --- a/profiler/src/main/java/com/navercorp/pinpoint/profiler/DefaultAgent.java +++ b/profiler/src/main/java/com/navercorp/pinpoint/profiler/DefaultAgent.java @@ -248,17 +248,21 @@ public class DefaultAgent implements Agent { } protected PinpointSocketFactory createPinpointSocketFactory(boolean isSupportServerMode) { - Map properties = this.agentInformation.toMap(); - PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); - pinpointSocketFactory.setTimeoutMillis(1000 * 5); - pinpointSocketFactory.setProperties(properties); + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + pinpointSocketFactory.setTimeoutMillis(1000 * 5); + + Map properties = this.agentInformation.toMap(); if (isSupportServerMode) { CommandDispatcher.Builder builder = new CommandDispatcher.Builder(); - pinpointSocketFactory.setMessageListener(builder.build()); + + properties.put(AgentHandShakePropertyType.SUPPORT_SERVER.getName(), true); + } else { + properties.put(AgentHandShakePropertyType.SUPPORT_SERVER.getName(), false); } + pinpointSocketFactory.setProperties(properties); return pinpointSocketFactory; } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java index fda385ca4..cf30aade9 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java @@ -19,6 +19,8 @@ import com.nhn.pinpoint.profiler.sender.TcpDataSender; import com.nhn.pinpoint.rpc.PinpointSocketException; import com.nhn.pinpoint.rpc.client.PinpointSocket; import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; @@ -276,8 +278,8 @@ public class AgentInfoSenderTest { } @Override - public int handleEnableWorker(Map arg0) { - return 0; + public HandShakeResponseCode handleHandShake(Map arg0) { + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } 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 68b83f4b3..808323cc8 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 @@ -11,6 +11,8 @@ import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.rpc.PinpointSocketException; import com.nhn.pinpoint.rpc.client.PinpointSocket; import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; @@ -46,8 +48,8 @@ public class TcpDataSenderReconnectTest { } @Override - public int handleEnableWorker(Map properties) { - return 0; + public HandShakeResponseCode handleHandShake(Map properties) { + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } }); server.bind(HOST, PORT); 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 7eacbcde0..5d8b723cb 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 @@ -17,6 +17,8 @@ import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.rpc.PinpointSocketException; import com.nhn.pinpoint.rpc.client.PinpointSocket; import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; @@ -56,8 +58,8 @@ public class TcpDataSenderTest { } @Override - public int handleEnableWorker(Map arg0) { - return 0; + public HandShakeResponseCode handleHandShake(Map arg0) { + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } }); server.bind(HOST, PORT); 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 6f87f4755..06a454bfc 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java @@ -5,7 +5,8 @@ import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; @@ -33,9 +34,9 @@ public class RequestResponseServerMessageListener implements ServerMessageListen } @Override - public int handleEnableWorker(Map properties) { - logger.info("handleEnableWorker {}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + public HandShakeResponseCode handleHandShake(Map properties) { + logger.info("handle handShake {}", properties); + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } 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 a2597f787..28c6c3e5c 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 @@ -37,13 +37,9 @@ public class PinpointSocket { public PinpointSocket(SocketHandler socketHandler) { AssertUtils.assertNotNull(socketHandler, "socketHandler"); - - if (socketHandler.isSupportServerMode()) { - socketHandler.turnOnServerMode(); - } - + socketHandler.doHandShake(); + this.socketHandler = socketHandler; - socketHandler.setPinpointSocket(this); } @@ -58,9 +54,7 @@ public class PinpointSocket { logger.warn("reconnectSocketHandler:{}", socketHandler); // Pinpoint 소켓 내부 객체가 되기전에 listener를 먼저 등록 - if (socketHandler.isSupportServerMode()) { - socketHandler.turnOnServerMode(); - } + socketHandler.doHandShake(); this.socketHandler = socketHandler; 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 e8460bce9..bae5d6e17 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 @@ -28,8 +28,9 @@ import com.nhn.pinpoint.rpc.PinpointSocketException; import com.nhn.pinpoint.rpc.ResponseMessage; import com.nhn.pinpoint.rpc.control.ProtocolException; import com.nhn.pinpoint.rpc.packet.ClientClosePacket; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakePacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakeResponsePacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; import com.nhn.pinpoint.rpc.packet.Packet; import com.nhn.pinpoint.rpc.packet.PacketType; import com.nhn.pinpoint.rpc.packet.PingPacket; @@ -59,7 +60,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke private static final long DEFAULT_TIMEOUTMILLIS = 3 * 1000; private static final long DEFAULT_ENABLE_WORKER_PACKET_DELAY = 60 * 1000 * 1; - private static final int DEFAULT_ENABLE_WORKER_PACKET_RETRY_COUNT = 3; + private static final int DEFAULT_ENABLE_WORKER_PACKET_RETRY_COUNT = Integer.MAX_VALUE; private final Logger logger = LoggerFactory.getLogger(this.getClass()); @@ -86,7 +87,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke 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."); - private final ChannelFutureListener enableWorkerWriteFailFutureListener = new WriteFailFutureListener(this.logger, "enableWorker write fail.", "enableWorker write success."); + private final ChannelFutureListener handShakeFailFutureListener = new WriteFailFutureListener(this.logger, "handShake write fail.", "handShake write success."); public PinpointSocketHandler(PinpointSocketFactory pinpointSocketFactory) { this(pinpointSocketFactory, DEFAULT_PING_DELAY, DEFAULT_ENABLE_WORKER_PACKET_DELAY, DEFAULT_TIMEOUTMILLIS); @@ -224,12 +225,12 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke write.addListener(pingWriteFailFutureListener); } - private class RegisterEnableWorkerPacketJob implements TimerTask { + private class HandShakeJob implements TimerTask { private final int maxRetryCount; private final AtomicInteger currentCount; - public RegisterEnableWorkerPacketJob(int maxRetryCount) { + public HandShakeJob(int maxRetryCount) { this.maxRetryCount = maxRetryCount; this.currentCount = new AtomicInteger(0); } @@ -237,7 +238,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke @Override public void run(Timeout timeout) throws Exception { if (timeout.isCancelled()) { - reservationEnableWorkerPacketJob(this); + reservationHandShakeJob(this); return; } if (isClosed()) { @@ -247,8 +248,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke if (state.getState() == State.RUN) { incrementCurrentRetryCount(); - sendEnableWorkerPacket(); - reservationEnableWorkerPacketJob(this); + sendHandShakePacket(); + reservationHandShakeJob(this); } } @@ -266,7 +267,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke } - private void reservationEnableWorkerPacketJob(RegisterEnableWorkerPacketJob task) { + private void reservationHandShakeJob(HandShakeJob task) { if (task.getCurrentRetryCount() >= task.getMaxRetryCount()) { return; } @@ -274,19 +275,20 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke this.channelTimer.newTimeout(task, enableWorkerPacketDelay, TimeUnit.MILLISECONDS); } - void sendEnableWorkerPacket() { + void sendHandShakePacket() { if (!isRun()) { return; } - logger.debug("write EnableWorkerPacket {}", channel); + Map properties = this.pinpointSocketFactory.getProperties(); + logger.debug("write HandShakePakcet channel:{}, property:{}..", channel, properties); try { - Map properties = this.pinpointSocketFactory.getProperties(); byte[] payload = ControlMessageEnDeconderUtils.encode(properties); - ControlEnableWorkerPacket packet = new ControlEnableWorkerPacket(payload); + + ControlHandShakePacket packet = new ControlHandShakePacket(payload); final ChannelFuture write = this.channel.write(packet); - write.addListener(enableWorkerWriteFailFutureListener); + write.addListener(handShakeFailFutureListener); } catch (ProtocolException e) { logger.warn(e.getMessage(), e); } @@ -445,8 +447,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke case PacketType.CONTROL_SERVER_CLOSE: messageReceivedServerClosed(e.getChannel()); return; - case PacketType.CONTROL_ENABLE_WORKER_CONFIRM: - messageReceivedEnableWorkerConfirm((ControlEnableWorkerConfirmPacket)message, e.getChannel()); + case PacketType.CONTROL_HANDSHAKE_RESPONSE: + messageReceivedHandShakeResponse((ControlHandShakeResponsePacket)message, e.getChannel()); return; default: logger.warn("unexpectedMessage received:{} address:{}", message, e.getRemoteAddress()); @@ -462,32 +464,38 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke state.setState(State.RECONNECT); } - private void messageReceivedEnableWorkerConfirm(ControlEnableWorkerConfirmPacket message, Channel channel) { - int code = getRegisterAgentConfirmPacketCode(message.getPayload()); + private void messageReceivedHandShakeResponse(ControlHandShakeResponsePacket message, Channel channel) { + HandShakeResponseCode code = getHandShakeResponseCode(message.getPayload()); - logger.info("EnableWorkerConfirm Packet({}) code={} received. {}", message, code, channel); + logger.info("HandShake Response Packet({}) code={} received. {}", message, code, channel); // reconnect 상태로 변경한다. - if (code == ControlEnableWorkerConfirmPacket.SUCCESS || code == ControlEnableWorkerConfirmPacket.ALREADY_REGISTER) { + if (code == HandShakeResponseCode.SUCCESS || code == HandShakeResponseCode.ALREADY_KNOWN) { + state.changeRunSimplexCommunication(); + } else if (code == HandShakeResponseCode.DUPLEX_COMMUNICATION || code == HandShakeResponseCode.ALREADY_DUPLEX_COMMUNICATION) { state.changeRunDuplexCommunication(); + } else if (code == HandShakeResponseCode.SIMPLEX_COMMUNICATION || code == HandShakeResponseCode.ALREADY_SIMPLEX_COMMUNICATION) { + state.changeRunSimplexCommunication(); } else { - logger.warn("Invalid EnableWorkerConfirm Packet ({}) code={} received. {}", message, code, channel); + logger.warn("Invalid HandShake Packet ({}) code={} received. {}", message, code, channel); } } - private int getRegisterAgentConfirmPacketCode(byte[] payload) { - Map result = null; + private HandShakeResponseCode getHandShakeResponseCode(byte[] payload) { try { - result = (Map) ControlMessageEnDeconderUtils.decode(payload); + Map result = (Map) ControlMessageEnDeconderUtils.decode(payload); + + int code = MapUtils.getInteger(result, ControlHandShakeResponsePacket.CODE, -1); + int subCode = MapUtils.getInteger(result, ControlHandShakeResponsePacket.SUB_CODE, -1); + + return HandShakeResponseCode.getValue(code, subCode); } catch (ProtocolException e) { logger.warn(e.getMessage(), e); } - - int code = MapUtils.getInteger(result, "code", -1); - - return code; + + return HandShakeResponseCode.UNKOWN_CODE; } - + @Override public void exceptionCaught(ChannelHandlerContext ctx, ExceptionEvent e) throws Exception { Throwable cause = e.getCause(); @@ -637,12 +645,12 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke } @Override - public void turnOnServerMode() { + public void doHandShake() { // MessageListener 등록시 EnableWorkerPacket전달 - sendEnableWorkerPacket(); + sendHandShakePacket(); - RegisterEnableWorkerPacketJob job = new RegisterEnableWorkerPacketJob(enableWorkerPacketRetryCount); - reservationEnableWorkerPacketJob(job); + HandShakeJob job = new HandShakeJob(enableWorkerPacketRetryCount); + reservationHandShakeJob(job); } class SocketHandlerContext { 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 950551fad..6765ed316 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 @@ -93,8 +93,8 @@ public class ReconnectStateSocketHandler implements SocketHandler { } @Override - public void turnOnServerMode() { - throw new UnsupportedOperationException(); + public void doHandShake() { +// throw new UnsupportedOperationException(); } } 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 5d2c4407c..9f69cf4e6 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 @@ -42,6 +42,6 @@ public interface SocketHandler { boolean isSupportServerMode(); - void turnOnServerMode(); + void doHandShake(); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/State.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/State.java index e4d5acd36..d85d93f7a 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/State.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/State.java @@ -20,9 +20,10 @@ public class State { public static final int INIT = 0; public static final int RUN = 1; public static final int RUN_DUPLEX_COMMUNICATION = 2; - public static final int CLOSED = 3; + public static final int RUN_SIMPLEX_COMMUNICATION = 3; + public static final int CLOSED = 4; // 이 상태가 있어야 되나? - public static final int RECONNECT = 4; + public static final int RECONNECT = 5; private final AtomicInteger state = new AtomicInteger(INIT); @@ -33,11 +34,11 @@ public class State { public boolean isRun() { int code = state.get(); - return code == RUN || code == RUN_DUPLEX_COMMUNICATION; + return isRun(code); } public boolean isRun(int code) { - return code == RUN || code == RUN_DUPLEX_COMMUNICATION; + return code == RUN || code == RUN_DUPLEX_COMMUNICATION || code == RUN_SIMPLEX_COMMUNICATION; } public boolean isClosed() { @@ -69,6 +70,22 @@ public class State { } throw new IllegalStateException("InvalidState current:" + getString(current) + " change:" + getString(RUN_DUPLEX_COMMUNICATION)); } + + public boolean changeRunSimplexCommunication() { + logger.debug("State Will Be Changed {}.", getString(RUN_SIMPLEX_COMMUNICATION)); + final int current = state.get(); + if (current == INIT) { + return this.state.compareAndSet(INIT, RUN_SIMPLEX_COMMUNICATION); + } else if(current == INIT_RECONNECT) { + return this.state.compareAndSet(INIT_RECONNECT, RUN_SIMPLEX_COMMUNICATION); + } else if (current == RUN) { + return this.state.compareAndSet(RUN, RUN_SIMPLEX_COMMUNICATION); + } else if (current == RUN_SIMPLEX_COMMUNICATION) { + return true; + } + throw new IllegalStateException("InvalidState current:" + getString(current) + " change:" + getString(RUN_SIMPLEX_COMMUNICATION)); + + } public boolean changeClosed(int before) { logger.debug("State Will Be Changed {} -> {}.", getString(before), getString(CLOSED)); @@ -97,6 +114,8 @@ public class State { return "RUN"; case RUN_DUPLEX_COMMUNICATION: return "RUN_DUPLEX_COMMUNICATION"; + case RUN_SIMPLEX_COMMUNICATION: + return "RUN_SIMPLEX_COMMUNICATION"; case CLOSED: return "CLOSED"; case RECONNECT: diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/codec/PacketDecoder.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/codec/PacketDecoder.java index edf5ca6b8..790940a0d 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/codec/PacketDecoder.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/codec/PacketDecoder.java @@ -10,8 +10,8 @@ import org.slf4j.LoggerFactory; import com.nhn.pinpoint.rpc.client.WriteFailFutureListener; import com.nhn.pinpoint.rpc.packet.ClientClosePacket; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakeResponsePacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakePacket; import com.nhn.pinpoint.rpc.packet.PacketType; import com.nhn.pinpoint.rpc.packet.PingPacket; import com.nhn.pinpoint.rpc.packet.PongPacket; @@ -78,9 +78,9 @@ public class PacketDecoder extends FrameDecoder { readPong(packetType, buffer); // pong 도 그냥 버리자. return null; - case PacketType.CONTROL_ENABLE_WORKER: + case PacketType.CONTROL_HANDSHAKE: return readEnableWorker(packetType, buffer); - case PacketType.CONTROL_ENABLE_WORKER_CONFIRM: + case PacketType.CONTROL_HANDSHAKE_RESPONSE: return readEnableWorkerConfirm(packetType, buffer); } logger.error("invalid packetType received. packetType:{}, channel:{}", packetType, channel); @@ -160,11 +160,11 @@ public class PacketDecoder extends FrameDecoder { } private Object readEnableWorker(short packetType, ChannelBuffer buffer) { - return ControlEnableWorkerPacket.readBuffer(packetType, buffer); + return ControlHandShakePacket.readBuffer(packetType, buffer); } private Object readEnableWorkerConfirm(short packetType, ChannelBuffer buffer) { - return ControlEnableWorkerConfirmPacket.readBuffer(packetType, buffer); + return ControlHandShakeResponsePacket.readBuffer(packetType, buffer); } } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlEnableWorkerPacket.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlHandShakePacket.java similarity index 68% rename from rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlEnableWorkerPacket.java rename to rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlHandShakePacket.java index e31c661f9..09c76d97a 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlEnableWorkerPacket.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlHandShakePacket.java @@ -6,34 +6,34 @@ import org.jboss.netty.buffer.ChannelBuffers; /** * @author koo.taejin */ -public class ControlEnableWorkerPacket extends ControlPacket { +public class ControlHandShakePacket extends ControlPacket { - public ControlEnableWorkerPacket(byte[] payload) { + public ControlHandShakePacket(byte[] payload) { super(payload); } - public ControlEnableWorkerPacket(int requestId, byte[] payload) { + public ControlHandShakePacket(int requestId, byte[] payload) { super(payload); setRequestId(requestId); } @Override public short getPacketType() { - return PacketType.CONTROL_ENABLE_WORKER; + return PacketType.CONTROL_HANDSHAKE; } @Override public ChannelBuffer toBuffer() { ChannelBuffer header = ChannelBuffers.buffer(2 + 4 + 4); - header.writeShort(PacketType.CONTROL_ENABLE_WORKER); + header.writeShort(PacketType.CONTROL_HANDSHAKE); header.writeInt(getRequestId()); return PayloadPacket.appendPayload(header, payload); } - public static ControlEnableWorkerPacket readBuffer(short packetType, ChannelBuffer buffer) { - assert packetType == PacketType.CONTROL_ENABLE_WORKER; + public static ControlHandShakePacket readBuffer(short packetType, ChannelBuffer buffer) { + assert packetType == PacketType.CONTROL_HANDSHAKE; if (buffer.readableBytes() < 8) { buffer.resetReaderIndex(); @@ -45,7 +45,7 @@ public class ControlEnableWorkerPacket extends ControlPacket { if (payload == null) { return null; } - final ControlEnableWorkerPacket helloPacket = new ControlEnableWorkerPacket(payload.array()); + final ControlHandShakePacket helloPacket = new ControlHandShakePacket(payload.array()); helloPacket.setRequestId(messageId); return helloPacket; } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlEnableWorkerConfirmPacket.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlHandShakeResponsePacket.java similarity index 58% rename from rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlEnableWorkerConfirmPacket.java rename to rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlHandShakeResponsePacket.java index bbcc5c465..364a054da 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlEnableWorkerConfirmPacket.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/ControlHandShakeResponsePacket.java @@ -6,40 +6,37 @@ import org.jboss.netty.buffer.ChannelBuffers; /** * @author koo.taejin */ -public class ControlEnableWorkerConfirmPacket extends ControlPacket { +public class ControlHandShakeResponsePacket extends ControlPacket { - public static final int SUCCESS = 0; - public static final int ALREADY_REGISTER = 1; - public static final int INVALID_PROPERTIES = 2; - public static final int ILLEGAL_PROTOCOL = 3; - public static final int UNKNOWN_ERROR = 4; - - public ControlEnableWorkerConfirmPacket(byte[] payload) { + public static final String CODE = "code"; + public static final String SUB_CODE = "subCode"; + + public ControlHandShakeResponsePacket(byte[] payload) { super(payload); } - public ControlEnableWorkerConfirmPacket(int requestId, byte[] payload) { + public ControlHandShakeResponsePacket(int requestId, byte[] payload) { super(payload); setRequestId(requestId); } @Override public short getPacketType() { - return PacketType.CONTROL_ENABLE_WORKER_CONFIRM; + return PacketType.CONTROL_HANDSHAKE_RESPONSE; } @Override public ChannelBuffer toBuffer() { ChannelBuffer header = ChannelBuffers.buffer(2 + 4 + 4); - header.writeShort(PacketType.CONTROL_ENABLE_WORKER_CONFIRM); + header.writeShort(PacketType.CONTROL_HANDSHAKE_RESPONSE); header.writeInt(getRequestId()); return PayloadPacket.appendPayload(header, payload); } - public static ControlEnableWorkerConfirmPacket readBuffer(short packetType, ChannelBuffer buffer) { - assert packetType == PacketType.CONTROL_ENABLE_WORKER_CONFIRM; + public static ControlHandShakeResponsePacket readBuffer(short packetType, ChannelBuffer buffer) { + assert packetType == PacketType.CONTROL_HANDSHAKE_RESPONSE; if (buffer.readableBytes() < 8) { buffer.resetReaderIndex(); @@ -51,7 +48,7 @@ public class ControlEnableWorkerConfirmPacket extends ControlPacket { if (payload == null) { return null; } - final ControlEnableWorkerConfirmPacket helloPacket = new ControlEnableWorkerConfirmPacket(payload.array()); + final ControlHandShakeResponsePacket helloPacket = new ControlHandShakeResponsePacket(payload.array()); helloPacket.setRequestId(messageId); return helloPacket; } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseCode.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseCode.java new file mode 100644 index 000000000..9c39c81d3 --- /dev/null +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseCode.java @@ -0,0 +1,71 @@ +package com.nhn.pinpoint.rpc.packet; + +public enum HandShakeResponseCode { + + SUCCESS(0, 0, "Success."), + SIMPLEX_COMMUNICATION(0, 1, "Simplex Connection successfully established."), + DUPLEX_COMMUNICATION(0, 2, "Duplex Connection successfully established."), + + ALREADY_KNOWN(1, 0, "Already Known."), + ALREADY_SIMPLEX_COMMUNICATION(1, 1, "Already Simplex Connection eastablished."), + ALREADY_DUPLEX_COMMUNICATION(1, 2, "Already Duplex Connection established."), + + PROPERTY_ERROR(2, 0, "Property error."), + + PROTOCOL_ERROR(3, 0, "Illegal protocol error."), + + UNKNOWN_ERROR(4, 0, "Unkown Error."), + + UNKOWN_CODE(-1, -1, "Unkown Code."); + + private final int code; + private final int subCode; + private final String codeMessage; + + private HandShakeResponseCode(int code, int subCode, String codeMessage) { + this.code = code; + this.subCode = subCode; + this.codeMessage = codeMessage; + } + + public int getCode() { + return code; + } + + public int getSubCode() { + return subCode; + } + + public String getCodeMessage() { + return codeMessage; + } + + public static HandShakeResponseCode getValue(int code, int subCode) { + for (HandShakeResponseCode value : HandShakeResponseCode.values()) { + if (code != value.getCode()) { + continue; + } + + if (subCode != value.getSubCode()) { + continue; + } + + return value; + } + + for (HandShakeResponseCode value : HandShakeResponseCode.values()) { + if (code != value.getCode()) { + continue; + } + + if (0 != value.getSubCode()) { + continue; + } + + return value; + } + + return UNKOWN_CODE; + } + +} diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseType.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseType.java new file mode 100644 index 000000000..bc035b59c --- /dev/null +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/HandShakeResponseType.java @@ -0,0 +1,40 @@ +package com.nhn.pinpoint.rpc.packet; + +public class HandShakeResponseType { + + public static class Success { + public static final int CODE = 0; + + public static final HandShakeResponseCode SUCCESS = HandShakeResponseCode.SUCCESS; + + public static final HandShakeResponseCode SIMPLEX_COMMUNICATION = HandShakeResponseCode.SIMPLEX_COMMUNICATION; + public static final HandShakeResponseCode DUPLEX_COMMUNICATION = HandShakeResponseCode.DUPLEX_COMMUNICATION; + } + + public static class AlreadyKnown { + public static final int CODE = 1; + + public static final HandShakeResponseCode ALREADY_KNOWN = HandShakeResponseCode.ALREADY_KNOWN; + public static final HandShakeResponseCode ALREADY_SIMPLEX_COMMUNICATION = HandShakeResponseCode.ALREADY_SIMPLEX_COMMUNICATION; + public static final HandShakeResponseCode ALREADY_DUPLEX_COMMUNICATION = HandShakeResponseCode.ALREADY_DUPLEX_COMMUNICATION; + } + + public static class PropertyError { + public static final int CODE = 2; + + public static final HandShakeResponseCode PROPERTY_ERROR = HandShakeResponseCode.PROPERTY_ERROR; + } + + public static class ProtocolError { + public static final int CODE = 3; + + public static final HandShakeResponseCode PROTOCOL_ERROR = HandShakeResponseCode.PROTOCOL_ERROR; + } + + public static class Error { + public static final int CODE = 4; + + public static final HandShakeResponseCode UNKNOWN_ERROR = HandShakeResponseCode.UNKNOWN_ERROR; + } + +} diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/PacketType.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/PacketType.java index e04aad129..48fbd0cfb 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/PacketType.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/packet/PacketType.java @@ -29,8 +29,8 @@ public class PacketType { public static final short CONTROL_SERVER_CLOSE = 110; // 컨트롤 패킷 - public static final short CONTROL_ENABLE_WORKER = 150; - public static final short CONTROL_ENABLE_WORKER_CONFIRM = 151; + public static final short CONTROL_HANDSHAKE = 150; + public static final short CONTROL_HANDSHAKE_RESPONSE = 151; // ping, pong의 경우 성능상 두고 다른 CONTROL은 이걸로 뺌 public static final short CONTROL_PING = 200; 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 458317038..90b03dfec 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 @@ -38,8 +38,10 @@ import com.nhn.pinpoint.common.util.PinpointThreadFactory; import com.nhn.pinpoint.rpc.PinpointSocketException; import com.nhn.pinpoint.rpc.client.WriteFailFutureListener; import com.nhn.pinpoint.rpc.control.ProtocolException; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakePacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakeResponsePacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.Packet; import com.nhn.pinpoint.rpc.packet.PacketType; import com.nhn.pinpoint.rpc.packet.PingPacket; @@ -229,26 +231,35 @@ public class PinpointServerSocket extends SimpleChannelHandler { case PacketType.APPLICATION_STREAM_PONG: handleStreamPacket((StreamPacket) message, channel); return; - case PacketType.CONTROL_ENABLE_WORKER: - int requestId = ((ControlEnableWorkerPacket)message).getRequestId(); + case PacketType.CONTROL_HANDSHAKE: + int requestId = ((ControlHandShakePacket)message).getRequestId(); - Map properties = decodeSocketProperties((ControlEnableWorkerPacket) message); + Map properties = decodeSocketProperties((ControlHandShakePacket) message); if (properties == null) { - sendEnableWorkerConfirmMessage(requestId, ControlEnableWorkerConfirmPacket.ILLEGAL_PROTOCOL, channel); + sendHandShakeResponseMessage(requestId, HandShakeResponseType.ProtocolError.PROTOCOL_ERROR, channel); return; } - channelContext.setChannelProperties(properties); + HandShakeResponseCode code = messageListener.handleHandShake(properties); - int returnCode = messageListener.handleEnableWorker(properties); - if (returnCode == ControlEnableWorkerConfirmPacket.SUCCESS) { - if (changeStateToRunDuplexCommunication(returnCode, channel)) { - sendEnableWorkerConfirmMessage(requestId, ControlEnableWorkerConfirmPacket.SUCCESS, channel); - } else { - sendEnableWorkerConfirmMessage(requestId, ControlEnableWorkerConfirmPacket.ALREADY_REGISTER, channel); - } + if (code.getCode() != HandShakeResponseType.Success.CODE) { + sendHandShakeResponseMessage(requestId, code, channel); } else { - sendEnableWorkerConfirmMessage(requestId, returnCode, channel); + boolean isSet = channelContext.setChannelProperties(properties); + + if (isSet) { + if (code == HandShakeResponseCode.DUPLEX_COMMUNICATION) { + changeStateToRunDuplexCommunication(code, channel); + } + sendHandShakeResponseMessage(requestId, code, channel); + } else { + if (code == HandShakeResponseCode.DUPLEX_COMMUNICATION) { + sendHandShakeResponseMessage(requestId, HandShakeResponseCode.ALREADY_DUPLEX_COMMUNICATION, channel); + } else { + sendHandShakeResponseMessage(requestId, HandShakeResponseCode.ALREADY_SIMPLEX_COMMUNICATION, channel); + } + } + } return; case PacketType.CONTROL_CLIENT_CLOSE: { @@ -274,7 +285,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { context.getStreamChannelManager().messageReceived(packet); } - private Map decodeSocketProperties(ControlEnableWorkerPacket message) { + private Map decodeSocketProperties(ControlHandShakePacket message) { try { byte[] payload = message.getPayload(); Map properties = (Map) ControlMessageEnDeconderUtils.decode(payload); @@ -286,10 +297,10 @@ public class PinpointServerSocket extends SimpleChannelHandler { return null; } - private boolean changeStateToRunDuplexCommunication(int returnCode, Channel channel) { + private boolean changeStateToRunDuplexCommunication(HandShakeResponseCode returnCode, Channel channel) { ChannelContext context = getChannelContext(channel); - if (returnCode == ControlEnableWorkerConfirmPacket.SUCCESS) { + if (returnCode == HandShakeResponseType.Success.DUPLEX_COMMUNICATION) { if (context.getCurrentStateCode() != PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION) { context.changeStateRunDuplexCommunication(); return true; @@ -299,19 +310,21 @@ public class PinpointServerSocket extends SimpleChannelHandler { return false; } - private void sendEnableWorkerConfirmMessage(int requestId, int returnCode, Channel channel) { + private void sendHandShakeResponseMessage(int requestId, HandShakeResponseCode handShakeResponseCode, Channel channel) { try { + logger.info("write HandShakeResponsePakcet channel:{}, HandShakeResponseCode:{}.", channel, handShakeResponseCode); + Map result = new HashMap(); - result.put("code", returnCode); + result.put(ControlHandShakeResponsePacket.CODE, handShakeResponseCode.getCode()); + result.put(ControlHandShakeResponsePacket.SUB_CODE, handShakeResponseCode.getSubCode()); byte[] resultPayload = ControlMessageEnDeconderUtils.encode(result); - ControlEnableWorkerConfirmPacket packet = new ControlEnableWorkerConfirmPacket(requestId, resultPayload); + ControlHandShakeResponsePacket packet = new ControlHandShakeResponsePacket(requestId, resultPayload); channel.write(packet); } catch (ProtocolException e) { logger.warn(e.getMessage(), e); } - } @Override 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 3ad87234e..1dad640f1 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 @@ -2,6 +2,7 @@ package com.nhn.pinpoint.rpc.server; import java.util.Map; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -14,6 +15,6 @@ public interface ServerMessageListener { // 외부 노출 Channel은 별도의 Tcp Channel로 감싸는걸로 변경할 것. void handleRequest(RequestPacket requestPacket, SocketChannel channel); - int handleEnableWorker(Map properties); + HandShakeResponseCode handleHandShake(Map properties); } 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 ea23e87e0..fa6dc4b21 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 @@ -5,7 +5,8 @@ import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -29,9 +30,9 @@ public class SimpleLoggingServerMessageListener implements ServerMessageListener } @Override - public int handleEnableWorker(Map properties) { + public HandShakeResponseCode handleHandShake(Map properties) { logger.info("handleEnableWorker {}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + return HandShakeResponseType.Success.SUCCESS; } } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/util/MapUtils.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/util/MapUtils.java index 0e90b8f16..1c8ebcb1b 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/util/MapUtils.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/util/MapUtils.java @@ -25,9 +25,26 @@ public class MapUtils { return (String) value; } - return null; + return defaultValue; } + + public static Boolean getBoolean(Map map, String key) { + return getBoolean(map, key, false); + } + + public static Boolean getBoolean(Map map, String key, Boolean defaultValue) { + if (map == null) { + return defaultValue; + } + + final Object value = map.get(key); + if (value instanceof Boolean) { + return (Boolean) value; + } + return defaultValue; + } + public static Integer getInteger(Map map, String key) { return getInteger(map, key, null); @@ -43,7 +60,7 @@ public class MapUtils { return (Integer) value; } - return null; + return defaultValue; } public static Long getLong(Map map, String key) { @@ -60,7 +77,7 @@ public class MapUtils { return (Long) value; } - return null; + return defaultValue; } } diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/AgentPropertiesType.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/AgentHandShakePropertyType.java similarity index 75% rename from rpc/src/test/java/com/navercorp/pinpoint/rpc/server/AgentPropertiesType.java rename to rpc/src/test/java/com/navercorp/pinpoint/rpc/server/AgentHandShakePropertyType.java index 8dd70f795..1cd5878e3 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/AgentPropertiesType.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/AgentHandShakePropertyType.java @@ -4,7 +4,9 @@ import java.util.Map; import com.nhn.pinpoint.rpc.util.ClassUtils; -public enum AgentPropertiesType { +public enum AgentHandShakePropertyType { + + SUPPORT_SERVER("supportServer", Boolean.class), HOSTNAME("hostName", String.class), IP("ip", String.class), @@ -19,7 +21,7 @@ public enum AgentPropertiesType { private final String name; private final Class clazzType; - private AgentPropertiesType(String name, Class clazzType) { + private AgentHandShakePropertyType(String name, Class clazzType) { this.name = name; this.clazzType = clazzType; } @@ -33,9 +35,13 @@ public enum AgentPropertiesType { } public static boolean hasAllType(Map properties) { - for (AgentPropertiesType type : AgentPropertiesType.values()) { + for (AgentHandShakePropertyType type : AgentHandShakePropertyType.values()) { Object value = properties.get(type.getName()); + if (type == SUPPORT_SERVER) { + continue; + } + if (value == null) { return false; } diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java index c5e3c308d..16d5eb570 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java @@ -17,8 +17,10 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.nhn.pinpoint.rpc.control.ProtocolException; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakePacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakeResponsePacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.ResponsePacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -117,7 +119,7 @@ public class ControlPacketServerTest { } // RegisterPacket 등록 성공 메시지를 여러번 보낼 경우 최초는 성공, 두번쨰는 이미 성공 code를 받는지 확인 - // 이후 메시지 전달 가능 확인 + // 이후 메시지 전달 가능 확인 이후도 똑같은 메시지 전달하는 것으로 변경 @Test public void registerAgentTest4() throws Exception { PinpointServerSocket pinpointServerSocket = new PinpointServerSocket(); @@ -156,7 +158,7 @@ public class ControlPacketServerTest { private int sendAndReceiveRegisterPacket(Socket socket, Map properties) throws ProtocolException, IOException { sendRegisterPacket(socket.getOutputStream(), properties); - ControlEnableWorkerConfirmPacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); + ControlHandShakeResponsePacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); Map result = (Map) ControlMessageEnDeconderUtils.decode(packet.getPayload()); return MapUtils.getInteger(result, "code", -1); @@ -170,7 +172,7 @@ public class ControlPacketServerTest { private void sendRegisterPacket(OutputStream outputStream, Map properties) throws ProtocolException, IOException { byte[] payload = ControlMessageEnDeconderUtils.encode(properties); - ControlEnableWorkerPacket packet = new ControlEnableWorkerPacket(1, payload); + ControlHandShakePacket packet = new ControlHandShakePacket(1, payload); ByteBuffer bb = packet.toBuffer().toByteBuffer(0, packet.toBuffer().writerIndex()); sendData(outputStream, bb.array()); @@ -189,14 +191,14 @@ public class ControlPacketServerTest { outputStream.flush(); } - private ControlEnableWorkerConfirmPacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { + private ControlHandShakeResponsePacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { byte[] payload = readData(inputStream); ChannelBuffer cb = ChannelBuffers.wrappedBuffer(payload); short packetType = cb.readShort(); - ControlEnableWorkerConfirmPacket packet = ControlEnableWorkerConfirmPacket.readBuffer(packetType, cb); + ControlHandShakeResponsePacket packet = ControlHandShakeResponsePacket.readBuffer(packetType, cb); return packet; } @@ -247,17 +249,17 @@ public class ControlPacketServerTest { } @Override - public int handleEnableWorker(Map properties) { + public HandShakeResponseCode handleHandShake(Map properties) { if (properties == null) { - return ControlEnableWorkerConfirmPacket.ILLEGAL_PROTOCOL; + return HandShakeResponseType.ProtocolError.PROTOCOL_ERROR; } - boolean hasAllType = AgentPropertiesType.hasAllType(properties); + boolean hasAllType = AgentHandShakePropertyType.hasAllType(properties); if (!hasAllType) { - return ControlEnableWorkerConfirmPacket.INVALID_PROPERTIES; + return HandShakeResponseType.PropertyError.PROPERTY_ERROR; } - return ControlEnableWorkerConfirmPacket.SUCCESS; + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventListnerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventListnerTest.java index 141241a08..ba8f378ca 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventListnerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventListnerTest.java @@ -16,8 +16,10 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.nhn.pinpoint.rpc.control.ProtocolException; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakeResponsePacket; +import com.nhn.pinpoint.rpc.packet.ControlHandShakePacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.ResponsePacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -63,7 +65,7 @@ public class EventListnerTest { private int sendAndReceiveRegisterPacket(Socket socket, Map properties) throws ProtocolException, IOException { sendRegisterPacket(socket.getOutputStream(), properties); - ControlEnableWorkerConfirmPacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); + ControlHandShakeResponsePacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); Map result = (Map) ControlMessageEnDeconderUtils.decode(packet.getPayload()); return MapUtils.getInteger(result, "code", -1); @@ -77,7 +79,7 @@ public class EventListnerTest { private void sendRegisterPacket(OutputStream outputStream, Map properties) throws ProtocolException, IOException { byte[] payload = ControlMessageEnDeconderUtils.encode(properties); - ControlEnableWorkerPacket packet = new ControlEnableWorkerPacket(1, payload); + ControlHandShakePacket packet = new ControlHandShakePacket(1, payload); ByteBuffer bb = packet.toBuffer().toByteBuffer(0, packet.toBuffer().writerIndex()); sendData(outputStream, bb.array()); @@ -96,14 +98,14 @@ public class EventListnerTest { outputStream.flush(); } - private ControlEnableWorkerConfirmPacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { + private ControlHandShakeResponsePacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { byte[] payload = readData(inputStream); ChannelBuffer cb = ChannelBuffers.wrappedBuffer(payload); short packetType = cb.readShort(); - ControlEnableWorkerConfirmPacket packet = ControlEnableWorkerConfirmPacket.readBuffer(packetType, cb); + ControlHandShakeResponsePacket packet = ControlHandShakeResponsePacket.readBuffer(packetType, cb); return packet; } @@ -169,9 +171,9 @@ public class EventListnerTest { } @Override - public int handleEnableWorker(Map properties) { - logger.info("handleEnableWorker {}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + public HandShakeResponseCode handleHandShake(Map properties) { + logger.info("handle HandShake {}", properties); + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/MessageListenerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/MessageListenerTest.java index a9a2ce4e9..02b30685d 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/MessageListenerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/MessageListenerTest.java @@ -17,8 +17,10 @@ import com.nhn.pinpoint.rpc.ResponseMessage; 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.PinpointSocketReconnectEventListener; import com.nhn.pinpoint.rpc.client.SimpleLoggingMessageListener; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.ResponsePacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -31,6 +33,7 @@ public class MessageListenerTest { public void serverMessageListenerTest1() throws InterruptedException { PinpointServerSocket ss = new PinpointServerSocket(); ss.bind("127.0.0.1", 10234); + ss.setMessageListener(new SimpleListener()); PinpointSocketFactory socketFactory1 = createPinpointSocketFactory(); socketFactory1.setMessageListener(new EchoMessageListener()); @@ -46,7 +49,7 @@ public class MessageListenerTest { Thread.sleep(500); List channelContextList = ss.getDuplexCommunicationChannelContext(); - if (channelContextList.size() != 1) { + if (channelContextList.size() != 2) { Assert.fail(); } @@ -64,6 +67,7 @@ public class MessageListenerTest { public void serverMessageListenerTest2() throws InterruptedException { PinpointServerSocket ss = new PinpointServerSocket(); ss.bind("127.0.0.1", 10234); + ss.setMessageListener(new SimpleListener()); EchoMessageListener echoMessageListener = new EchoMessageListener(); @@ -104,6 +108,7 @@ public class MessageListenerTest { public void serverMessageListenerTest3() throws InterruptedException { PinpointServerSocket ss = new PinpointServerSocket(); ss.bind("127.0.0.1", 10234); + ss.setMessageListener(new SimpleListener()); PinpointSocketFactory socketFactory1 = createPinpointSocketFactory(); EchoMessageListener echoMessageListener1 = new EchoMessageListener(); @@ -149,6 +154,7 @@ public class MessageListenerTest { public void serverMessageListenerTest4() throws InterruptedException { PinpointServerSocket ss = new PinpointServerSocket(); ss.bind("127.0.0.1", 10234); + ss.setMessageListener(new SimpleListener()); Map params = getParams(); PinpointSocketFactory socketFactory = createPinpointSocketFactory(params); @@ -161,10 +167,10 @@ public class MessageListenerTest { Thread.sleep(500); - ChannelContext channelContext = getChannelContext("application", "agent", (Long) params.get(AgentPropertiesType.START_TIMESTAMP.getName()), ss.getDuplexCommunicationChannelContext()); + ChannelContext channelContext = getChannelContext("application", "agent", (Long) params.get(AgentHandShakePropertyType.START_TIMESTAMP.getName()), ss.getDuplexCommunicationChannelContext()); Assert.assertNotNull(channelContext); - channelContext = getChannelContext("application", "agent", (Long) params.get(AgentPropertiesType.START_TIMESTAMP.getName()) + 1, ss.getDuplexCommunicationChannelContext()); + channelContext = getChannelContext("application", "agent", (Long) params.get(AgentHandShakePropertyType.START_TIMESTAMP.getName()) + 1, ss.getDuplexCommunicationChannelContext()); Assert.assertNull(channelContext); socket.close(); @@ -178,6 +184,7 @@ public class MessageListenerTest { public void serverMessageListenerTest5() throws InterruptedException { PinpointServerSocket ss = new PinpointServerSocket(); ss.bind("127.0.0.1", 10234); + ss.setMessageListener(new SimpleListener()); PinpointSocketFactory socketFactory = createPinpointSocketFactory(); socketFactory.setMessageListener(SimpleLoggingMessageListener.LISTENER); @@ -189,7 +196,7 @@ public class MessageListenerTest { Thread.sleep(500); List channelContextList = ss.getDuplexCommunicationChannelContext(); - if (channelContextList.size() != 0) { + if (channelContextList.size() != 1) { Assert.fail(); } @@ -224,7 +231,8 @@ public class MessageListenerTest { Assert.fail(); } - if (serverListener.getReceiveEnableWorkerPacketCount() != 4) { + System.out.println(serverListener.getReceiveEnableWorkerPacketCount()); + if (serverListener.getReceiveEnableWorkerPacketCount() < 8) { Assert.fail(); } @@ -297,9 +305,9 @@ public class MessageListenerTest { private final AtomicInteger receiveEnableWorkerPacketCount = new AtomicInteger(); @Override - public int handleEnableWorker(Map properties) { + public HandShakeResponseCode handleHandShake(Map properties) { receiveEnableWorkerPacketCount.incrementAndGet(); - return ControlEnableWorkerConfirmPacket.UNKNOWN_ERROR; + return HandShakeResponseType.Error.UNKNOWN_ERROR; } public int getReceiveEnableWorkerPacketCount() { @@ -327,15 +335,15 @@ public class MessageListenerTest { if (eachContext.getCurrentStateCode() == PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION) { Map agentProperties = eachContext.getChannelProperties(); - if (!applicationName.equals(agentProperties.get(AgentPropertiesType.APPLICATION_NAME.getName()))) { + if (!applicationName.equals(agentProperties.get(AgentHandShakePropertyType.APPLICATION_NAME.getName()))) { continue; } - if (!agentId.equals(agentProperties.get(AgentPropertiesType.AGENT_ID.getName()))) { + if (!agentId.equals(agentProperties.get(AgentHandShakePropertyType.AGENT_ID.getName()))) { continue; } - if (startTimeMillis != (Long) agentProperties.get(AgentPropertiesType.START_TIMESTAMP.getName())) { + if (startTimeMillis != (Long) agentProperties.get(AgentHandShakePropertyType.START_TIMESTAMP.getName())) { continue; } @@ -356,6 +364,16 @@ public class MessageListenerTest { } } + private class SimpleListener extends SimpleLoggingServerMessageListener { + + @Override + public HandShakeResponseCode handleHandShake(Map properties) { + logger.info("handleEnableWorker {}", properties); + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; + + } + + } } 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 cc9c5e905..de1660312 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,7 +7,8 @@ import java.util.Map; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -34,9 +35,9 @@ public class TestSeverMessageListener implements ServerMessageListener { } @Override - public int handleEnableWorker(Map properties) { - logger.debug("handleEnableWorker properties:{} channel:{}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + public HandShakeResponseCode handleHandShake(Map properties) { + logger.debug("handle handShake properties:{} channel:{}", properties); + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } public byte[] getOpen() { diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseDeserializerFactory.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseDeserializerFactory.java index 4ee9b0360..7d126de6e 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseDeserializerFactory.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseDeserializerFactory.java @@ -11,8 +11,6 @@ public final class CommandHeaderTBaseDeserializerFactory implements Deserializer private final DeserializerFactory factory; public CommandHeaderTBaseDeserializerFactory(String version) { - System.out.println(version); - TBaseLocator commandTbaseLocator = new TCommandRegistry(TCommandTypeVersion.getVersion(version)); TProtocolFactory protocolFactory = new TCompactProtocol.Factory(); diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseSerializerFactory.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseSerializerFactory.java index 6dcbbcc42..7d2a79fb7 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseSerializerFactory.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/CommandHeaderTBaseSerializerFactory.java @@ -17,9 +17,6 @@ public final class CommandHeaderTBaseSerializerFactory implements SerializerFact } public CommandHeaderTBaseSerializerFactory(String version, int outputStreamSize) { - System.out.println(version); - - TBaseLocator commandTbaseLocator = new TCommandRegistry(TCommandTypeVersion.getVersion(version)); TProtocolFactory protocolFactory = new TCompactProtocol.Factory(); 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 8c1aad263..d43bbcd14 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 @@ -15,15 +15,14 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.nhn.pinpoint.common.util.NetUtils; -import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseCode; +import com.nhn.pinpoint.rpc.packet.HandShakeResponseType; 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.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; @@ -169,9 +168,9 @@ public class PinpointSocketManager { } @Override - public int handleEnableWorker(Map properties) { - logger.warn("do handleEnableWorker {}", properties); - return ControlEnableWorkerConfirmPacket.SUCCESS; + public HandShakeResponseCode handleHandShake(Map properties) { + logger.warn("do handShake {}", properties); + return HandShakeResponseType.Success.DUPLEX_COMMUNICATION; } } 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 44cd47609..f6d67ff32 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 @@ -77,7 +77,6 @@ public class ClusterTest { ZooKeeper zookeeper = new ZooKeeper("127.0.0.1:22213", 5000, null); getNodeAndCompareContents(zookeeper); - } // ApplicationContext 설정에 맞게 등록이 되는지