Merge branch 'master' of sunsh318/pinpoint-2

from pull-request 189

* refs/heads/master:
  #81. always added agent's identity to channel
This commit is contained in:
koo-taejin
2014-11-25 17:19:57 +09:00
40 changed files with 447 additions and 224 deletions
@@ -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<Object, Object> 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.");
}
@@ -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<Object, Object> 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<Object, Object> 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);
@@ -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<Object, Object> 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;
@@ -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<Object, Object> 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;
}
@@ -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);
@@ -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;
}
}
@@ -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;
}
}
@@ -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;
}
}
@@ -48,6 +48,14 @@
</bean>
<!-- Cluster 관련된 Bean들 -->
<bean id="commandHeaderTBaseSerializerFactory" class="com.nhn.pinpoint.thrift.io.CommandHeaderTBaseSerializerFactory">
<constructor-arg value="#{T(com.nhn.pinpoint.common.Version).VERSION}" />
</bean>
<bean id="commandHeaderTBaseDeserializerFactory" class="com.nhn.pinpoint.thrift.io.CommandHeaderTBaseDeserializerFactory">
<constructor-arg value="#{T(com.nhn.pinpoint.common.Version).VERSION}" />
</bean>
<bean id="clusterPointRouter" class="com.nhn.pinpoint.collector.cluster.ClusterPointRouter">
</bean>
@@ -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) {
@@ -86,14 +86,14 @@ public class AgentInformation {
public Map<String, Object> toMap() {
Map<String, Object> map = new HashMap<String, Object>();
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;
}
@@ -248,17 +248,21 @@ public class DefaultAgent implements Agent {
}
protected PinpointSocketFactory createPinpointSocketFactory(boolean isSupportServerMode) {
Map<String, Object> properties = this.agentInformation.toMap();
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
pinpointSocketFactory.setTimeoutMillis(1000 * 5);
pinpointSocketFactory.setProperties(properties);
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
pinpointSocketFactory.setTimeoutMillis(1000 * 5);
Map<String, Object> 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;
}
@@ -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;
}
}
@@ -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);
@@ -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);
@@ -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;
}
}
@@ -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;
@@ -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<String, Object> properties = this.pinpointSocketFactory.getProperties();
logger.debug("write HandShakePakcet channel:{}, property:{}..", channel, properties);
try {
Map<String, Object> 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 {
@@ -93,8 +93,8 @@ public class ReconnectStateSocketHandler implements SocketHandler {
}
@Override
public void turnOnServerMode() {
throw new UnsupportedOperationException();
public void doHandShake() {
// throw new UnsupportedOperationException();
}
}
@@ -42,6 +42,6 @@ public interface SocketHandler {
boolean isSupportServerMode();
void turnOnServerMode();
void doHandShake();
}
@@ -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:
@@ -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);
}
}
@@ -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;
}
@@ -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;
}
@@ -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;
}
}
@@ -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;
}
}
@@ -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;
@@ -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<Object, Object> properties = decodeSocketProperties((ControlEnableWorkerPacket) message);
Map<Object, Object> 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<Object, Object> decodeSocketProperties(ControlEnableWorkerPacket message) {
private Map<Object, Object> decodeSocketProperties(ControlHandShakePacket message) {
try {
byte[] payload = message.getPayload();
Map<Object, Object> 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<String, Object> result = new HashMap<String, Object>();
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
@@ -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);
}
@@ -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;
}
}
@@ -25,9 +25,26 @@ public class MapUtils {
return (String) value;
}
return null;
return defaultValue;
}
public static Boolean getBoolean(Map<Object, Object> map, String key) {
return getBoolean(map, key, false);
}
public static Boolean getBoolean(Map<Object, Object> 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<Object, Object> 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<Object, Object> map, String key) {
@@ -60,7 +77,7 @@ public class MapUtils {
return (Long) value;
}
return null;
return defaultValue;
}
}
@@ -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;
}
@@ -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<Object, Object> result = (Map<Object, Object>) 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;
}
}
@@ -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<String, Object> properties) throws ProtocolException, IOException {
sendRegisterPacket(socket.getOutputStream(), properties);
ControlEnableWorkerConfirmPacket packet = receiveRegisterConfirmPacket(socket.getInputStream());
ControlHandShakeResponsePacket packet = receiveRegisterConfirmPacket(socket.getInputStream());
Map<Object, Object> result = (Map<Object, Object>) ControlMessageEnDeconderUtils.decode(packet.getPayload());
return MapUtils.getInteger(result, "code", -1);
@@ -77,7 +79,7 @@ public class EventListnerTest {
private void sendRegisterPacket(OutputStream outputStream, Map<String, Object> 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;
}
}
@@ -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<ChannelContext> 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<ChannelContext> 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;
}
}
}
@@ -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() {
@@ -11,8 +11,6 @@ public final class CommandHeaderTBaseDeserializerFactory implements Deserializer
private final DeserializerFactory<HeaderTBaseDeserializer> factory;
public CommandHeaderTBaseDeserializerFactory(String version) {
System.out.println(version);
TBaseLocator commandTbaseLocator = new TCommandRegistry(TCommandTypeVersion.getVersion(version));
TProtocolFactory protocolFactory = new TCompactProtocol.Factory();
@@ -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();
@@ -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;
}
}
@@ -77,7 +77,6 @@ public class ClusterTest {
ZooKeeper zookeeper = new ZooKeeper("127.0.0.1:22213", 5000, null);
getNodeAndCompareContents(zookeeper);
}
// ApplicationContext 설정에 맞게 등록이 되는지