#4 Profiler - Collector간의 컨트롤 메시지 생성

여러 곳에서 사용할수 있게 rpc 코드를 좀 변경함

1) RegisterAgent -> Client Worker Enable
2) RegisterAgentConfirm -> Client Worker Enable Confirm
3) Confirm 조건은 외부로 뺌
This commit is contained in:
koo-taejin
2014-08-19 11:07:49 +09:00
parent 3a70220c5c
commit 8b28fc2ca8
19 changed files with 245 additions and 182 deletions
@@ -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;
}
}
@@ -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));
}
}
@@ -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();
}
}
@@ -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;
@@ -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);
}
}
@@ -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;
}
@@ -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;
}
@@ -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;
@@ -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;
}
@@ -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<ChannelContext> getRegisterAgentChannelContext() {
public List<ChannelContext> getDuplexCommunicationChannelContext() {
List<ChannelContext> channelContextList = new ArrayList<ChannelContext>();
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;
@@ -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);
}
@@ -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<PinpointServerSocketStateCode> 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;
}
@@ -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);
}
@@ -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;
}
}
@@ -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() {
@@ -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;
}
}
@@ -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<ChannelContext> channelContextList = ss.getRegisterAgentChannelContext();
List<ChannelContext> 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<ChannelContext> channelContextList = ss.getRegisterAgentChannelContext();
List<ChannelContext> channelContextList = ss.getDuplexCommunicationChannelContext();
if (channelContextList.size() != 1) {
Assert.fail();
}
@@ -112,7 +111,7 @@ public class MessageListenerTest {
Thread.sleep(500);
List<ChannelContext> channelContextList = ss.getRegisterAgentChannelContext();
List<ChannelContext> 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;
}
@@ -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());
@@ -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]);