Merge pull request #201 from koo-taejin/#136

Improve state synchronization between server and client. #136
This commit is contained in:
koo-taejin
2015-03-03 17:33:23 +09:00
20 changed files with 820 additions and 520 deletions
@@ -40,8 +40,8 @@ import com.navercorp.pinpoint.collector.cluster.zookeeper.job.Job;
import com.navercorp.pinpoint.collector.cluster.zookeeper.job.UpdateJob;
import com.navercorp.pinpoint.collector.receiver.tcp.AgentHandshakePropertyType;
import com.navercorp.pinpoint.common.util.PinpointThreadFactory;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.server.PinpointServer;
import com.navercorp.pinpoint.rpc.server.PinpointServerStateCode;
import com.navercorp.pinpoint.rpc.util.MapUtils;
/**
@@ -187,7 +187,7 @@ public class ZookeeperLatestJobWorker implements Runnable {
}
for (PinpointServer pinpointServer : pinpointServerRepository) {
if (PinpointServerStateCode.isFinished(pinpointServer.getCurrentStateCode())) {
if (SocketStateCode.isClosed(pinpointServer.getCurrentStateCode())) {
logger.info("LeakDetector Find Leak PinpointServer={}.", pinpointServer);
putJob(new DeleteJob(pinpointServer));
}
@@ -202,8 +202,8 @@ public class ZookeeperLatestJobWorker implements Runnable {
public boolean handleUpdate(UpdateJob job) {
PinpointServer pinpointServer = job.getPinpointServer();
PinpointServerStateCode code = pinpointServer.getCurrentStateCode();
if (PinpointServerStateCode.isFinished(code)) {
SocketStateCode code = pinpointServer.getCurrentStateCode();
if (SocketStateCode.isClosed(code)) {
putJob(new DeleteJob(pinpointServer));
return false;
}
@@ -26,15 +26,15 @@ import org.apache.commons.lang.StringUtils;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.navercorp.pinpoint.collector.cluster.PinpointServerClusterPoint;
import com.navercorp.pinpoint.collector.cluster.ClusterPointRepository;
import com.navercorp.pinpoint.collector.cluster.PinpointServerClusterPoint;
import com.navercorp.pinpoint.collector.cluster.WorkerState;
import com.navercorp.pinpoint.collector.cluster.WorkerStateContext;
import com.navercorp.pinpoint.collector.cluster.zookeeper.job.DeleteJob;
import com.navercorp.pinpoint.collector.cluster.zookeeper.job.UpdateJob;
import com.navercorp.pinpoint.collector.receiver.tcp.AgentHandshakePropertyType;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.server.PinpointServer;
import com.navercorp.pinpoint.rpc.server.PinpointServerStateCode;
import com.navercorp.pinpoint.rpc.server.handler.ChannelStateChangeEventHandler;
import com.navercorp.pinpoint.rpc.util.MapUtils;
@@ -113,7 +113,7 @@ public class ZookeeperProfilerClusterManager implements ChannelStateChangeEventH
}
@Override
public void eventPerformed(PinpointServer pinpointServer, PinpointServerStateCode stateCode) {
public void eventPerformed(PinpointServer pinpointServer, SocketStateCode stateCode) {
if (workerState.isStarted()) {
logger.info("eventPerformed PinpointServer={}, State={}", pinpointServer, stateCode);
@@ -124,12 +124,12 @@ public class ZookeeperProfilerClusterManager implements ChannelStateChangeEventH
return;
}
if (PinpointServerStateCode.RUN_DUPLEX == stateCode) {
if (SocketStateCode.RUN_DUPLEX == stateCode) {
UpdateJob job = new UpdateJob(pinpointServer, new byte[0]);
worker.putJob(job);
profileCluster.addClusterPoint(new PinpointServerClusterPoint(pinpointServer));
} else if (PinpointServerStateCode.isFinished(stateCode)) {
} else if (SocketStateCode.isClosed(stateCode)) {
DeleteJob job = new DeleteJob(pinpointServer);
worker.putJob(job);
@@ -143,7 +143,7 @@ public class ZookeeperProfilerClusterManager implements ChannelStateChangeEventH
}
@Override
public void exceptionCaught(PinpointServer pinpointServer, PinpointServerStateCode stateCode, Throwable e) {
public void exceptionCaught(PinpointServer pinpointServer, SocketStateCode stateCode, Throwable e) {
logger.warn("ZookeeperProfilerClusterManager exceptionCaught() (pinpointServer:{}, PinpointServerStateCode:{}). Error: {}.", pinpointServer, stateCode, e.getMessage(), e);
}
@@ -53,7 +53,7 @@ public class PinpointSocket {
public PinpointSocket(SocketHandler socketHandler) {
AssertUtils.assertNotNull(socketHandler, "socketHandler");
socketHandler.doHandshake();
socketHandler.doHandshake();
this.socketHandler = socketHandler;
socketHandler.setPinpointSocket(this);
@@ -87,7 +87,7 @@ public class PinpointSocket {
return false;
}
return this.reconnectEventListeners.add(eventListener);
return this.reconnectEventListeners.add(eventListener);
}
public boolean removePinpointSocketReconnectEventListener(PinpointSocketReconnectEventListener eventListener) {
@@ -398,6 +398,14 @@ public class PinpointSocketFactory {
public MessageListener getMessageListener() {
return messageListener;
}
public MessageListener getMessageListener(MessageListener defaultMessageListener) {
if (messageListener == null) {
return defaultMessageListener;
}
return messageListener;
}
public void setMessageListener(MessageListener messageListener) {
AssertUtils.assertNotNull(messageListener, "messageListener must not be null");
@@ -408,7 +416,15 @@ public class PinpointSocketFactory {
public ServerStreamChannelMessageListener getServerStreamChannelMessageListener() {
return serverStreamChannelMessageListener;
}
public ServerStreamChannelMessageListener getServerStreamChannelMessageListener(ServerStreamChannelMessageListener defaultStreamMessageListener) {
if (serverStreamChannelMessageListener == null) {
return defaultStreamMessageListener;
}
return serverStreamChannelMessageListener;
}
public void setServerStreamChannelMessageListener(ServerStreamChannelMessageListener serverStreamChannelMessageListener) {
AssertUtils.assertNotNull(messageListener, "messageListener must not be null");
@@ -118,6 +118,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
HashedWheelTimer timer = TimerFactory.createHashedWheelTimer("Pinpoint-SocketHandler-Timer", 100, TimeUnit.MILLISECONDS, 512);
timer.start();
this.channelTimer = timer;
this.pinpointSocketFactory = pinpointSocketFactory;
this.requestManager = new RequestManager(timer, timeoutMillis);
@@ -125,24 +126,10 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
this.handshakeRetryInterval = handshakeRetryInterval;
this.timeoutMillis = timeoutMillis;
MessageListener messageLisener = pinpointSocketFactory.getMessageListener();
if (messageLisener != null) {
this.messageListener = messageLisener;
} else {
this.messageListener = SimpleLoggingMessageListener.LISTENER;
}
ServerStreamChannelMessageListener serverStreamChannelMessageListener = pinpointSocketFactory.getServerStreamChannelMessageListener();
if (serverStreamChannelMessageListener != null) {
this.serverStreamChannelMessageListener = serverStreamChannelMessageListener;
} else {
this.serverStreamChannelMessageListener = DisabledServerStreamChannelMessageListener.INSTANCE;
}
this.messageListener = pinpointSocketFactory.getMessageListener(SimpleLoggingMessageListener.LISTENER);
this.serverStreamChannelMessageListener = pinpointSocketFactory.getServerStreamChannelMessageListener(DisabledServerStreamChannelMessageListener.INSTANCE);
this.objectUniqName = ClassUtils.simpleClassNameAndHashCodeString(this);
pinpointSocketFactory.getServerStreamChannelMessageListener();
this.handshaker = new PinpointClientSocketHandshaker(channelTimer, (int) handshakeRetryInterval, maxHandshakeCount);
}
@@ -198,12 +185,12 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
StreamChannelManager streamChannelManager = new StreamChannelManager(channel, IDGenerator.createOddIdGenerator(), serverStreamChannelMessageListener);
SocketHandlerContext context = new SocketHandlerContext(channel, streamChannelManager);
PinpointSocketHandlerContext context = new PinpointSocketHandlerContext(channel, streamChannelManager);
channel.setAttachment(context);
}
private SocketHandlerContext getChannelContext(Channel channel) {
return (SocketHandlerContext) channel.getAttachment();
private PinpointSocketHandlerContext getChannelContext(Channel channel) {
return (PinpointSocketHandlerContext) channel.getAttachment();
}
@Override
@@ -246,8 +233,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
return;
}
logger.debug("writePing {}", channel);
ChannelFuture write = this.channel.write(PingPacket.PING_PACKET);
write.addListener(pingWriteFailFutureListener);
write0(PingPacket.PING_PACKET, pingWriteFailFutureListener);
}
public void sendPing() {
@@ -255,10 +242,11 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
return;
}
logger.debug("sendPing {}", channel);
ChannelFuture write = this.channel.write(PingPacket.PING_PACKET);
write.awaitUninterruptibly();
if (!write.isSuccess()) {
Throwable cause = write.getCause();
ChannelFuture future = write0(PingPacket.PING_PACKET);
future.awaitUninterruptibly();
if (!future.isSuccess()) {
Throwable cause = future.getCause();
throw new PinpointSocketException("send ping failed. Error:" + cause.getMessage(), cause);
}
logger.debug("sendPing success {}", channel);
@@ -327,7 +315,7 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
ensureOpen();
SendPacket send = new SendPacket(bytes);
return this.channel.write(send);
return write0(send);
}
public Future<ResponseMessage> request(byte[] bytes) {
@@ -343,13 +331,9 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
}
RequestPacket request = new RequestPacket(bytes);
final Channel channel = this.channel;
final ChannelWriteFailListenableFuture<ResponseMessage> messageFuture = this.requestManager.register(request, this.timeoutMillis);
ChannelFuture write = channel.write(request);
write.addListener(messageFuture);
write0(request, messageFuture);
return messageFuture;
}
@@ -358,8 +342,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
ensureOpen();
final Channel channel = this.channel;
SocketHandlerContext context = getChannelContext(channel);
return context.getStreamChannelManager().openStreamChannel(payload, clientStreamChannelMessageListener);
PinpointSocketHandlerContext context = getChannelContext(channel);
return context.createStream(payload, clientStreamChannelMessageListener);
}
@Override
@@ -367,8 +351,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
ensureOpen();
final Channel channel = this.channel;
SocketHandlerContext context = getChannelContext(channel);
return context.getStreamChannelManager().findStreamChannel(streamChannelId);
PinpointSocketHandlerContext context = getChannelContext(channel);
return context.getStreamChannel(streamChannelId);
}
@Override
@@ -395,8 +379,8 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
case PacketType.APPLICATION_STREAM_RESPONSE:
case PacketType.APPLICATION_STREAM_PING:
case PacketType.APPLICATION_STREAM_PONG:
SocketHandlerContext context = getChannelContext(channel);
context.getStreamChannelManager().messageReceived((StreamPacket) message);
PinpointSocketHandlerContext context = getChannelContext(channel);
context.handleStreamEvent((StreamPacket) message);
return;
case PacketType.CONTROL_SERVER_CLOSE:
messageReceivedServerClosed(e.getChannel());
@@ -520,9 +504,9 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
}
// stream channel clear and send stream close packet
SocketHandlerContext context = getChannelContext(channel);
PinpointSocketHandlerContext context = getChannelContext(channel);
if (context != null) {
context.getStreamChannelManager().close();
context.closeAllStreamChannel();
}
}
@@ -582,6 +566,19 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
connectFuture.setResult(Result.FAIL);
}
}
private ChannelFuture write0(Object message) {
return write0(message, null);
}
private ChannelFuture write0(Object message, ChannelFutureListener futureListener) {
ChannelFuture future = channel.write(message);
if (futureListener != null) {
future.addListener(futureListener);
}
return future;
}
// Calling this method on a closed SocketHandler has no effect.
private void releaseResource() {
@@ -616,28 +613,4 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
handshaker.handshakeStart(channel, handshakeData);
}
static class SocketHandlerContext {
private final Channel channel;
private final StreamChannelManager streamChannelManager;
public SocketHandlerContext(Channel channel, StreamChannelManager streamChannelManager) {
if (channel == null) {
throw new NullPointerException("channel must not be null");
}
if (streamChannelManager == null) {
throw new NullPointerException("streamChannelManager must not be null");
}
this.channel = channel;
this.streamChannelManager = streamChannelManager;
}
public Channel getChannel() {
return channel;
}
public StreamChannelManager getStreamChannelManager() {
return streamChannelManager;
}
}
}
}
@@ -0,0 +1,65 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.client;
import org.jboss.netty.channel.Channel;
import com.navercorp.pinpoint.rpc.packet.stream.StreamPacket;
import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
import com.navercorp.pinpoint.rpc.stream.StreamChannelContext;
import com.navercorp.pinpoint.rpc.stream.StreamChannelManager;
/**
* @author Taejin Koo
*/
public class PinpointSocketHandlerContext {
private final Channel channel;
private final StreamChannelManager streamChannelManager;
public PinpointSocketHandlerContext(Channel channel, StreamChannelManager streamChannelManager) {
if (channel == null) {
throw new NullPointerException("channel must not be null");
}
if (streamChannelManager == null) {
throw new NullPointerException("streamChannelManager must not be null");
}
this.channel = channel;
this.streamChannelManager = streamChannelManager;
}
public Channel getChannel() {
return channel;
}
public ClientStreamChannelContext createStream(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) {
return streamChannelManager.openStreamChannel(payload, clientStreamChannelMessageListener);
}
public void handleStreamEvent(StreamPacket message) {
streamChannelManager.messageReceived(message);
}
public void closeAllStreamChannel() {
streamChannelManager.close();
}
public StreamChannelContext getStreamChannel(int streamChannelId) {
return streamChannelManager.findStreamChannel(streamChannelId);
}
}
@@ -0,0 +1,134 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.common;
/**
* @author Taejin Koo
*/
public class SocketState {
private SocketStateCode beforeState = SocketStateCode.NONE;
private SocketStateCode currentState = SocketStateCode.NONE;
public synchronized SocketStateChangeResult changeState(SocketStateCode nextState) {
boolean enable = this.currentState.canChangeState(nextState);
if (enable) {
this.beforeState = this.currentState;
this.currentState = nextState;
return new SocketStateChangeResult(true, beforeState, currentState, nextState);
}
return new SocketStateChangeResult(false, beforeState, currentState, nextState);
}
public SocketStateChangeResult stateToBeingConnect() {
SocketStateCode nextState = SocketStateCode.BEING_CONNECT;
return changeState(nextState);
}
public SocketStateChangeResult stateToConnected() {
SocketStateCode nextState = SocketStateCode.CONNECTED;
return changeState(nextState);
}
public SocketStateChangeResult stateToConnectFailed() {
SocketStateCode nextState = SocketStateCode.CONNECT_FAILED;
return changeState(nextState);
}
public SocketStateChangeResult stateToIgnore() {
SocketStateCode nextState = SocketStateCode.IGNORE;
return changeState(nextState);
}
public SocketStateChangeResult stateToRunWithoutHandshake() {
SocketStateCode nextState = SocketStateCode.RUN_WITHOUT_HANDSHAKE;
return changeState(nextState);
}
public SocketStateChangeResult stateToRunSimplex() {
SocketStateCode nextState = SocketStateCode.RUN_SIMPLEX;
return changeState(nextState);
}
public SocketStateChangeResult stateToRunDuplex() {
SocketStateCode nextState = SocketStateCode.RUN_DUPLEX;
return changeState(nextState);
}
public SocketStateChangeResult stateToBeingCloseByClient() {
SocketStateCode nextState = SocketStateCode.BEING_CLOSE_BY_CLIENT;
return changeState(nextState);
}
public SocketStateChangeResult stateToClosedByClient() {
SocketStateCode nextState = SocketStateCode.CLOSED_BY_CLIENT;
return changeState(nextState);
}
public SocketStateChangeResult stateToUnexpectedCloseByClient() {
SocketStateCode nextState = SocketStateCode.UNEXPECTED_CLOSE_BY_CLIENT;
return changeState(nextState);
}
public SocketStateChangeResult stateToBeingCloseByServer() {
SocketStateCode nextState = SocketStateCode.BEING_CLOSE_BY_SERVER;
return changeState(nextState);
}
public SocketStateChangeResult stateToClosedByServer() {
SocketStateCode nextState = SocketStateCode.CLOSED_BY_SERVER;
return changeState(nextState);
}
public SocketStateChangeResult stateToUnexpectedCloseByServer() {
SocketStateCode nextState = SocketStateCode.UNEXPECTED_CLOSE_BY_SERVER;
return changeState(nextState);
}
public SocketStateChangeResult stateToUnkownError() {
SocketStateCode nextState = SocketStateCode.ERROR_UNKOWN;
return changeState(nextState);
}
public synchronized SocketStateCode getCurrentState() {
return currentState;
}
@Override
public String toString() {
SocketStateCode beforeState;
SocketStateCode currentState;
synchronized (this) {
beforeState = this.beforeState;
currentState = this.currentState;
}
StringBuilder toString = new StringBuilder();
toString.append(this.getClass().getSimpleName());
toString.append("(");
toString.append(beforeState);
toString.append("->");
toString.append(currentState);
toString.append(")");
return toString.toString();
}
}
@@ -0,0 +1,78 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.common;
/**
* @author Taejin Koo
*/
public class SocketStateChangeResult {
private final boolean isChange;
private final SocketStateCode beforeState;
private final SocketStateCode currentState;
private final SocketStateCode updateWantedState;
public SocketStateChangeResult(boolean isChange, SocketStateCode beforeState, SocketStateCode currentState, SocketStateCode updateWantedState) {
this.isChange = isChange;
this.beforeState = beforeState;
this.currentState = currentState;
this.updateWantedState = updateWantedState;
}
public boolean isChange() {
return isChange;
}
public SocketStateCode getBeforeState() {
return beforeState;
}
public SocketStateCode getCurrentState() {
return currentState;
}
public SocketStateCode getUpdateWantedState() {
return updateWantedState;
}
@Override
public String toString() {
StringBuilder toString = new StringBuilder();
toString.append("Socket state change ");
if (isChange) {
toString.append("success");
} else {
toString.append("fail");
}
toString.append("(updateWanted:");
toString.append(updateWantedState);
toString.append(" ,before:");
toString.append(beforeState);
toString.append(" ,current:");
toString.append(currentState);
toString.append(").");
return toString.toString();
}
}
@@ -0,0 +1,176 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.common;
import java.util.HashSet;
import java.util.Set;
/**
* @author Taejin Koo
*/
public enum SocketStateCode {
NONE((byte) 1),
BEING_CONNECT((byte)2, NONE),
CONNECTED((byte) 3, NONE, BEING_CONNECT),
CONNECT_FAILED((byte)6, BEING_CONNECT),
IGNORE((byte) 9, CONNECTED),
RUN_WITHOUT_HANDSHAKE((byte) 10, CONNECTED),
RUN_SIMPLEX((byte) 11, RUN_WITHOUT_HANDSHAKE),
RUN_DUPLEX((byte) 12, RUN_WITHOUT_HANDSHAKE),
BEING_CLOSE_BY_CLIENT((byte) 20, RUN_WITHOUT_HANDSHAKE, RUN_SIMPLEX, RUN_DUPLEX),
CLOSED_BY_CLIENT((byte) 22, NONE, BEING_CLOSE_BY_CLIENT),
UNEXPECTED_CLOSE_BY_CLIENT((byte) 26, NONE, CONNECTED, RUN_WITHOUT_HANDSHAKE, RUN_SIMPLEX, RUN_DUPLEX),
BEING_CLOSE_BY_SERVER((byte) 30, RUN_WITHOUT_HANDSHAKE, RUN_SIMPLEX, RUN_DUPLEX),
CLOSED_BY_SERVER((byte) 32, NONE, BEING_CLOSE_BY_SERVER),
UNEXPECTED_CLOSE_BY_SERVER((byte) 36, NONE, CONNECTED, RUN_WITHOUT_HANDSHAKE, RUN_SIMPLEX, RUN_DUPLEX),
ERROR_UNKOWN((byte) 40),
ERROR_ILLEGAL_STATE_CHANGE((byte) 41);
private final byte id;
private final Set<SocketStateCode> validBeforeStateSet;
private SocketStateCode(byte id, SocketStateCode... validBeforeStates) {
this.id = id;
this.validBeforeStateSet = new HashSet<SocketStateCode>();
if (validBeforeStates != null) {
for (SocketStateCode eachStateCode : validBeforeStates) {
this.validBeforeStateSet.add(eachStateCode);
}
}
}
public boolean canChangeState(SocketStateCode nextState) {
Set<SocketStateCode> validBeforeStateSet = nextState.getValidBeforeStateSet();
if (validBeforeStateSet.contains(this)) {
return true;
}
return isError(nextState);
}
public boolean isBeforeConnected() {
return isBeforeConnected(this);
}
public static boolean isBeforeConnected(SocketStateCode code) {
switch (code) {
case NONE:
case BEING_CONNECT:
return true;
default:
return false;
}
}
public boolean isRun() {
return isRun(this);
}
public static boolean isRun(SocketStateCode code) {
switch (code) {
case RUN_WITHOUT_HANDSHAKE:
case RUN_SIMPLEX:
case RUN_DUPLEX:
return true;
default:
return false;
}
}
public boolean isRunDuplex() {
return isRunDuplex(this);
}
public static boolean isRunDuplex(SocketStateCode code) {
switch (code) {
case RUN_DUPLEX:
return true;
default:
return false;
}
}
public boolean onDisconnect() {
return onDisconnect(this);
}
public static boolean onDisconnect(SocketStateCode code) {
switch (code) {
case BEING_CLOSE_BY_CLIENT:
case BEING_CLOSE_BY_SERVER:
return true;
default:
return false;
}
}
public boolean isClosed() {
return isClosed(this);
}
public static boolean isClosed(SocketStateCode code) {
switch (code) {
case CLOSED_BY_CLIENT:
case UNEXPECTED_CLOSE_BY_CLIENT:
case CLOSED_BY_SERVER:
case UNEXPECTED_CLOSE_BY_SERVER:
case ERROR_UNKOWN:
case ERROR_ILLEGAL_STATE_CHANGE:
return true;
default:
return false;
}
}
private boolean isError(SocketStateCode code) {
switch (code) {
case ERROR_ILLEGAL_STATE_CHANGE:
case ERROR_UNKOWN:
return true;
default:
return false;
}
}
public static SocketStateCode getStateCode(byte id) {
SocketStateCode[] allStateCodes = SocketStateCode.values();
for (SocketStateCode code : allStateCodes) {
if (code.id == id) {
return code;
}
}
return null;
}
public byte getId() {
return id;
}
private Set<SocketStateCode> getValidBeforeStateSet() {
return validBeforeStateSet;
}
}
@@ -36,6 +36,9 @@ import com.navercorp.pinpoint.rpc.Future;
import com.navercorp.pinpoint.rpc.ResponseMessage;
import com.navercorp.pinpoint.rpc.client.RequestManager;
import com.navercorp.pinpoint.rpc.client.WriteFailFutureListener;
import com.navercorp.pinpoint.rpc.common.SocketState;
import com.navercorp.pinpoint.rpc.common.SocketStateChangeResult;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.control.ProtocolException;
import com.navercorp.pinpoint.rpc.packet.ControlHandshakePacket;
import com.navercorp.pinpoint.rpc.packet.ControlHandshakeResponsePacket;
@@ -68,7 +71,7 @@ public class DefaultPinpointServer implements PinpointServer {
private final Channel channel;
private final RequestManager requestManager;
private final PinpointServerState state;
private final SocketState state;
private final ServerMessageListener messageListener;
@@ -107,7 +110,7 @@ public class DefaultPinpointServer implements PinpointServer {
RequestManager requestManager = new RequestManager(serverConfig.getRequestManagerTimer(), serverConfig.getDefaultRequestTimeout());
this.requestManager = requestManager;
this.state = new PinpointServerState();
this.state = new SocketState();
this.objectUniqName = ClassUtils.simpleClassNameAndHashCodeString(this);
@@ -116,14 +119,37 @@ public class DefaultPinpointServer implements PinpointServer {
}
public void start() {
changeStateToRunWithoutHandshake();
logger.info("{} start() started. channel:{}.", objectUniqName, channel);
stateToConnected();
stateToRunWithoutHandshake();
logger.info("{} start() completed.", objectUniqName);
}
public void stop() {
if (PinpointServerStateCode.BEING_SHUTDOWN == getCurrentStateCode()) {
changeStateToShutdown();
logger.info("{} stop() started. channel:{}.", objectUniqName, channel);
stop(false);
logger.info("{} stop() completed.", objectUniqName);
}
public void stop(boolean serverStop) {
SocketStateCode currentStateCode = state.getCurrentState();
if (SocketStateCode.BEING_CLOSE_BY_SERVER == currentStateCode) {
stateToClosed();
} else if (SocketStateCode.BEING_CLOSE_BY_CLIENT == currentStateCode) {
stateToClosedByPeer();
} else if (SocketStateCode.isRun(currentStateCode) && serverStop) {
stateToUnexpectedClosed();
} else if (SocketStateCode.isRun(currentStateCode) ) {
stateToUnexpectedClosedByPeer();
} else if (SocketStateCode.isClosed(currentStateCode)){
logger.warn("{} stop(). Socket has closed state({}).", objectUniqName, currentStateCode);
} else {
changeStateToUnexpectedShutdown();
stateToErrorUnknown();
logger.warn("{} stop(). Socket has unexpected state.", objectUniqName, currentStateCode);
}
if (this.channel.isConnected()) {
@@ -191,11 +217,20 @@ public class DefaultPinpointServer implements PinpointServer {
@Override
public ClientStreamChannelContext createStream(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) {
return streamChannelManager.openStreamChannel(payload, clientStreamChannelMessageListener);
logger.info("{} createStream() started.", objectUniqName);
ClientStreamChannelContext streamChannel = streamChannelManager.openStreamChannel(payload, clientStreamChannelMessageListener);
logger.info("{} createStream() completed.", objectUniqName);
return streamChannel;
}
public void closeAllStreamChannel() {
logger.info("{} closeAllStreamChannel() started.", objectUniqName);
streamChannelManager.close();
logger.info("{} closeAllStreamChannel() completed.", objectUniqName);
}
@Override
@@ -218,33 +253,31 @@ public class DefaultPinpointServer implements PinpointServer {
}
public ChannelFuture sendClosePacket() {
logger.info("sendServerClosedPacket start");
logger.info("{} sendClosePacket() started.", objectUniqName);
PinpointServerStateCode errorCode = changeStateBeingShutdown();
if (errorCode == null) {
SocketStateChangeResult stateChangeResult = stateToBeingClose();
if (stateChangeResult.isChange()) {
final ChannelFuture writeFuture = this.channel.write(ServerClosePacket.DEFAULT_SERVER_CLOSE_PACKET);
writeFuture.addListener(serverCloseWriteListener);
logger.info("sendServerClosedPacket end");
logger.info("{} sendClosePacket() completed.", objectUniqName);
return writeFuture;
} else {
logger.info("sendServerClosedPacket fail. Error: change state failed.");
logger.info("{} sendClosePacket() failed. Error:{}.", objectUniqName, stateChangeResult);
return null;
}
}
@Override
public void messageReceived(Object message) {
// TODO Auto-generated method stub
if (!PinpointServerStateCode.isRun(getCurrentStateCode())) {
if (!isEnableCommunication()) {
// FIXME need change rules.
// as-is : do nothing when state is not run.
// candidate : close channel when state is not run.
logger.warn("{} messageReceived:{} from IllegalState this message will be ignore.", this, message);
logger.warn("{} messageReceived() failed. Error: Illegal state this message({}) will be ignore.", objectUniqName, message);
return;
}
final short packetType = getPacketType(message);
switch (packetType) {
case PacketType.APPLICATION_SEND: {
@@ -310,25 +343,37 @@ public class DefaultPinpointServer implements PinpointServer {
}
private void handleHandshake(ControlHandshakePacket handshakepacket) {
logger.info("{} handleHandshake() started. Packet:{}", objectUniqName, handshakepacket);
int requestId = handshakepacket.getRequestId();
Map<Object, Object> handshakeData = decodeHandshakePacket(handshakepacket);
HandshakeResponseCode responseCode = messageListener.handleHandshake(handshakeData);
boolean isFirst = setChannelProperties(handshakeData);
if (isFirst) {
if (HandshakeResponseCode.DUPLEX_COMMUNICATION == responseCode) {
changeStateToRunDuplex(PinpointServerStateCode.RUN_DUPLEX);
stateToRunDuplex();
} else if (HandshakeResponseCode.SIMPLEX_COMMUNICATION == responseCode) {
changeStateToRunSimplex(PinpointServerStateCode.RUN_SIMPLEX);
stateToRunSimplex();
}
}
logger.info("{} handleHandshake(). ResponseCode:{}", objectUniqName, responseCode);
Map<String, Object> responseData = createHandshakeResponse(responseCode, isFirst);
sendHandshakeResponse0(requestId, responseData);
logger.info("{} handleHandshake() completed.", objectUniqName);
}
private void handleClosePacket(Channel channel) {
logger.debug("handleClosePacket channel:{}", channel);
changeStateBeingShutdown(PinpointServerStateCode.BEING_SHUTDOWN);
logger.info("{} handleClosePacket() started.", objectUniqName);
SocketStateChangeResult stateChangeResult = stateToBeingCloseByPeer();
if (!stateChangeResult.isChange()) {
logger.info("{} handleClosePacket() failed. Error: {}", objectUniqName, stateChangeResult);
} else {
logger.info("{} handleClosePacket() completed.", objectUniqName);
}
}
private Map<String, Object> createHandshakeResponse(HandshakeResponseCode responseCode, boolean isFirst) {
@@ -354,8 +399,6 @@ public class DefaultPinpointServer implements PinpointServer {
private void sendHandshakeResponse0(int requestId, Map<String, Object> data) {
try {
logger.info("write HandshakeResponsePakcet. channel:{}, HandshakeResponseCode:{}.", channel, data);
byte[] resultPayload = ControlMessageEncodingUtils.encode(data);
ControlHandshakeResponsePacket packet = new ControlHandshakeResponsePacket(requestId, resultPayload);
@@ -378,79 +421,94 @@ public class DefaultPinpointServer implements PinpointServer {
}
@Override
public PinpointServerStateCode getCurrentStateCode() {
public SocketStateCode getCurrentStateCode() {
return state.getCurrentState();
}
private PinpointServerStateCode changeStateToRunWithoutHandshake(PinpointServerStateCode... skipLogicStateList) {
PinpointServerStateCode nextState = PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE;
return change0(nextState, skipLogicStateList);
private SocketStateChangeResult stateToConnected() {
SocketStateCode nextState = SocketStateCode.CONNECTED;
return stateTo(nextState);
}
private SocketStateChangeResult stateToRunWithoutHandshake() {
SocketStateCode nextState = SocketStateCode.RUN_WITHOUT_HANDSHAKE;
return stateTo(nextState);
}
private PinpointServerStateCode changeStateToRunSimplex(PinpointServerStateCode... skipLogicStateList) {
PinpointServerStateCode nextState = PinpointServerStateCode.RUN_SIMPLEX;
return change0(nextState, skipLogicStateList);
private SocketStateChangeResult stateToRunSimplex() {
SocketStateCode nextState = SocketStateCode.RUN_SIMPLEX;
return stateTo(nextState);
}
private SocketStateChangeResult stateToRunDuplex() {
SocketStateCode nextState = SocketStateCode.RUN_DUPLEX;
return stateTo(nextState);
}
private PinpointServerStateCode changeStateToRunDuplex(PinpointServerStateCode... skipLogicStateList) {
PinpointServerStateCode nextState = PinpointServerStateCode.RUN_DUPLEX;
return change0(nextState, skipLogicStateList);
private SocketStateChangeResult stateToBeingClose() {
SocketStateCode nextState = SocketStateCode.BEING_CLOSE_BY_SERVER;
return stateTo(nextState);
}
private SocketStateChangeResult stateToBeingCloseByPeer() {
SocketStateCode nextState = SocketStateCode.BEING_CLOSE_BY_CLIENT;
return stateTo(nextState);
}
private SocketStateChangeResult stateToClosed() {
SocketStateCode nextState = SocketStateCode.CLOSED_BY_SERVER;
return stateTo(nextState);
}
private SocketStateChangeResult stateToClosedByPeer() {
SocketStateCode nextState = SocketStateCode.CLOSED_BY_CLIENT;
return stateTo(nextState);
}
private SocketStateChangeResult stateToUnexpectedClosed() {
SocketStateCode nextState = SocketStateCode.UNEXPECTED_CLOSE_BY_SERVER;
return stateTo(nextState);
}
private SocketStateChangeResult stateToUnexpectedClosedByPeer() {
SocketStateCode nextState = SocketStateCode.UNEXPECTED_CLOSE_BY_CLIENT;
return stateTo(nextState);
}
private SocketStateChangeResult stateToErrorUnknown() {
SocketStateCode nextState = SocketStateCode.ERROR_UNKOWN;
return stateTo(nextState);
}
private PinpointServerStateCode changeStateBeingShutdown(PinpointServerStateCode... skipLogicStateList) {
PinpointServerStateCode nextState = PinpointServerStateCode.BEING_SHUTDOWN;
return change0(nextState, skipLogicStateList);
}
private SocketStateChangeResult stateTo(SocketStateCode nextState) {
logger.debug("{} stateTo() started. to:{}", objectUniqName, nextState);
private PinpointServerStateCode changeStateToShutdown(PinpointServerStateCode... skipLogicStateList) {
PinpointServerStateCode nextState = PinpointServerStateCode.SHUTDOWN;
return change0(nextState, skipLogicStateList);
}
private PinpointServerStateCode changeStateToUnexpectedShutdown(PinpointServerStateCode... skipLogicStateList) {
PinpointServerStateCode nextState = PinpointServerStateCode.UNEXPECTED_SHUTDOWN;
return change0(nextState, skipLogicStateList);
}
private PinpointServerStateCode change0(PinpointServerStateCode nextState, PinpointServerStateCode... skipLogicStateList) {
logger.debug("Channel({}) state will be changed {}.", channel, nextState);
PinpointServerStateCode beforeState = state.changeState(nextState, skipLogicStateList);
if (beforeState == null) {
SocketStateChangeResult stateChangeResult = state.changeState(nextState);
if (stateChangeResult.isChange()) {
executeChangeEventHandler(this, nextState);
} else if (!isSkipedChange0(beforeState, skipLogicStateList)) {
executeChangeEventHandler(this, state.getCurrentState());
}
return beforeState;
}
logger.info("{} stateTo() completed. {}", objectUniqName, stateChangeResult);
private boolean isSkipedChange0(PinpointServerStateCode beforeState, PinpointServerStateCode... skipLogicStateList) {
if (skipLogicStateList != null) {
for (PinpointServerStateCode skipLogicState : skipLogicStateList) {
if (beforeState == skipLogicState) {
return true;
}
}
}
return false;
return stateChangeResult;
}
private void executeChangeEventHandler(DefaultPinpointServer pinpointServer, PinpointServerStateCode stateCode) {
private void executeChangeEventHandler(DefaultPinpointServer pinpointServer, SocketStateCode nextState) {
for (ChannelStateChangeEventHandler eachListener : this.stateChangeEventListeners) {
try {
eachListener.eventPerformed(this, stateCode);
eachListener.eventPerformed(this, nextState);
} catch (Exception e) {
eachListener.exceptionCaught(this, stateCode, e);
eachListener.exceptionCaught(this, nextState, e);
}
}
}
public boolean isEnableCommunication() {
return PinpointServerStateCode.isRun(getCurrentStateCode());
return SocketStateCode.isRun(getCurrentStateCode());
}
public boolean isEnableDuplexCommunication() {
return PinpointServerStateCode.isRunDuplexCommunication(getCurrentStateCode());
return SocketStateCode.isRunDuplex(getCurrentStateCode());
}
@Override
@@ -20,6 +20,7 @@ import java.net.SocketAddress;
import java.util.Map;
import com.navercorp.pinpoint.rpc.Future;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.packet.RequestPacket;
import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
@@ -39,7 +40,7 @@ public interface PinpointServer {
void messageReceived(Object message);
PinpointServerStateCode getCurrentStateCode();
SocketStateCode getCurrentStateCode();
SocketAddress getRemoteAddress();
@@ -285,7 +285,6 @@ public class PinpointServerAcceptor implements PinpointServerConfig {
}
healthCheckTimer.stop();
sendServerClosePacket();
closePinpointServer();
if (serverChannel != null) {
@@ -302,7 +301,7 @@ public class PinpointServerAcceptor implements PinpointServerConfig {
requestManagerTimer.stop();
}
private void sendServerClosePacket() {
private void closePinpointServer() {
for (Channel channel : channelGroup) {
DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment();
@@ -312,14 +311,6 @@ public class PinpointServerAcceptor implements PinpointServerConfig {
}
}
private void closePinpointServer() {
for (Channel channel : channelGroup) {
DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment();
pinpointServer.sendClosePacket();
}
}
public List<PinpointServer> getWritableServerList() {
List<PinpointServer> pinpointServerList = new ArrayList<PinpointServer>();
@@ -374,7 +365,7 @@ public class PinpointServerAcceptor implements PinpointServerConfig {
DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment();
if (pinpointServer != null) {
pinpointServer.stop();
pinpointServer.stop(released);
}
super.channelDisconnected(ctx, e);
@@ -1,89 +0,0 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.server;
import java.util.Arrays;
import java.util.List;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
/**
* @author koo.taejin
*/
class PinpointServerState {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private PinpointServerStateCode beforeState = PinpointServerStateCode.NONE;
private PinpointServerStateCode currentState = PinpointServerStateCode.NONE;
private synchronized PinpointServerStateCode changeState0(PinpointServerStateCode nextState, List<PinpointServerStateCode> skipLogicStateList, boolean throwException) {
if (skipLogicStateList != null) {
for (PinpointServerStateCode skipLogicState : skipLogicStateList) {
if (this.currentState == skipLogicState) {
return currentState;
}
}
}
boolean enable = this.currentState.canChangeState(nextState);
if (enable) {
this.beforeState = this.currentState;
this.currentState = nextState;
return null;
}
// if state can't be changed, just log.
// no problem because the state of socket has been already closed.
PinpointServerStateCode checkBefore = this.beforeState;
PinpointServerStateCode checkCurrent = this.currentState;
String errorMessage = cannotChangeMessage(checkBefore, checkCurrent, nextState);
this.beforeState = this.currentState;
this.currentState = PinpointServerStateCode.ERROR_ILLEGAL_STATE_CHANGE;
if (throwException) {
throw new IllegalStateException(errorMessage);
} else {
logger.warn(errorMessage);
return beforeState;
}
}
/**
* @return <tt>null</tt> if state changed expected value. or
* <tt>currentState</tt> if state do not changed.
*/
public PinpointServerStateCode changeState(PinpointServerStateCode nextState, PinpointServerStateCode... skipLogicStateList) {
return changeState0(nextState, Arrays.asList(skipLogicStateList), false);
}
public PinpointServerStateCode changeStateThrowWhenFailed(PinpointServerStateCode nextState, PinpointServerStateCode... skipLogicStateList) {
return changeState0(nextState, Arrays.asList(skipLogicStateList), true);
}
private String cannotChangeMessage(PinpointServerStateCode checkBefore, PinpointServerStateCode checkCurrent, PinpointServerStateCode nextState) {
return "Can not change State(current:" + checkCurrent + " before:" + checkBefore + " next:" + nextState + ")";
}
public synchronized PinpointServerStateCode getCurrentState() {
return currentState;
}
}
@@ -1,110 +0,0 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.server;
import java.util.HashSet;
import java.util.Set;
/**
* @author koo.taejin
*/
public enum PinpointServerStateCode {
// NONE : No event
// RUN_WITHOUT_HANDSHAKE : can send message only to server without handshake each other.
// RUN_SIMPLEX : can send message only to server
// RUN_DUPLEX_COMMUNICATION : can communicate each other by full-duplex
// BEING_SHUTDOWN : received a close packet from a peer first and releasing resources
// SHUTDOWN : has been closed
// UNEXPECTED_SHUTDOWN : has not received a close packet from a peer but a peer has been shutdown
NONE(),
RUN_WITHOUT_HANDSHAKE(NONE), //Simplex Communication
RUN_SIMPLEX(NONE, RUN_WITHOUT_HANDSHAKE),
RUN_DUPLEX(NONE, RUN_WITHOUT_HANDSHAKE),
BEING_SHUTDOWN(RUN_SIMPLEX, RUN_DUPLEX, RUN_WITHOUT_HANDSHAKE),
SHUTDOWN(RUN_SIMPLEX, RUN_DUPLEX, RUN_WITHOUT_HANDSHAKE, BEING_SHUTDOWN),
UNEXPECTED_SHUTDOWN(RUN_SIMPLEX, RUN_DUPLEX, RUN_WITHOUT_HANDSHAKE),
// need messages to close a connection from server to agent
// for example, checked all of needed things followed by HELLO, if a same agent name exists, have to notify that to agent or not?
ERROR_UNKOWN(RUN_SIMPLEX, RUN_DUPLEX, RUN_WITHOUT_HANDSHAKE),
ERROR_ILLEGAL_STATE_CHANGE(NONE, RUN_SIMPLEX, RUN_DUPLEX, RUN_WITHOUT_HANDSHAKE, BEING_SHUTDOWN, SHUTDOWN);
private final Set<PinpointServerStateCode> validBeforeStateSet;
private PinpointServerStateCode(PinpointServerStateCode... validBeforeStates) {
this.validBeforeStateSet = new HashSet<PinpointServerStateCode>();
if (validBeforeStates != null) {
for (PinpointServerStateCode eachStateCode : validBeforeStates) {
getValidBeforeStateSet().add(eachStateCode);
}
}
}
public boolean canChangeState(PinpointServerStateCode nextState) {
Set<PinpointServerStateCode> validBeforeStateSet = nextState.getValidBeforeStateSet();
if (validBeforeStateSet.contains(this)) {
return true;
}
return false;
}
public Set<PinpointServerStateCode> getValidBeforeStateSet() {
return validBeforeStateSet;
}
public static boolean isRun(PinpointServerStateCode code) {
if (code == RUN_SIMPLEX || code == RUN_DUPLEX || code == RUN_WITHOUT_HANDSHAKE) {
return true;
}
return false;
}
public static boolean isRunDuplexCommunication(PinpointServerStateCode code) {
if (code == RUN_DUPLEX) {
return true;
}
return false;
}
public static boolean isFinished(PinpointServerStateCode code) {
if (code == SHUTDOWN || code == UNEXPECTED_SHUTDOWN || code == ERROR_UNKOWN || code == ERROR_ILLEGAL_STATE_CHANGE) {
return true;
}
return false;
}
public static PinpointServerStateCode getStateCode(String name) {
PinpointServerStateCode[] allStateCodes = PinpointServerStateCode.values();
for (PinpointServerStateCode code : allStateCodes) {
if (code.name().equalsIgnoreCase(name)) {
return code;
}
}
return null;
}
}
@@ -16,16 +16,16 @@
package com.navercorp.pinpoint.rpc.server.handler;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.server.PinpointServer;
import com.navercorp.pinpoint.rpc.server.PinpointServerStateCode;
/**
* @author koo.taejin
*/
public interface ChannelStateChangeEventHandler {
void eventPerformed(PinpointServer pinpointServer, PinpointServerStateCode stateCode) throws Exception;
void eventPerformed(PinpointServer pinpointServer, SocketStateCode stateCode) throws Exception;
void exceptionCaught(PinpointServer pinpointServer, PinpointServerStateCode stateCode, Throwable e);
void exceptionCaught(PinpointServer pinpointServer, SocketStateCode stateCode, Throwable e);
}
@@ -19,8 +19,8 @@ package com.navercorp.pinpoint.rpc.server.handler;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.server.PinpointServer;
import com.navercorp.pinpoint.rpc.server.PinpointServerStateCode;
/**
* @author koo.taejin
@@ -32,12 +32,12 @@ public class DoNothingChannelStateEventHandler implements ChannelStateChangeEven
public static final ChannelStateChangeEventHandler INSTANCE = new DoNothingChannelStateEventHandler();
@Override
public void eventPerformed(PinpointServer pinpointServer, PinpointServerStateCode stateCode) {
public void eventPerformed(PinpointServer pinpointServer, SocketStateCode stateCode) {
logger.info("{} eventPerformed(). pinpointServer:{}, code:{}", this.getClass().getSimpleName(), pinpointServer, stateCode);
}
@Override
public void exceptionCaught(PinpointServer pinpointServer, PinpointServerStateCode stateCode, Throwable e) {
public void exceptionCaught(PinpointServer pinpointServer, SocketStateCode stateCode, Throwable e) {
logger.warn("{} exceptionCaught(). pinpointServer:{}, code:{}. Error: {}.", this.getClass().getSimpleName(), pinpointServer, stateCode, e.getMessage(), e);
}
@@ -21,8 +21,8 @@ import java.util.concurrent.Executor;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.server.PinpointServer;
import com.navercorp.pinpoint.rpc.server.PinpointServerStateCode;
/**
* @author koo.taejin
@@ -40,7 +40,7 @@ public abstract class ExecutionChannelStateChangeEventHandler implements Channel
}
@Override
public void eventPerformed(PinpointServer pinpointServer, PinpointServerStateCode stateCode) {
public void eventPerformed(PinpointServer pinpointServer, SocketStateCode stateCode) {
logger.info("{} eventPerformed {}:{}", this.getClass().getSimpleName(), pinpointServer, stateCode);
Execution execution = new Execution(pinpointServer, stateCode);
@@ -49,9 +49,9 @@ public abstract class ExecutionChannelStateChangeEventHandler implements Channel
private class Execution implements Runnable {
private final PinpointServer pinpointServer;
private final PinpointServerStateCode stateCode;
private final SocketStateCode stateCode;
public Execution(PinpointServer pinpointServer, PinpointServerStateCode stateCode) {
public Execution(PinpointServer pinpointServer, SocketStateCode stateCode) {
this.pinpointServer = pinpointServer;
this.stateCode = stateCode;
}
@@ -31,6 +31,7 @@ import org.junit.Test;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.control.ProtocolException;
import com.navercorp.pinpoint.rpc.packet.ControlHandshakePacket;
import com.navercorp.pinpoint.rpc.packet.ControlHandshakeResponsePacket;
@@ -72,10 +73,10 @@ public class EventHandlerTest {
try {
socket = new Socket("127.0.0.1", bindPort);
sendAndReceiveSimplePacket(socket);
Assert.assertEquals(eventHandler.getCode(), PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE);
Assert.assertEquals(eventHandler.getCode(), SocketStateCode.RUN_WITHOUT_HANDSHAKE);
int code = sendAndReceiveRegisterPacket(socket, PinpointRPCTestUtils.getParams());
Assert.assertEquals(eventHandler.getCode(), PinpointServerStateCode.RUN_DUPLEX);
Assert.assertEquals(eventHandler.getCode(), SocketStateCode.RUN_DUPLEX);
sendAndReceiveSimplePacket(socket);
} finally {
@@ -214,18 +215,18 @@ public class EventHandlerTest {
class EventHandler implements ChannelStateChangeEventHandler {
private PinpointServerStateCode code;
private SocketStateCode code;
@Override
public void eventPerformed(PinpointServer pinpointServer, PinpointServerStateCode stateCode) {
public void eventPerformed(PinpointServer pinpointServer, SocketStateCode stateCode) {
this.code = stateCode;
}
@Override
public void exceptionCaught(PinpointServer pinpointServer, PinpointServerStateCode stateCode, Throwable e) {
public void exceptionCaught(PinpointServer pinpointServer, SocketStateCode stateCode, Throwable e) {
}
public PinpointServerStateCode getCode() {
public SocketStateCode getCode() {
return code;
}
}
@@ -235,12 +236,12 @@ public class EventHandlerTest {
private int errorCount = 0;
@Override
public void eventPerformed(PinpointServer pinpointServer, PinpointServerStateCode stateCode) throws Exception {
public void eventPerformed(PinpointServer pinpointServer, SocketStateCode stateCode) throws Exception {
throw new Exception("always error.");
}
@Override
public void exceptionCaught(PinpointServer pinpointServer, PinpointServerStateCode stateCode, Throwable e) {
public void exceptionCaught(PinpointServer pinpointServer, SocketStateCode stateCode, Throwable e) {
errorCount++;
}
@@ -1,144 +0,0 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.server;
import org.junit.Assert;
import org.junit.Test;
import com.navercorp.pinpoint.rpc.server.PinpointServerState;
import com.navercorp.pinpoint.rpc.server.PinpointServerStateCode;
/**
* @author koo.taejin
*/
public class PinpointServerSocketStateTest {
// basic type of connection's lifecycle between peers.
// RUN -> RUN_DUPLEX_COMMUNICATION -> BEING_SHUTDOWN -> connection closed
@Test
public void changeStateTest1() {
PinpointServerState state = new PinpointServerState();
state.changeState(PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE);
Assert.assertEquals(PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE, state.getCurrentState());
state.changeState(PinpointServerStateCode.RUN_DUPLEX);
Assert.assertEquals(PinpointServerStateCode.RUN_DUPLEX, state.getCurrentState());
state.changeState(PinpointServerStateCode.BEING_SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.BEING_SHUTDOWN, state.getCurrentState());
state.changeState(PinpointServerStateCode.SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.SHUTDOWN, state.getCurrentState());
}
// basic type of connection's lifecycle between peers.
// RUN_DUPLEX_COMMUNICATION -> RUN_DUPLEX_COMMUNICATION -> BEING_SHUTDOWN -> connection closed
@Test
public void changeStateTest2() {
PinpointServerState state = new PinpointServerState();
state.changeState(PinpointServerStateCode.RUN_DUPLEX);
Assert.assertEquals(PinpointServerStateCode.RUN_DUPLEX, state.getCurrentState());
PinpointServerStateCode currentState = state.changeState(PinpointServerStateCode.BEING_SHUTDOWN);
Assert.assertNull(currentState);
Assert.assertEquals(PinpointServerStateCode.BEING_SHUTDOWN, state.getCurrentState());
currentState = state.changeState(PinpointServerStateCode.BEING_SHUTDOWN, PinpointServerStateCode.BEING_SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.BEING_SHUTDOWN, currentState);
Assert.assertEquals(PinpointServerStateCode.BEING_SHUTDOWN, state.getCurrentState());
state.changeState(PinpointServerStateCode.SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.SHUTDOWN, state.getCurrentState());
}
@Test
public void changeStateTest3() {
PinpointServerState state = new PinpointServerState();
state.changeState(PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE);
Assert.assertEquals(PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE, state.getCurrentState());
state.changeState(PinpointServerStateCode.UNEXPECTED_SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.UNEXPECTED_SHUTDOWN, state.getCurrentState());
}
@Test
public void changeStateTest4() {
PinpointServerState state = new PinpointServerState();
state.changeState(PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE);
Assert.assertEquals(PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE, state.getCurrentState());
state.changeState(PinpointServerStateCode.SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.SHUTDOWN, state.getCurrentState());
}
@Test
public void changeStateTest5() {
PinpointServerState state = new PinpointServerState();
state.changeState(PinpointServerStateCode.RUN_DUPLEX);
Assert.assertEquals(PinpointServerStateCode.RUN_DUPLEX, state.getCurrentState());
state.changeState(PinpointServerStateCode.SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.SHUTDOWN, state.getCurrentState());
}
@Test
public void invalidChangeStateTest1() {
PinpointServerState state = new PinpointServerState();
PinpointServerStateCode beforeCode = state.changeState(PinpointServerStateCode.BEING_SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.NONE, beforeCode);
}
@Test(expected = IllegalStateException.class)
public void invalidChangeStateTest2() {
PinpointServerState state = new PinpointServerState();
state.changeStateThrowWhenFailed(PinpointServerStateCode.BEING_SHUTDOWN);
}
@Test
public void invalidChangeStateTest3() {
PinpointServerState state = new PinpointServerState();
state.changeState(PinpointServerStateCode.RUN_DUPLEX);
Assert.assertEquals(PinpointServerStateCode.RUN_DUPLEX, state.getCurrentState());
state.changeStateThrowWhenFailed(PinpointServerStateCode.BEING_SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.BEING_SHUTDOWN, state.getCurrentState());
PinpointServerStateCode beforeCode = state.changeState(PinpointServerStateCode.UNEXPECTED_SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.BEING_SHUTDOWN, beforeCode);
}
@Test(expected = IllegalStateException.class)
public void invalidChangeStateTest4() {
PinpointServerState state = new PinpointServerState();
state.changeState(PinpointServerStateCode.RUN_DUPLEX);
Assert.assertEquals(PinpointServerStateCode.RUN_DUPLEX, state.getCurrentState());
state.changeStateThrowWhenFailed(PinpointServerStateCode.BEING_SHUTDOWN);
Assert.assertEquals(PinpointServerStateCode.BEING_SHUTDOWN, state.getCurrentState());
PinpointServerStateCode beforeCode = state.changeStateThrowWhenFailed(PinpointServerStateCode.UNEXPECTED_SHUTDOWN);
}
}
@@ -0,0 +1,150 @@
/*
* Copyright 2014 NAVER Corp.
*
* Licensed under the Apache License, Version 2.0 (the "License");
* you may not use this file except in compliance with the License.
* You may obtain a copy of the License at
*
* http://www.apache.org/licenses/LICENSE-2.0
*
* Unless required by applicable law or agreed to in writing, software
* distributed under the License is distributed on an "AS IS" BASIS,
* WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
* See the License for the specific language governing permissions and
* limitations under the License.
*/
package com.navercorp.pinpoint.rpc.server;
import java.io.IOException;
import java.net.Socket;
import java.util.List;
import java.util.Map;
import org.jboss.netty.buffer.ChannelBuffer;
import org.junit.Assert;
import org.junit.BeforeClass;
import org.junit.Test;
import com.navercorp.pinpoint.rpc.client.PinpointSocket;
import com.navercorp.pinpoint.rpc.client.PinpointSocketFactory;
import com.navercorp.pinpoint.rpc.common.SocketStateCode;
import com.navercorp.pinpoint.rpc.control.ProtocolException;
import com.navercorp.pinpoint.rpc.packet.ControlHandshakePacket;
import com.navercorp.pinpoint.rpc.util.ControlMessageEncodingUtils;
import com.navercorp.pinpoint.rpc.util.PinpointRPCTestUtils;
/**
* @author Taejin Koo
*/
public class PinpointServerStateTest {
private static int bindPort;
@BeforeClass
public static void setUp() throws IOException {
bindPort = PinpointRPCTestUtils.findAvailablePort();
}
@Test
public void closeByPeerTest() throws InterruptedException {
PinpointServerAcceptor serverAcceptor = null;
try {
serverAcceptor = PinpointRPCTestUtils.createPinpointServerFactory(bindPort, PinpointRPCTestUtils.createEchoServerListener());
PinpointSocketFactory clientSocketFactory1 = PinpointRPCTestUtils.createSocketFactory(PinpointRPCTestUtils.getParams(), PinpointRPCTestUtils.createEchoClientListener());
PinpointSocket pinpointSocket = clientSocketFactory1.connect("127.0.0.1", bindPort);
Thread.sleep(1000);
List<PinpointServer> pinpointServerList = serverAcceptor.getWritableServerList();
PinpointServer pinpointServer = pinpointServerList.get(0);
Assert.assertEquals(SocketStateCode.RUN_DUPLEX, pinpointServer.getCurrentStateCode());
pinpointSocket.close();
Thread.sleep(1000);
Assert.assertEquals(SocketStateCode.CLOSED_BY_CLIENT, pinpointServer.getCurrentStateCode());
} finally {
PinpointRPCTestUtils.close(serverAcceptor);
}
}
@Test
public void closeTest() throws InterruptedException {
PinpointServerAcceptor serverAcceptor = null;
try {
serverAcceptor = PinpointRPCTestUtils.createPinpointServerFactory(bindPort, PinpointRPCTestUtils.createEchoServerListener());
PinpointSocketFactory clientSocketFactory1 = PinpointRPCTestUtils.createSocketFactory(PinpointRPCTestUtils.getParams(), PinpointRPCTestUtils.createEchoClientListener());
PinpointSocket pinpointSocket = clientSocketFactory1.connect("127.0.0.1", bindPort);
Thread.sleep(1000);
List<PinpointServer> pinpointServerList = serverAcceptor.getWritableServerList();
PinpointServer pinpointServer = pinpointServerList.get(0);
Assert.assertEquals(SocketStateCode.RUN_DUPLEX, pinpointServer.getCurrentStateCode());
serverAcceptor.close();
Thread.sleep(1000);
Assert.assertEquals(SocketStateCode.CLOSED_BY_SERVER, pinpointServer.getCurrentStateCode());
} finally {
PinpointRPCTestUtils.close(serverAcceptor);
}
}
@Test
public void unexpecteCloseByPeerTest() throws InterruptedException, IOException, ProtocolException {
PinpointServerAcceptor serverAcceptor = null;
try {
serverAcceptor = PinpointRPCTestUtils.createPinpointServerFactory(bindPort, PinpointRPCTestUtils.createEchoServerListener());
PinpointSocketFactory clientSocketFactory1 = PinpointRPCTestUtils.createSocketFactory(PinpointRPCTestUtils.getParams(), PinpointRPCTestUtils.createEchoClientListener());
PinpointSocket pinpointSocket = clientSocketFactory1.connect("127.0.0.1", bindPort);
Thread.sleep(1000);
List<PinpointServer> pinpointServerList = serverAcceptor.getWritableServerList();
PinpointServer pinpointServer = pinpointServerList.get(0);
Assert.assertEquals(SocketStateCode.RUN_DUPLEX, pinpointServer.getCurrentStateCode());
((DefaultPinpointServer)pinpointServer).stop(true);
Thread.sleep(1000);
Assert.assertEquals(SocketStateCode.UNEXPECTED_CLOSE_BY_SERVER, pinpointServer.getCurrentStateCode());
} finally {
PinpointRPCTestUtils.close(serverAcceptor);
}
}
@Test
public void unexpecteCloseTest() throws InterruptedException, IOException, ProtocolException {
PinpointServerAcceptor serverAcceptor = null;
try {
serverAcceptor = PinpointRPCTestUtils.createPinpointServerFactory(bindPort, PinpointRPCTestUtils.createEchoServerListener());
Socket socket = new Socket("127.0.0.1", bindPort);
socket.getOutputStream().write(createHandshakePayload(PinpointRPCTestUtils.getParams()));
socket.getOutputStream().flush();
Thread.sleep(1000);
List<PinpointServer> pinpointServerList = serverAcceptor.getWritableServerList();
PinpointServer pinpointServer = pinpointServerList.get(0);
Assert.assertEquals(SocketStateCode.RUN_DUPLEX, pinpointServer.getCurrentStateCode());
socket.close();
Thread.sleep(1000);
Assert.assertEquals(SocketStateCode.UNEXPECTED_CLOSE_BY_CLIENT, pinpointServer.getCurrentStateCode());
} finally {
PinpointRPCTestUtils.close(serverAcceptor);
}
}
private byte[] createHandshakePayload(Map<String, Object> data) throws ProtocolException {
byte[] payload = ControlMessageEncodingUtils.encode(data);
ControlHandshakePacket handshakePacket = new ControlHandshakePacket(payload);
ChannelBuffer channelBuffer = handshakePacket.toBuffer();
return channelBuffer.toByteBuffer().array();
}
}