diff --git a/src/main/java/com/nhn/pinpoint/rpc/RequestResponseServerMessageListener.java b/src/main/java/com/nhn/pinpoint/rpc/RequestResponseServerMessageListener.java index d1e658450..5766a649c 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/RequestResponseServerMessageListener.java +++ b/src/main/java/com/nhn/pinpoint/rpc/RequestResponseServerMessageListener.java @@ -1,11 +1,15 @@ package com.nhn.pinpoint.rpc; +import java.util.Map; + +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.StreamPacket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -33,8 +37,13 @@ public class RequestResponseServerMessageListener implements ServerMessageListen @Override public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { - logger.info("handlerStream {}", streamChannel, streamChannel); + logger.info("handlerStream {} {}", streamChannel, streamChannel); } + @Override + public int handleEnableWorker(Map properties) { + logger.info("handleEnableWorker {}", properties); + return ControlEnableWorkerConfirmPacket.SUCCESS; + } } diff --git a/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketFactory.java b/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketFactory.java index 22368c7d4..a4d1fdd61 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketFactory.java +++ b/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketFactory.java @@ -45,12 +45,12 @@ public class PinpointSocketFactory { private static final int DEFAULT_CONNECT_TIMEOUT = 5000; private static final long DEFAULT_TIMEOUTMILLIS = 3 * 1000; private static final long DEFAULT_PING_DELAY = 60 * 1000 * 5; - private static final long DEFAULT_REGISTER_AGENT_PACKET_DELAY = 60 * 1000 * 1; + private static final long DEFAULT_ENABLE_WORKER_PACKET_DELAY = 60 * 1000 * 1; private volatile boolean released; private ClientBootstrap bootstrap; - private Map agentProperties = Collections.EMPTY_MAP; + private Map properties = Collections.EMPTY_MAP; private long reconnectDelay = 3 * 1000; private final Timer timer; @@ -58,7 +58,7 @@ public class PinpointSocketFactory { // 이 값이 짧아야 될 필요가 없음. client에서 server로 가는 핑 주기를 짧게 유지한다고 해서. // 연결끊김이 빨랑 디텍트 되는게 아님. 오히려 server에서 client의 ping주기를 짧게 해야 디텍트 속도가 빨라짐. private long pingDelay = DEFAULT_PING_DELAY; - private long registerAgentPacketDelay = DEFAULT_REGISTER_AGENT_PACKET_DELAY; + private long enableWorkerPacketDelay = DEFAULT_ENABLE_WORKER_PACKET_DELAY; private long timeoutMillis = DEFAULT_TIMEOUTMILLIS; @@ -142,15 +142,15 @@ public class PinpointSocketFactory { this.pingDelay = pingDelay; } - public long getRegisterAgentPacketDelay() { - return registerAgentPacketDelay; + public long getEnableWorkerPacketDelay() { + return enableWorkerPacketDelay; } - public void setRegisterAgentPacketDelay(long registerAgentPacketDelay) { - if (registerAgentPacketDelay < 0) { - throw new IllegalArgumentException("registerAgentPacketDelay cannot be a negative number"); + public void setEnableWorkerPacketDelay(long enableWorkerPacketDelay) { + if (enableWorkerPacketDelay < 0) { + throw new IllegalArgumentException("EnableWorkerPacketDelay cannot be a negative number"); } - this.registerAgentPacketDelay = registerAgentPacketDelay; + this.enableWorkerPacketDelay = enableWorkerPacketDelay; } public long getTimeoutMillis() { @@ -368,21 +368,21 @@ public class PinpointSocketFactory { // stop 뭔가 취소를 해야 되나?? } - public Map getAgentProperties() { - return agentProperties; + public Map getProperties() { + return properties; } - public void setAgentProperties(Map agentProperties) { + public void setProperties(Map agentProperties) { if (agentProperties == null) { return; } - if (this.agentProperties != Collections.EMPTY_MAP) { + if (this.properties != Collections.EMPTY_MAP) { logger.warn("Properties variable alreay registered."); return; } - this.agentProperties = Collections.unmodifiableMap(CopyUtils.mediumCopyMap(agentProperties)); + this.properties = Collections.unmodifiableMap(CopyUtils.mediumCopyMap(agentProperties)); } } diff --git a/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketHandler.java b/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketHandler.java index 191b999b9..48a8da45b 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketHandler.java +++ b/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocketHandler.java @@ -27,8 +27,8 @@ 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.ControlRegisterAgentConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlRegisterAgentPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; import com.nhn.pinpoint.rpc.packet.Packet; import com.nhn.pinpoint.rpc.packet.PacketType; import com.nhn.pinpoint.rpc.packet.PingPacket; @@ -49,7 +49,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke private static final long DEFAULT_PING_DELAY = 60 * 1000 * 5; private static final long DEFAULT_TIMEOUTMILLIS = 3 * 1000; - private static final long DEFAULT_REGISTER_AGENT_PACKET_DELAY = 60 * 1000 * 1; + private static final long DEFAULT_ENABLE_WORKER_PACKET_DELAY = 60 * 1000 * 1; private final Logger logger = LoggerFactory.getLogger(this.getClass()); @@ -60,7 +60,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke private long timeoutMillis = DEFAULT_TIMEOUTMILLIS; private long pingDelay = DEFAULT_PING_DELAY; - private long registerAgentPacketDelay = DEFAULT_REGISTER_AGENT_PACKET_DELAY; + private long enableWorkerPacketDelay = DEFAULT_ENABLE_WORKER_PACKET_DELAY; private final Timer channelTimer; @@ -75,10 +75,10 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke private final ChannelFutureListener sendWriteFailFutureListener = new WriteFailFutureListener(this.logger, "send() write fail.", "send() write fail."); public PinpointSocketHandler(PinpointSocketFactory pinpointSocketFactory, Map agentProperties) { - this(pinpointSocketFactory, DEFAULT_PING_DELAY, DEFAULT_REGISTER_AGENT_PACKET_DELAY, DEFAULT_TIMEOUTMILLIS); + this(pinpointSocketFactory, DEFAULT_PING_DELAY, DEFAULT_ENABLE_WORKER_PACKET_DELAY, DEFAULT_TIMEOUTMILLIS); } - public PinpointSocketHandler(PinpointSocketFactory pinpointSocketFactory, long pingDelay, long registerAgentPacketDelay, long timeoutMillis) { + public PinpointSocketHandler(PinpointSocketFactory pinpointSocketFactory, long pingDelay, long enableWorkerPacketDelay, long timeoutMillis) { if (pinpointSocketFactory == null) { throw new NullPointerException("pinpointSocketFactory must not be null"); } @@ -90,7 +90,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke this.requestManager = new RequestManager(timer); this.streamChannelManager = new StreamChannelManager(); this.pingDelay = pingDelay; - this.registerAgentPacketDelay = registerAgentPacketDelay; + this.enableWorkerPacketDelay = enableWorkerPacketDelay; this.timeoutMillis = timeoutMillis; } @@ -140,8 +140,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke this.messageListener = messageListener; // MessageHandler가 걸릴 경우 Register Agent Packet 전달 - sendRegisterAgentPacket(); - registerRegisterAgentPacketTask(); + sendEnableWorkerPacket(); + reservationEnableWorkerPacketJob(); } @Override @@ -192,12 +192,12 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke write.addListener(pingWriteFailFutureListener); } - private class RegisterAgentPacketTask implements TimerTask { + private class RegisterEnableWorkerPacketJob implements TimerTask { @Override public void run(Timeout timeout) throws Exception { if (timeout.isCancelled()) { - newRegisterAgentPacketTimeout(this); + reservationEnableWorkerPacketJob(this); return; } if (isClosed()) { @@ -205,34 +205,34 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke } if (state.getState() == State.RUN_WITHOUT_REGISTER) { - sendRegisterAgentPacket(); - newRegisterAgentPacketTimeout(this); + sendEnableWorkerPacket(); + reservationEnableWorkerPacketJob(this); } } } - private void registerRegisterAgentPacketTask() { - final RegisterAgentPacketTask task = new RegisterAgentPacketTask(); - newRegisterAgentPacketTimeout(task); + private void reservationEnableWorkerPacketJob() { + final RegisterEnableWorkerPacketJob job = new RegisterEnableWorkerPacketJob(); + reservationEnableWorkerPacketJob(job); } - private void newRegisterAgentPacketTimeout(RegisterAgentPacketTask task) { - this.channelTimer.newTimeout(task, registerAgentPacketDelay, TimeUnit.MILLISECONDS); + private void reservationEnableWorkerPacketJob(RegisterEnableWorkerPacketJob task) { + this.channelTimer.newTimeout(task, enableWorkerPacketDelay, TimeUnit.MILLISECONDS); } - void sendRegisterAgentPacket() { + void sendEnableWorkerPacket() { if (!isRun()) { return; } - logger.debug("write RegisterAgentPacket {}", channel); + logger.debug("write EnableWorkerPacket {}", channel); byte[] payload; try { - Map properties = this.pinpointSocketFactory.getAgentProperties(); + Map properties = this.pinpointSocketFactory.getProperties(); payload = ControlMessageEnDeconderUtils.encode(properties); - ControlRegisterAgentPacket packet = new ControlRegisterAgentPacket(payload); + ControlEnableWorkerPacket packet = new ControlEnableWorkerPacket(payload); ChannelFuture write = this.channel.write(packet); } catch (ProtocolException e) { logger.warn(e.getMessage(), e); @@ -380,8 +380,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke case PacketType.CONTROL_SERVER_CLOSE: messageReceivedServerClosed(e.getChannel()); return; - case PacketType.CONTROL_REGISTER_AGENT_CONFIRM: - messageReceivedRegisterAgentConfirm((ControlRegisterAgentConfirmPacket)message, e.getChannel()); + case PacketType.CONTROL_ENABLE_WORKER_CONFIRM: + messageReceivedEnableWorkerConfirm((ControlEnableWorkerConfirmPacket)message, e.getChannel()); return; default: logger.warn("unexpectedMessage received:{} address:{}", message, e.getRemoteAddress()); @@ -397,13 +397,13 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke state.setState(State.RECONNECT); } - private void messageReceivedRegisterAgentConfirm(ControlRegisterAgentConfirmPacket message, Channel channel) { + private void messageReceivedEnableWorkerConfirm(ControlEnableWorkerConfirmPacket message, Channel channel) { int code = getRegisterAgnetConfirmPacketCode(message.getPayload()); - logger.info("RegisterAgentConfirm Packet({}) code={} received. {}", message, code, channel); + logger.info("EnableWorkerConfirm Packet({}) code={} received. {}", message, code, channel); // reconnect 상태로 변경한다. - if (code == ControlRegisterAgentConfirmPacket.SUCCESS || code == ControlRegisterAgentConfirmPacket.ALREADY_REGISTER) { + if (code == ControlEnableWorkerConfirmPacket.SUCCESS || code == ControlEnableWorkerConfirmPacket.ALREADY_REGISTER) { state.changeRun(); } } diff --git a/src/main/java/com/nhn/pinpoint/rpc/client/SocketClientPipelineFactory.java b/src/main/java/com/nhn/pinpoint/rpc/client/SocketClientPipelineFactory.java index 1bd7b013c..f21e331d5 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/client/SocketClientPipelineFactory.java +++ b/src/main/java/com/nhn/pinpoint/rpc/client/SocketClientPipelineFactory.java @@ -32,9 +32,9 @@ public class SocketClientPipelineFactory implements ChannelPipelineFactory { pipeline.addLast("encoder", new PacketEncoder()); pipeline.addLast("decoder", new PacketDecoder()); long pingDelay = pinpointSocketFactory.getPingDelay(); - long registerAgentPacketDelay = pinpointSocketFactory.getRegisterAgentPacketDelay(); + long enableWorkerPacketDelay = pinpointSocketFactory.getEnableWorkerPacketDelay(); long timeoutMillis = pinpointSocketFactory.getTimeoutMillis(); - PinpointSocketHandler pinpointSocketHandler = new PinpointSocketHandler(pinpointSocketFactory, pingDelay, registerAgentPacketDelay, timeoutMillis); + PinpointSocketHandler pinpointSocketHandler = new PinpointSocketHandler(pinpointSocketFactory, pingDelay, enableWorkerPacketDelay, timeoutMillis); pipeline.addLast("writeTimeout", new WriteTimeoutHandler(pinpointSocketHandler.getChannelTimer(), 3000, TimeUnit.MILLISECONDS)); pipeline.addLast("socketHandler", pinpointSocketHandler); return pipeline; diff --git a/src/main/java/com/nhn/pinpoint/rpc/codec/PacketDecoder.java b/src/main/java/com/nhn/pinpoint/rpc/codec/PacketDecoder.java index 0bca8f7c3..acba4fdd1 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/codec/PacketDecoder.java +++ b/src/main/java/com/nhn/pinpoint/rpc/codec/PacketDecoder.java @@ -58,11 +58,10 @@ public class PacketDecoder extends FrameDecoder { readPong(packetType, buffer); // pong 도 그냥 버리자. return null; - case PacketType.CONTROL_REGISTER_AGENT: - return readRegisterAgent(packetType, buffer); - case PacketType.CONTROL_REGISTER_AGENT_CONFIRM: - return readRegisterAgentConfirm(packetType, buffer); - + case PacketType.CONTROL_ENABLE_WORKER: + return readEnableWorker(packetType, buffer); + case PacketType.CONTROL_ENABLE_WORKER_CONFIRM: + return readEnableWorkerConfirm(packetType, buffer); } logger.error("invalid packetType received. packetType:{}, channel:{}", packetType, channel); channel.close(); @@ -130,12 +129,12 @@ public class PacketDecoder extends FrameDecoder { return StreamClosePacket.readBuffer(packetType, buffer); } - private Object readRegisterAgent(short packetType, ChannelBuffer buffer) { - return ControlRegisterAgentPacket.readBuffer(packetType, buffer); + private Object readEnableWorker(short packetType, ChannelBuffer buffer) { + return ControlEnableWorkerPacket.readBuffer(packetType, buffer); } - private Object readRegisterAgentConfirm(short packetType, ChannelBuffer buffer) { - return ControlRegisterAgentConfirmPacket.readBuffer(packetType, buffer); + private Object readEnableWorkerConfirm(short packetType, ChannelBuffer buffer) { + return ControlEnableWorkerConfirmPacket.readBuffer(packetType, buffer); } } diff --git a/src/main/java/com/nhn/pinpoint/rpc/packet/ControlRegisterAgentConfirmPacket.java b/src/main/java/com/nhn/pinpoint/rpc/packet/ControlEnableWorkerConfirmPacket.java similarity index 65% rename from src/main/java/com/nhn/pinpoint/rpc/packet/ControlRegisterAgentConfirmPacket.java rename to src/main/java/com/nhn/pinpoint/rpc/packet/ControlEnableWorkerConfirmPacket.java index 8288bd4e0..b4f059ee1 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/packet/ControlRegisterAgentConfirmPacket.java +++ b/src/main/java/com/nhn/pinpoint/rpc/packet/ControlEnableWorkerConfirmPacket.java @@ -6,39 +6,40 @@ import org.jboss.netty.buffer.ChannelBuffers; /** * @author koo.taejin */ -public class ControlRegisterAgentConfirmPacket extends ControlPacket { +public class ControlEnableWorkerConfirmPacket 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 ControlRegisterAgentConfirmPacket(byte[] payload) { + public ControlEnableWorkerConfirmPacket(byte[] payload) { super(payload); } - public ControlRegisterAgentConfirmPacket(int requestId, byte[] payload) { + public ControlEnableWorkerConfirmPacket(int requestId, byte[] payload) { super(payload); setRequestId(requestId); } @Override public short getPacketType() { - return PacketType.CONTROL_REGISTER_AGENT_CONFIRM; + return PacketType.CONTROL_ENABLE_WORKER_CONFIRM; } @Override public ChannelBuffer toBuffer() { ChannelBuffer header = ChannelBuffers.buffer(2 + 4 + 4); - header.writeShort(PacketType.CONTROL_REGISTER_AGENT_CONFIRM); + header.writeShort(PacketType.CONTROL_ENABLE_WORKER_CONFIRM); header.writeInt(getRequestId()); return PayloadPacket.appendPayload(header, payload); } - public static ControlRegisterAgentConfirmPacket readBuffer(short packetType, ChannelBuffer buffer) { - assert packetType == PacketType.CONTROL_REGISTER_AGENT_CONFIRM; + public static ControlEnableWorkerConfirmPacket readBuffer(short packetType, ChannelBuffer buffer) { + assert packetType == PacketType.CONTROL_ENABLE_WORKER_CONFIRM; if (buffer.readableBytes() < 8) { buffer.resetReaderIndex(); @@ -50,7 +51,7 @@ public class ControlRegisterAgentConfirmPacket extends ControlPacket { if (payload == null) { return null; } - final ControlRegisterAgentConfirmPacket helloPacket = new ControlRegisterAgentConfirmPacket(payload.array()); + final ControlEnableWorkerConfirmPacket helloPacket = new ControlEnableWorkerConfirmPacket(payload.array()); helloPacket.setRequestId(messageId); return helloPacket; } diff --git a/src/main/java/com/nhn/pinpoint/rpc/packet/ControlRegisterAgentPacket.java b/src/main/java/com/nhn/pinpoint/rpc/packet/ControlEnableWorkerPacket.java similarity index 65% rename from src/main/java/com/nhn/pinpoint/rpc/packet/ControlRegisterAgentPacket.java rename to src/main/java/com/nhn/pinpoint/rpc/packet/ControlEnableWorkerPacket.java index 4bfe41e2c..a177fdfb1 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/packet/ControlRegisterAgentPacket.java +++ b/src/main/java/com/nhn/pinpoint/rpc/packet/ControlEnableWorkerPacket.java @@ -6,34 +6,34 @@ import org.jboss.netty.buffer.ChannelBuffers; /** * @author koo.taejin */ -public class ControlRegisterAgentPacket extends ControlPacket { +public class ControlEnableWorkerPacket extends ControlPacket { - public ControlRegisterAgentPacket(byte[] payload) { + public ControlEnableWorkerPacket(byte[] payload) { super(payload); } - public ControlRegisterAgentPacket(int requestId, byte[] payload) { + public ControlEnableWorkerPacket(int requestId, byte[] payload) { super(payload); setRequestId(requestId); } @Override public short getPacketType() { - return PacketType.CONTROL_REGISTER_AGENT; + return PacketType.CONTROL_ENABLE_WORKER; } @Override public ChannelBuffer toBuffer() { ChannelBuffer header = ChannelBuffers.buffer(2 + 4 + 4); - header.writeShort(PacketType.CONTROL_REGISTER_AGENT); + header.writeShort(PacketType.CONTROL_ENABLE_WORKER); header.writeInt(getRequestId()); return PayloadPacket.appendPayload(header, payload); } - public static ControlRegisterAgentPacket readBuffer(short packetType, ChannelBuffer buffer) { - assert packetType == PacketType.CONTROL_REGISTER_AGENT; + public static ControlEnableWorkerPacket readBuffer(short packetType, ChannelBuffer buffer) { + assert packetType == PacketType.CONTROL_ENABLE_WORKER; if (buffer.readableBytes() < 8) { buffer.resetReaderIndex(); @@ -45,7 +45,7 @@ public class ControlRegisterAgentPacket extends ControlPacket { if (payload == null) { return null; } - final ControlRegisterAgentPacket helloPacket = new ControlRegisterAgentPacket(payload.array()); + final ControlEnableWorkerPacket helloPacket = new ControlEnableWorkerPacket(payload.array()); helloPacket.setRequestId(messageId); return helloPacket; } diff --git a/src/main/java/com/nhn/pinpoint/rpc/packet/PacketType.java b/src/main/java/com/nhn/pinpoint/rpc/packet/PacketType.java index 7b6045df9..783c17fb5 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/packet/PacketType.java +++ b/src/main/java/com/nhn/pinpoint/rpc/packet/PacketType.java @@ -26,8 +26,8 @@ public class PacketType { public static final short CONTROL_SERVER_CLOSE = 110; // 컨트롤 패킷 - public static final short CONTROL_REGISTER_AGENT = 150; - public static final short CONTROL_REGISTER_AGENT_CONFIRM = 151; + public static final short CONTROL_ENABLE_WORKER = 150; + public static final short CONTROL_ENABLE_WORKER_CONFIRM = 151; // ping, pong의 경우 성능상 두고 다른 CONTROL은 이걸로 뺌 public static final short CONTROL_PING = 200; diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java b/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java index e13d7a7a4..b52b64431 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java @@ -21,7 +21,7 @@ public class ChannelContext { private final SocketChannelStateChangeEventListener stateChangeEventListener; - private volatile Map agentProperties = Collections.EMPTY_MAP; + private volatile Map channelProperties = Collections.EMPTY_MAP; public ChannelContext(SocketChannel socketChannel, ServerStreamChannelManager streamChannelManager) { this(socketChannel, streamChannelManager, DoNothingChannelStateEventListener.INSTANCE); @@ -63,10 +63,10 @@ public class ChannelContext { } } - public void changeStateRunWithoutRegister() { - logger.debug("Channel({}) state will be changed {}.", socketChannel, PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER); - if (state.changeStateRunWithoutRegister()) { - stateChangeEventListener.eventPerformed(this, PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER); + public void changeStateRunDuplexCommunication() { + logger.debug("Channel({}) state will be changed {}.", socketChannel, PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION); + if (state.changeStateRunDuplexCommunication()) { + stateChangeEventListener.eventPerformed(this, PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION); } } @@ -99,26 +99,26 @@ public class ChannelContext { } public String getVersion() { - return MapUtils.get(agentProperties, AgentPropertiesType.VERSION.getName(), String.class, "UNKNOWN"); + return MapUtils.get(channelProperties, AgentPropertiesType.VERSION.getName(), String.class, "UNKNOWN"); } - public Map getAgentProperties() { - return agentProperties; + public Map getChannelProperties() { + return channelProperties; } - public boolean setAgentProperties(Map agentProperties) { - if (agentProperties == null) { + public boolean setChannelProperties(Map properties) { + if (properties == null) { return false; } - synchronized (agentProperties) { - if (this.agentProperties == Collections.EMPTY_MAP) { - this.agentProperties = Collections.unmodifiableMap(CopyUtils.mediumCopyMap(agentProperties)); + synchronized (properties) { + if (this.channelProperties == Collections.EMPTY_MAP) { + this.channelProperties = Collections.unmodifiableMap(CopyUtils.mediumCopyMap(properties)); return true; } } - logger.warn("Already Register AgentProperties.({}).", this.agentProperties); + logger.warn("Already Register ChannelProperties.({}).", this.channelProperties); return false; } diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java index 45234db9b..f19a14407 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java @@ -3,8 +3,8 @@ package com.nhn.pinpoint.rpc.server; import java.net.InetAddress; import java.net.InetSocketAddress; import java.util.ArrayList; +import java.util.Collections; import java.util.HashMap; -import java.util.Iterator; import java.util.List; import java.util.Map; import java.util.concurrent.ExecutorService; @@ -39,8 +39,8 @@ 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.ControlRegisterAgentConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlRegisterAgentPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; import com.nhn.pinpoint.rpc.packet.Packet; import com.nhn.pinpoint.rpc.packet.PacketType; import com.nhn.pinpoint.rpc.packet.PingPacket; @@ -211,8 +211,28 @@ public class PinpointServerSocket extends SimpleChannelHandler { case PacketType.APPLICATION_STREAM_RESPONSE: handleStreamPacket((StreamPacket) message, channel); return; - case PacketType.CONTROL_REGISTER_AGENT: - handleRegisterAgent((ControlRegisterAgentPacket) message, channel); + case PacketType.CONTROL_ENABLE_WORKER: + int requestId = ((ControlEnableWorkerPacket)message).getRequestId(); + + Map properties = decodeSocketProperties((ControlEnableWorkerPacket) message); + if (properties == null) { + sendEnableWorkerConfirmMessage(requestId, ControlEnableWorkerConfirmPacket.ILLEGAL_PROTOCOL, channel); + return; + } + + ChannelContext channelContext = getChannelContext(channel); + channelContext.setChannelProperties(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); + } + } else { + sendEnableWorkerConfirmMessage(requestId, returnCode, channel); + } return; case PacketType.CONTROL_CLIENT_CLOSE: { closeChannel(channel); @@ -260,55 +280,48 @@ public class PinpointServerSocket extends SimpleChannelHandler { logger.warn("invalid streamPacket. channel:{}", channel); } } - - private void handleRegisterAgent(ControlRegisterAgentPacket message, Channel channel) { - ChannelContext context = getChannelContext(channel); - - int code = registerAgent(context, message); - - if (code == ControlRegisterAgentConfirmPacket.SUCCESS) { - context.changeStateRun(); - } - - try { - Map result = new HashMap(); - result.put("code", code); - - byte[] resultPayload = ControlMessageEnDeconderUtils.encode(result); - ControlRegisterAgentConfirmPacket packet = new ControlRegisterAgentConfirmPacket(message.getRequestId(), resultPayload); - - channel.write(packet); - } catch (ProtocolException e) { - logger.warn(e.getMessage(), e); - } - } - private int registerAgent(ChannelContext context, ControlRegisterAgentPacket message) { + private Map decodeSocketProperties(ControlEnableWorkerPacket message) { Map properties = null; try { byte[] payload = message.getPayload(); properties = (Map) ControlMessageEnDeconderUtils.decode(payload); + return properties; } catch (ProtocolException e) { logger.warn(e.getMessage(), e); } - if (properties == null) { - return ControlRegisterAgentConfirmPacket.ILLEGAL_PROTOCOL; - } - - boolean hasAllType = AgentPropertiesType.hasAllType(properties); - if (!hasAllType) { - return ControlRegisterAgentConfirmPacket.INVALID_PROPERTIES; - } - - boolean isSuccess = context.setAgentProperties(properties); - if (isSuccess) { - return ControlRegisterAgentConfirmPacket.SUCCESS; - } else { - return ControlRegisterAgentConfirmPacket.ALREADY_REGISTER; - } + return null; } + private boolean changeStateToRunDuplexCommunication(int returnCode, Channel channel) { + ChannelContext context = getChannelContext(channel); + + if (returnCode == ControlEnableWorkerConfirmPacket.SUCCESS) { + if (context.getCurrentStateCode() != PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION) { + context.changeStateRunDuplexCommunication(); + return true; + } + } + + return false; + } + + private void sendEnableWorkerConfirmMessage(int requestId, int returnCode, Channel channel) { + try { + Map result = new HashMap(); + result.put("code", returnCode); + + byte[] resultPayload = ControlMessageEnDeconderUtils.encode(result); + ControlEnableWorkerConfirmPacket packet = new ControlEnableWorkerConfirmPacket(requestId, resultPayload); + + channel.write(packet); + } catch (ProtocolException e) { + logger.warn(e.getMessage(), e); + } + + } + @Override public void channelOpen(ChannelHandlerContext ctx, ChannelStateEvent e) throws Exception { final Channel channel = e.getChannel(); @@ -337,7 +350,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { prepareChannel(channel); ChannelContext channelContext = getChannelContext(channel); - channelContext.changeStateRunWithoutRegister(); + channelContext.changeStateRun(); super.channelConnected(ctx, e); } @@ -529,13 +542,13 @@ public class PinpointServerSocket extends SimpleChannelHandler { logger.info("sendServerClosedPacket end"); } - public List getRegisterAgentChannelContext() { + public List getDuplexCommunicationChannelContext() { List channelContextList = new ArrayList(); for (Channel channel : channelGroup) { ChannelContext context = getChannelContext(channel); - if (context.getCurrentStateCode() == PinpointServerSocketStateCode.RUN) { + if (context.getCurrentStateCode() == PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION) { channelContextList.add(context); } } @@ -543,7 +556,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { return channelContextList; } - public ChannelContext getRegisterAgentChannelContext(String applicationName, String agentId, long startTimeMillis) { + public ChannelContext getDuplexChannelContext(String applicationName, String agentId, long startTimeMillis) { if (applicationName == null) { return null; } @@ -561,8 +574,8 @@ public class PinpointServerSocket extends SimpleChannelHandler { for (Channel channel : channelGroup) { ChannelContext context = getChannelContext(channel); - if (context.getCurrentStateCode() == PinpointServerSocketStateCode.RUN) { - Map agentProperties = context.getAgentProperties(); + if (context.getCurrentStateCode() == PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION) { + Map agentProperties = context.getChannelProperties(); if (!applicationName.equals(agentProperties.get(AgentPropertiesType.APPLICATION_NAME.getName()))) { continue; diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java index ac40eb417..4def32d61 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java @@ -51,11 +51,11 @@ public class PinpointServerSocketState { public boolean changeStateRun() { return setSessionState(PinpointServerSocketStateCode.RUN); } - - public boolean changeStateRunWithoutRegister() { - return setSessionState(PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER); - } + public boolean changeStateRunDuplexCommunication() { + return setSessionState(PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION); + } + public boolean changeStateBeingShutdown() { return setSessionState(PinpointServerSocketStateCode.BEING_SHUTDOWN); } diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java index aeceebe15..da92b31e5 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java @@ -17,16 +17,16 @@ public enum PinpointServerSocketStateCode { // UNEXPECTED_SHUTDOWN : CLOSE 등의 명령을 받지 못한 상태에서 상대방이 연결을 종료하였을떄 NONE(), - RUN_WITHOUT_REGISTER(NONE), - RUN(NONE, RUN_WITHOUT_REGISTER), - BEING_SHUTDOWN(RUN, RUN_WITHOUT_REGISTER), - SHUTDOWN(RUN, RUN_WITHOUT_REGISTER, BEING_SHUTDOWN), - UNEXPECTED_SHUTDOWN(RUN, RUN_WITHOUT_REGISTER), + RUN(NONE), //Simplex Communication + RUN_DUPLEX_COMMUNICATION(NONE, RUN), + BEING_SHUTDOWN(RUN_DUPLEX_COMMUNICATION, RUN), + SHUTDOWN(RUN_DUPLEX_COMMUNICATION, RUN, BEING_SHUTDOWN), + UNEXPECTED_SHUTDOWN(RUN_DUPLEX_COMMUNICATION, RUN), // 서버쪽에서 먼저 연결을 끊자는 메시지도 필요하다. // 예를 들어 HELLO 이후 다 확인했는데, 같은 Agent명이 있으면(?) 이걸 사용자에게 말해야 할까? 아닐까? 알림 등 - ERROR_UNKOWN(RUN, RUN_WITHOUT_REGISTER), - ERROR_ILLEGAL_STATE_CHANGE(NONE, RUN, RUN_WITHOUT_REGISTER, BEING_SHUTDOWN); + ERROR_UNKOWN(RUN_DUPLEX_COMMUNICATION, RUN), + ERROR_ILLEGAL_STATE_CHANGE(NONE, RUN_DUPLEX_COMMUNICATION, RUN, BEING_SHUTDOWN); private final Set validBeforeStateSet; @@ -55,7 +55,7 @@ public enum PinpointServerSocketStateCode { } public static boolean isRun(PinpointServerSocketStateCode code) { - if (code == RUN || code == RUN_WITHOUT_REGISTER) { + if (code == RUN_DUPLEX_COMMUNICATION || code == RUN) { return true; } diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/ServerMessageListener.java b/src/main/java/com/nhn/pinpoint/rpc/server/ServerMessageListener.java index e912acc96..82d693ca7 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/ServerMessageListener.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/ServerMessageListener.java @@ -1,5 +1,7 @@ package com.nhn.pinpoint.rpc.server; +import java.util.Map; + import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.StreamPacket; @@ -14,5 +16,7 @@ public interface ServerMessageListener { void handleRequest(RequestPacket requestPacket, SocketChannel channel); void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel); + + int handleEnableWorker(Map properties); } diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java b/src/main/java/com/nhn/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java index f6f42d7bc..42d0b72c3 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java @@ -1,8 +1,12 @@ package com.nhn.pinpoint.rpc.server; +import java.util.Map; + +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.StreamPacket; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -30,6 +34,11 @@ public class SimpleLoggingServerMessageListener implements ServerMessageListener public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { logger.info("handlerStream {} {}", streamChannel, streamChannel); } - + + @Override + public int handleEnableWorker(Map properties) { + logger.info("handleEnableWorker {}", properties); + return ControlEnableWorkerConfirmPacket.SUCCESS; + } } diff --git a/src/test/java/com/nhn/pinpoint/rpc/server/ControlPacketServerTest.java b/src/test/java/com/nhn/pinpoint/rpc/server/ControlPacketServerTest.java index df143325e..99e2d4fe0 100644 --- a/src/test/java/com/nhn/pinpoint/rpc/server/ControlPacketServerTest.java +++ b/src/test/java/com/nhn/pinpoint/rpc/server/ControlPacketServerTest.java @@ -17,8 +17,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.nhn.pinpoint.rpc.control.ProtocolException; -import com.nhn.pinpoint.rpc.packet.ControlRegisterAgentConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlRegisterAgentPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.ResponsePacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -157,7 +157,7 @@ public class ControlPacketServerTest { private int sendAndReceiveRegisterPacket(Socket socket, Map properties) throws ProtocolException, IOException { sendRegisterPacket(socket.getOutputStream(), properties); - ControlRegisterAgentConfirmPacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); + ControlEnableWorkerConfirmPacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); Map result = (Map) ControlMessageEnDeconderUtils.decode(packet.getPayload()); return MapUtils.get(result, "code", Integer.class, -1); @@ -171,7 +171,7 @@ public class ControlPacketServerTest { private void sendRegisterPacket(OutputStream outputStream, Map properties) throws ProtocolException, IOException { byte[] payload = ControlMessageEnDeconderUtils.encode(properties); - ControlRegisterAgentPacket packet = new ControlRegisterAgentPacket(1, payload); + ControlEnableWorkerPacket packet = new ControlEnableWorkerPacket(1, payload); ByteBuffer bb = packet.toBuffer().toByteBuffer(0, packet.toBuffer().writerIndex()); sendData(outputStream, bb.array()); @@ -190,14 +190,14 @@ public class ControlPacketServerTest { outputStream.flush(); } - private ControlRegisterAgentConfirmPacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { + private ControlEnableWorkerConfirmPacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { byte[] payload = readData(inputStream); ChannelBuffer cb = ChannelBuffers.wrappedBuffer(payload); short packetType = cb.readShort(); - ControlRegisterAgentConfirmPacket packet = ControlRegisterAgentConfirmPacket.readBuffer(packetType, cb); + ControlEnableWorkerConfirmPacket packet = ControlEnableWorkerConfirmPacket.readBuffer(packetType, cb); return packet; } @@ -243,14 +243,29 @@ public class ControlPacketServerTest { @Override public void handleRequest(RequestPacket requestPacket, SocketChannel channel) { - logger.info("handlerRequest {}", requestPacket, channel); + logger.info("handlerRequest {} {}", requestPacket, channel); channel.sendResponseMessage(requestPacket, requestPacket.getPayload()); } @Override public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { + logger.info("handleStream {} {}", streamPacket, streamChannel); } + + @Override + public int handleEnableWorker(Map properties) { + if (properties == null) { + return ControlEnableWorkerConfirmPacket.ILLEGAL_PROTOCOL; + } + + boolean hasAllType = AgentPropertiesType.hasAllType(properties); + if (!hasAllType) { + return ControlEnableWorkerConfirmPacket.INVALID_PROPERTIES; + } + + return ControlEnableWorkerConfirmPacket.SUCCESS; + } } private Map getParams() { diff --git a/src/test/java/com/nhn/pinpoint/rpc/server/EventListnerTest.java b/src/test/java/com/nhn/pinpoint/rpc/server/EventListnerTest.java index 96935e72a..05e1638cb 100644 --- a/src/test/java/com/nhn/pinpoint/rpc/server/EventListnerTest.java +++ b/src/test/java/com/nhn/pinpoint/rpc/server/EventListnerTest.java @@ -17,8 +17,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; import com.nhn.pinpoint.rpc.control.ProtocolException; -import com.nhn.pinpoint.rpc.packet.ControlRegisterAgentConfirmPacket; -import com.nhn.pinpoint.rpc.packet.ControlRegisterAgentPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerConfirmPacket; +import com.nhn.pinpoint.rpc.packet.ControlEnableWorkerPacket; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.ResponsePacket; import com.nhn.pinpoint.rpc.packet.SendPacket; @@ -46,10 +46,10 @@ public class EventListnerTest { try { socket = new Socket("127.0.0.1", 22234); sendAndReceiveSimplePacket(socket); - Assert.assertEquals(eventListner.getCode(), PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER); + Assert.assertEquals(eventListner.getCode(), PinpointServerSocketStateCode.RUN); int code= sendAndReceiveRegisterPacket(socket, getParams()); - Assert.assertEquals(eventListner.getCode(), PinpointServerSocketStateCode.RUN); + Assert.assertEquals(eventListner.getCode(), PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION); sendAndReceiveSimplePacket(socket); } finally { @@ -65,7 +65,7 @@ public class EventListnerTest { private int sendAndReceiveRegisterPacket(Socket socket, Map properties) throws ProtocolException, IOException { sendRegisterPacket(socket.getOutputStream(), properties); - ControlRegisterAgentConfirmPacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); + ControlEnableWorkerConfirmPacket packet = receiveRegisterConfirmPacket(socket.getInputStream()); Map result = (Map) ControlMessageEnDeconderUtils.decode(packet.getPayload()); return MapUtils.get(result, "code", Integer.class, -1); @@ -79,7 +79,7 @@ public class EventListnerTest { private void sendRegisterPacket(OutputStream outputStream, Map properties) throws ProtocolException, IOException { byte[] payload = ControlMessageEnDeconderUtils.encode(properties); - ControlRegisterAgentPacket packet = new ControlRegisterAgentPacket(1, payload); + ControlEnableWorkerPacket packet = new ControlEnableWorkerPacket(1, payload); ByteBuffer bb = packet.toBuffer().toByteBuffer(0, packet.toBuffer().writerIndex()); sendData(outputStream, bb.array()); @@ -98,14 +98,14 @@ public class EventListnerTest { outputStream.flush(); } - private ControlRegisterAgentConfirmPacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { + private ControlEnableWorkerConfirmPacket receiveRegisterConfirmPacket(InputStream inputStream) throws ProtocolException, IOException { byte[] payload = readData(inputStream); ChannelBuffer cb = ChannelBuffers.wrappedBuffer(payload); short packetType = cb.readShort(); - ControlRegisterAgentConfirmPacket packet = ControlRegisterAgentConfirmPacket.readBuffer(packetType, cb); + ControlEnableWorkerConfirmPacket packet = ControlEnableWorkerConfirmPacket.readBuffer(packetType, cb); return packet; } @@ -174,6 +174,12 @@ public class EventListnerTest { public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) { } + + @Override + public int handleEnableWorker(Map properties) { + logger.info("handleEnableWorker {}", properties); + return ControlEnableWorkerConfirmPacket.SUCCESS; + } } diff --git a/src/test/java/com/nhn/pinpoint/rpc/server/MessageListenerTest.java b/src/test/java/com/nhn/pinpoint/rpc/server/MessageListenerTest.java index e942774e4..8b8b34598 100644 --- a/src/test/java/com/nhn/pinpoint/rpc/server/MessageListenerTest.java +++ b/src/test/java/com/nhn/pinpoint/rpc/server/MessageListenerTest.java @@ -1,6 +1,5 @@ package com.nhn.pinpoint.rpc.server; -import java.io.IOException; import java.util.ArrayList; import java.util.HashMap; import java.util.List; @@ -40,7 +39,7 @@ public class MessageListenerTest { Thread.sleep(500); - List channelContextList = ss.getRegisterAgentChannelContext(); + List channelContextList = ss.getDuplexCommunicationChannelContext(); if (channelContextList.size() != 1) { Assert.fail(); } @@ -68,7 +67,7 @@ public class MessageListenerTest { PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, echoMessageListener); Thread.sleep(500); - List channelContextList = ss.getRegisterAgentChannelContext(); + List channelContextList = ss.getDuplexCommunicationChannelContext(); if (channelContextList.size() != 1) { Assert.fail(); } @@ -112,7 +111,7 @@ public class MessageListenerTest { Thread.sleep(500); - List channelContextList = ss.getRegisterAgentChannelContext(); + List channelContextList = ss.getDuplexCommunicationChannelContext(); if (channelContextList.size() != 2) { Assert.fail(); } @@ -151,10 +150,10 @@ public class MessageListenerTest { Thread.sleep(500); - ChannelContext channelContext = ss.getRegisterAgentChannelContext("application", "agent", (Long) params.get(AgentPropertiesType.START_TIMESTAMP.getName())); + ChannelContext channelContext = ss.getDuplexChannelContext("application", "agent", (Long) params.get(AgentPropertiesType.START_TIMESTAMP.getName())); Assert.assertNotNull(channelContext); - channelContext = ss.getRegisterAgentChannelContext("application", "agent", (Long) params.get(AgentPropertiesType.START_TIMESTAMP.getName()) + 1); + channelContext = ss.getDuplexChannelContext("application", "agent", (Long) params.get(AgentPropertiesType.START_TIMESTAMP.getName()) + 1); Assert.assertNull(channelContext); socket.close(); @@ -171,7 +170,7 @@ public class MessageListenerTest { private PinpointSocketFactory createPinpointSocketFactory(Map param) { PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); - pinpointSocketFactory.setAgentProperties(param); + pinpointSocketFactory.setProperties(param); return pinpointSocketFactory; } diff --git a/src/test/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateTest.java b/src/test/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateTest.java index 082a4af7a..046ac5743 100644 --- a/src/test/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateTest.java +++ b/src/test/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateTest.java @@ -14,12 +14,12 @@ public class PinpointServerSocketStateTest { public void changeStateTest1() { PinpointServerSocketState state = new PinpointServerSocketState(); - state.changeStateRunWithoutRegister(); - Assert.assertEquals(PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER, state.getCurrentState()); - state.changeStateRun(); Assert.assertEquals(PinpointServerSocketStateCode.RUN, state.getCurrentState()); + state.changeStateRunDuplexCommunication(); + Assert.assertEquals(PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION, state.getCurrentState()); + state.changeStateBeingShutdown(); Assert.assertEquals(PinpointServerSocketStateCode.BEING_SHUTDOWN, state.getCurrentState()); @@ -33,8 +33,8 @@ public class PinpointServerSocketStateTest { public void changeStateTest2() { PinpointServerSocketState state = new PinpointServerSocketState(); - state.changeStateRun(); - Assert.assertEquals(PinpointServerSocketStateCode.RUN, state.getCurrentState()); + state.changeStateRunDuplexCommunication(); + Assert.assertEquals(PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION, state.getCurrentState()); state.changeStateBeingShutdown(); Assert.assertEquals(PinpointServerSocketStateCode.BEING_SHUTDOWN, state.getCurrentState()); @@ -49,8 +49,8 @@ public class PinpointServerSocketStateTest { public void changeStateTest3() { PinpointServerSocketState state = new PinpointServerSocketState(); - state.changeStateRunWithoutRegister(); - Assert.assertEquals(PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER, state.getCurrentState()); + state.changeStateRun(); + Assert.assertEquals(PinpointServerSocketStateCode.RUN, state.getCurrentState()); state.changeStateUnexpectedShutdown(); Assert.assertEquals(PinpointServerSocketStateCode.UNEXPECTED_SHUTDOWN, state.getCurrentState()); @@ -62,8 +62,8 @@ public class PinpointServerSocketStateTest { public void changeStateTest4() { PinpointServerSocketState state = new PinpointServerSocketState(); - state.changeStateRunWithoutRegister(); - Assert.assertEquals(PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER, state.getCurrentState()); + state.changeStateRun(); + Assert.assertEquals(PinpointServerSocketStateCode.RUN, state.getCurrentState()); state.changeStateShutdown(); Assert.assertEquals(PinpointServerSocketStateCode.SHUTDOWN, state.getCurrentState()); @@ -73,8 +73,8 @@ public class PinpointServerSocketStateTest { public void changeStateTest5() { PinpointServerSocketState state = new PinpointServerSocketState(); - state.changeStateRun(); - Assert.assertEquals(PinpointServerSocketStateCode.RUN, state.getCurrentState()); + state.changeStateRunDuplexCommunication(); + Assert.assertEquals(PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION, state.getCurrentState()); state.changeStateShutdown(); Assert.assertEquals(PinpointServerSocketStateCode.SHUTDOWN, state.getCurrentState()); @@ -93,8 +93,8 @@ public class PinpointServerSocketStateTest { public void invalidChangeStateTest2() { PinpointServerSocketState state = new PinpointServerSocketState(); - state.changeStateRun(); - Assert.assertEquals(PinpointServerSocketStateCode.RUN, state.getCurrentState()); + state.changeStateRunDuplexCommunication(); + Assert.assertEquals(PinpointServerSocketStateCode.RUN_DUPLEX_COMMUNICATION, state.getCurrentState()); state.changeStateBeingShutdown(); Assert.assertEquals(PinpointServerSocketStateCode.BEING_SHUTDOWN, state.getCurrentState()); diff --git a/src/test/java/com/nhn/pinpoint/rpc/server/TestSeverMessageListener.java b/src/test/java/com/nhn/pinpoint/rpc/server/TestSeverMessageListener.java index 3c02e9c05..8f3e17d8b 100644 --- a/src/test/java/com/nhn/pinpoint/rpc/server/TestSeverMessageListener.java +++ b/src/test/java/com/nhn/pinpoint/rpc/server/TestSeverMessageListener.java @@ -2,11 +2,13 @@ package com.nhn.pinpoint.rpc.server; import com.nhn.pinpoint.rpc.TestByteUtils; import com.nhn.pinpoint.rpc.packet.*; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.util.ArrayList; import java.util.List; +import java.util.Map; /** * @author emeroad @@ -48,6 +50,12 @@ public class TestSeverMessageListener implements ServerMessageListener { } } + + @Override + public int handleEnableWorker(Map properties) { + logger.debug("handleEnableWorker properties:{} channel:{}", properties); + return ControlEnableWorkerConfirmPacket.SUCCESS; + } private void sendClose(ServerStreamChannel streamChannel) { sendMessageList.add(new byte[0]);