diff --git a/src/main/java/com/nhn/pinpoint/rpc/DefaultFuture.java b/src/main/java/com/nhn/pinpoint/rpc/DefaultFuture.java index 31e60305c..3364ce02f 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/DefaultFuture.java +++ b/src/main/java/com/nhn/pinpoint/rpc/DefaultFuture.java @@ -1,8 +1,5 @@ package com.nhn.pinpoint.rpc; -import java.util.ArrayList; -import java.util.List; - import org.jboss.netty.util.Timeout; import org.jboss.netty.util.TimerTask; import org.slf4j.Logger; @@ -10,7 +7,6 @@ import org.slf4j.LoggerFactory; /** * @author emeroad - * @author koo.taejin */ public class DefaultFuture implements TimerTask, Future { @@ -26,7 +22,8 @@ public class DefaultFuture implements TimerTask, Future { private Timeout timeout; private FailureEventHandler failureEventHandler; - private List> listeners = new ArrayList>(); + private FutureListener listener; + public DefaultFuture() { this(3000); @@ -125,13 +122,11 @@ public class DefaultFuture implements TimerTask, Future { } private void notifyListener() { - for (FutureListener listener : this.listeners) { - if (listener != null) { - fireOnComplete(listener); - } - } - - this.listeners = new ArrayList>(); + FutureListener listener = this.listener; + if (listener != null) { + fireOnComplete(listener); + this.listener = null; + } } protected void notifyFailureHandle() { @@ -144,8 +139,8 @@ public class DefaultFuture implements TimerTask, Future { } @Override - public boolean addListener(FutureListener listener) { - if (listener == null) { + public boolean setListener(FutureListener listener) { + if (listener == null) { throw new NullPointerException("listener"); } @@ -154,13 +149,13 @@ public class DefaultFuture implements TimerTask, Future { if (ready) { alreadyReady = true; } else { - this.listeners.add(listener); + this.listener = listener; } } if (alreadyReady) { - fireOnComplete(listener); + fireOnComplete(listener); } return !alreadyReady; } diff --git a/src/main/java/com/nhn/pinpoint/rpc/Future.java b/src/main/java/com/nhn/pinpoint/rpc/Future.java index 21e54b7bf..22c2a7a74 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/Future.java +++ b/src/main/java/com/nhn/pinpoint/rpc/Future.java @@ -2,7 +2,6 @@ package com.nhn.pinpoint.rpc; /** * @author emeroad - * @author koo.taejin */ public interface Future { @@ -14,7 +13,7 @@ public interface Future { boolean isSuccess(); - boolean addListener(FutureListener listener); + boolean setListener(FutureListener listener); boolean await(long timeoutMillis); diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java b/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java index cecf28cba..a29f7c87f 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/ChannelContext.java @@ -50,8 +50,38 @@ public class ChannelContext { return socketChannel; } - public PinpointServerSocketState getState() { - return state; + public PinpointServerSocketStateCode getCurrentStateCode() { + return state.getCurrentState(); + } + + public void changeStateRun() { + logger.debug("Channel({}) state will be changed {}.", channel, PinpointServerSocketStateCode.RUN); + state.changeStateRun(); + } + + public void changeStateRunWithoutRegister() { + logger.debug("Channel({}) state will be changed {}.", channel, PinpointServerSocketStateCode.RUN_WITHOUT_REGISTER); + state.changeStateRunWithoutRegister(); + } + + public void changeStateBeingShutdown() { + logger.debug("Channel({}) state will be changed {}.", channel, PinpointServerSocketStateCode.BEING_SHUTDOWN); + state.changeStateBeingShutdown(); + } + + public void changeStateShutdown() { + logger.debug("Channel({}) state will be changed {}.", channel, PinpointServerSocketStateCode.SHUTDOWN); + state.changeStateShutdown(); + } + + public void changeStateUnexpectedShutdown() { + logger.debug("Channel({}) state will be changed {}.", channel, PinpointServerSocketStateCode.UNEXPECTED_SHUTDOWN); + state.changeStateUnexpectedShutdown(); + } + + public void changeStateUnkownError() { + logger.debug("Channel({}) state will be changed {}.", channel, PinpointServerSocketStateCode.ERROR_UNKOWN); + state.changeStateUnkownError(); } public AgentProperties getAgentProperties() { diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java index 8e911bbec..d5826a69d 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocket.java @@ -1,17 +1,23 @@ package com.nhn.pinpoint.rpc.server; -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.*; -import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils; -import com.nhn.pinpoint.rpc.util.CpuUtils; -import com.nhn.pinpoint.rpc.util.LoggerFactorySetup; -import com.nhn.pinpoint.rpc.util.TimerFactory; +import java.net.InetAddress; +import java.net.InetSocketAddress; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; import org.jboss.netty.bootstrap.ServerBootstrap; -import org.jboss.netty.channel.*; +import org.jboss.netty.channel.Channel; +import org.jboss.netty.channel.ChannelFuture; +import org.jboss.netty.channel.ChannelFutureListener; +import org.jboss.netty.channel.ChannelHandlerContext; +import org.jboss.netty.channel.ChannelPipelineFactory; +import org.jboss.netty.channel.ChannelStateEvent; +import org.jboss.netty.channel.ExceptionEvent; +import org.jboss.netty.channel.MessageEvent; +import org.jboss.netty.channel.SimpleChannelHandler; import org.jboss.netty.channel.group.ChannelGroup; import org.jboss.netty.channel.group.ChannelGroupFuture; import org.jboss.netty.channel.group.ChannelGroupFutureListener; @@ -26,13 +32,25 @@ import org.jboss.netty.util.TimerTask; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.net.InetAddress; -import java.net.InetSocketAddress; -import java.util.HashMap; -import java.util.Map; -import java.util.concurrent.ExecutorService; -import java.util.concurrent.Executors; -import java.util.concurrent.TimeUnit; +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.Packet; +import com.nhn.pinpoint.rpc.packet.PacketType; +import com.nhn.pinpoint.rpc.packet.PingPacket; +import com.nhn.pinpoint.rpc.packet.RequestPacket; +import com.nhn.pinpoint.rpc.packet.SendPacket; +import com.nhn.pinpoint.rpc.packet.ServerClosePacket; +import com.nhn.pinpoint.rpc.packet.StreamClosePacket; +import com.nhn.pinpoint.rpc.packet.StreamCreatePacket; +import com.nhn.pinpoint.rpc.packet.StreamPacket; +import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils; +import com.nhn.pinpoint.rpc.util.CpuUtils; +import com.nhn.pinpoint.rpc.util.LoggerFactorySetup; +import com.nhn.pinpoint.rpc.util.TimerFactory; /** * @author emeroad @@ -188,7 +206,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { private void closeChannel(Channel channel) { logger.debug("received ClientClosePacket {}", channel); ChannelContext channelContext = getChannelContext(channel); - channelContext.getState().changeStateBeingShutdown(); + channelContext.changeStateBeingShutdown(); // 상대방이 닫는거에 반응해서 socket을 닫도록 하자. // channel.close(); @@ -232,7 +250,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { boolean isSuccess = context.setAgentProperties(new AgentProperties(properties)); // 이미 등록되어 있다면 상태를 변경하지 않음 if (isSuccess) { - context.getState().changeStateRun(); + context.changeStateRun(); } logger.debug("Channel({}) State changed to Run.", channel); } catch (ProtocolException e) { @@ -280,7 +298,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { prepareChannel(channel); ChannelContext channelContext = getChannelContext(channel); - channelContext.getState().changeStateRunWithoutRegister(); + channelContext.changeStateRunWithoutRegister(); super.channelConnected(ctx, e); } @@ -289,12 +307,12 @@ public class PinpointServerSocket extends SimpleChannelHandler { public void channelDisconnected(ChannelHandlerContext ctx, ChannelStateEvent e) throws Exception { final Channel channel = e.getChannel(); final ChannelContext channelContext = getChannelContext(channel); - PinpointServerSocketState state = channelContext.getState(); + PinpointServerSocketStateCode currentStateCode = channelContext.getCurrentStateCode(); - if (state.getCurrentState() != PinpointServerSocketStateCode.BEING_SHUTDOWN) { - state.changeStateShutdown(); + if (currentStateCode != PinpointServerSocketStateCode.BEING_SHUTDOWN) { + channelContext.changeStateShutdown(); } else { - state.changeStateUnexpectedShutdown(); + channelContext.changeStateUnexpectedShutdown(); } if (logger.isDebugEnabled()) { @@ -309,18 +327,18 @@ public class PinpointServerSocket extends SimpleChannelHandler { final Channel channel = e.getChannel(); final ChannelContext channelContext = getChannelContext(channel); - PinpointServerSocketState state = channelContext.getState(); + PinpointServerSocketStateCode currentStateCode = channelContext.getCurrentStateCode(); - if (state.getCurrentState() != PinpointServerSocketStateCode.BEING_SHUTDOWN) { + if (currentStateCode == PinpointServerSocketStateCode.BEING_SHUTDOWN) { if (logger.isDebugEnabled()) { logger.debug("client channelClosed. normal closed. {}", channel); } - state.changeStateShutdown(); + channelContext.changeStateShutdown(); } else if(released) { if (logger.isDebugEnabled()) { logger.debug("client channelClosed. server shutdown. {}", channel); } - state.changeStateShutdown(); + channelContext.changeStateShutdown(); } else { boolean check = checkIgnoreAddress(channel); if (check) { @@ -328,7 +346,7 @@ public class PinpointServerSocket extends SimpleChannelHandler { } else { logger.debug("checkAddress, Client channelClosed channelClosed {}", channel); } - state.changeStateUnexpectedShutdown(); + channelContext.changeStateUnexpectedShutdown(); } channelContext.closeAllStreamChannel(); } diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java index 12abde377..9c0c7814c 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketState.java @@ -19,6 +19,18 @@ public class PinpointServerSocketState { if (enable) { this.beforeState = this.currentState; this.currentState = state; + } else if (PinpointServerSocketStateCode.isFinished(this.currentState)) { + // 상태가 더 이상 변경할수 없는 것들은 로그만 출력 + // 이미 종료 상태이기 때문에 이렇게 처리해도 큰 문제가 없음 + PinpointServerSocketStateCode checkBefore = this.beforeState; + PinpointServerSocketStateCode checkCurrent = this.currentState; + + String errorMessage = cannotChangeMessage(checkBefore, checkCurrent, state); + + this.beforeState = this.currentState; + this.currentState = PinpointServerSocketStateCode.ERROR_ILLEGAL_STATE_CHANGE; + + logger.warn(errorMessage); } else { PinpointServerSocketStateCode checkBefore = this.beforeState; PinpointServerSocketStateCode checkCurrent = this.currentState; @@ -59,7 +71,11 @@ public class PinpointServerSocketState { } private String errorMessage(PinpointServerSocketStateCode checkBefore, PinpointServerSocketStateCode checkCurrent, PinpointServerSocketStateCode nextState) { - return "Invalid State(current:" + checkCurrent + " before:" + checkBefore + " next:" + nextState; + return "Invalid State(current:" + checkCurrent + " before:" + checkBefore + " next:" + nextState + ")"; + } + + private String cannotChangeMessage(PinpointServerSocketStateCode checkBefore, PinpointServerSocketStateCode checkCurrent, PinpointServerSocketStateCode nextState) { + return "Can not change State(current:" + checkCurrent + " before:" + checkBefore + " next:" + nextState + ")"; } public PinpointServerSocketStateCode getCurrentState() { diff --git a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java index 44cded1f2..d9878b1cf 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java +++ b/src/main/java/com/nhn/pinpoint/rpc/server/PinpointServerSocketStateCode.java @@ -26,7 +26,7 @@ public enum PinpointServerSocketStateCode { // 서버쪽에서 먼저 연결을 끊자는 메시지도 필요하다. // 예를 들어 HELLO 이후 다 확인했는데, 같은 Agent명이 있으면(?) 이걸 사용자에게 말해야 할까? 아닐까? 알림 등 ERROR_UNKOWN(RUN, RUN_WITHOUT_REGISTER), - ERROR_ILLEGAL_STATE_CHANGE(NONE, RUN, RUN_WITHOUT_REGISTER, BEING_SHUTDOWN, SHUTDOWN); + ERROR_ILLEGAL_STATE_CHANGE(NONE, RUN, RUN_WITHOUT_REGISTER, BEING_SHUTDOWN); private final Set validBeforeStateSet; @@ -61,5 +61,12 @@ public enum PinpointServerSocketStateCode { return false; } + + public static boolean isFinished(PinpointServerSocketStateCode code) { + if (code == SHUTDOWN || code == UNEXPECTED_SHUTDOWN || code == ERROR_UNKOWN || code == ERROR_ILLEGAL_STATE_CHANGE) { + return true; + } + return false; + } } diff --git a/src/test/java/com/nhn/pinpoint/rpc/FutureTest.java b/src/test/java/com/nhn/pinpoint/rpc/FutureTest.java index d1c6b6810..956f23ca4 100644 --- a/src/test/java/com/nhn/pinpoint/rpc/FutureTest.java +++ b/src/test/java/com/nhn/pinpoint/rpc/FutureTest.java @@ -16,18 +16,17 @@ public class FutureTest { DefaultFuture future = new DefaultFuture(); SimpleListener listener1 = new SimpleListener(); - SimpleListener listener2 = new SimpleListener(); - future.addListener(listener1); - future.addListener(listener2); + future.setListener(listener1); +// future.addListener(listener2); Assert.assertFalse(listener1.isFinished()); - Assert.assertFalse(listener2.isFinished()); +// Assert.assertFalse(listener2.isFinished()); future.setResult("Hello"); Assert.assertTrue(listener1.isFinished()); - Assert.assertTrue(listener2.isFinished()); +// Assert.assertTrue(listener2.isFinished()); } @Test @@ -38,7 +37,7 @@ public class FutureTest { future.setResult("Hello"); - future.addListener(listener); + future.setListener(listener); Assert.assertTrue(listener.isFinished()); }