diff --git a/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocket.java b/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocket.java index e48ce1004..2a2d5d22f 100644 --- a/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocket.java +++ b/src/main/java/com/nhn/pinpoint/rpc/client/PinpointSocket.java @@ -1,5 +1,8 @@ package com.nhn.pinpoint.rpc.client; +import java.util.ArrayList; +import java.util.List; + import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -21,7 +24,8 @@ public class PinpointSocket { private volatile boolean closed; - private volatile PinpointSocketReconnectEventListener reconnectEventListener; + private List reconnectEventListeners = new ArrayList(); + public PinpointSocket(SocketHandler socketHandler) { if (socketHandler == null) { @@ -29,12 +33,10 @@ public class PinpointSocket { } this.socketHandler = socketHandler; socketHandler.setPinpointSocket(this); - this.reconnectEventListener = new DummyPinpointSocketReconnectEventListener(); } public PinpointSocket() { this.socketHandler = new ReconnectStateSocketHandler(); - this.reconnectEventListener = new DummyPinpointSocketReconnectEventListener(); } void reconnectSocketHandler(SocketHandler socketHandler) { @@ -48,24 +50,50 @@ public class PinpointSocket { } logger.warn("reconnectSocketHandler:{}", socketHandler); this.socketHandler = socketHandler; - getPinpointSocketReconnectEventListener().reconnectPerformed(this); + + notifyReconnectEvent(); } - // reconnectEventListener의 경우 직접 생성자 호출시에 Dummy를 포함하고 있으며, + // reconnectEventListener의 경우 직접 생성자 호출시에 Dummy를 포함하고 있으며, // setter를 통해서도 접근을 못하게 하기 때문에 null이 아닌 것이 보장됨 - public boolean setPinpointSocketReconnectEventListener(PinpointSocketReconnectEventListener reconnectEventListener) { - if (reconnectEventListener == null) { + public boolean addPinpointSocketReconnectEventListener(PinpointSocketReconnectEventListener eventListener) { + if (eventListener == null) { return false; } - this.reconnectEventListener = reconnectEventListener; - return true; + synchronized (this) { + return this.reconnectEventListeners.add(eventListener); + } } - - private PinpointSocketReconnectEventListener getPinpointSocketReconnectEventListener() { - return reconnectEventListener; + + public boolean removePinpointSocketReconnectEventListener(PinpointSocketReconnectEventListener eventListener) { + if (eventListener == null) { + return false; + } + synchronized (this) { + return this.reconnectEventListeners.remove(eventListener); + } } + private List getPinpointSocketReconnectEventListener() { + List result = new ArrayList(); + synchronized (this) { + for (PinpointSocketReconnectEventListener eventListener : this.reconnectEventListeners) { + result.add(eventListener); + } + } + + return result; + } + + private void notifyReconnectEvent() { + List reconnectEventListeners = getPinpointSocketReconnectEventListener(); + + for (PinpointSocketReconnectEventListener eachListener : reconnectEventListeners) { + eachListener.reconnectPerformed(this); + } + } + public void sendSync(byte[] bytes) { ensureOpen(); socketHandler.sendSync(bytes); diff --git a/src/test/java/com/nhn/pinpoint/rpc/client/ReconnectTest.java b/src/test/java/com/nhn/pinpoint/rpc/client/ReconnectTest.java index 6d467db18..6e36a5676 100644 --- a/src/test/java/com/nhn/pinpoint/rpc/client/ReconnectTest.java +++ b/src/test/java/com/nhn/pinpoint/rpc/client/ReconnectTest.java @@ -6,12 +6,14 @@ import com.nhn.pinpoint.rpc.ResponseMessage; import com.nhn.pinpoint.rpc.TestByteUtils; import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.TestSeverMessageListener; + import org.junit.Assert; import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; import java.io.IOException; +import java.util.concurrent.atomic.AtomicBoolean; /** @@ -30,6 +32,7 @@ public class ReconnectTest { serverSocket.setMessageListener(new TestSeverMessageListener()); serverSocket.bind("localhost", PORT); + final AtomicBoolean reconnectPerformed = new AtomicBoolean(false); final PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); pinpointSocketFactory.setReconnectDelay(200); @@ -37,6 +40,15 @@ public class ReconnectTest { PinpointServerSocket newServerSocket = null; try { PinpointSocket socket = pinpointSocketFactory.connect("localhost", 10234); + socket.addPinpointSocketReconnectEventListener(new PinpointSocketReconnectEventListener() { + + @Override + public void reconnectPerformed(PinpointSocket socket) { + reconnectPerformed.set(true); + } + + }); + serverSocket.close(); logger.info("server.close()---------------------------"); Thread.sleep(1000); @@ -68,7 +80,8 @@ public class ReconnectTest { } pinpointSocketFactory.release(); } - + + Assert.assertTrue(reconnectPerformed.get()); } @Test