mirror of
https://github.com/wahyd4/pinpoint.git
synced 2026-08-17 08:46:22 +10:00
[구태진] [pinpoint-rpc-2] ReconnectEventListener 생성 및 PinpointSocket에 저장함
1. Event를 여러개 등록할수 있게끔 변경 git-svn-id: http://svn.bds.nhncorp.com/pe/pinpoint-rpc/trunk@3615 84d0f5b1-2673-498c-a247-62c4ff18d310
This commit is contained in:
@@ -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<PinpointSocketReconnectEventListener> reconnectEventListeners = new ArrayList<PinpointSocketReconnectEventListener>();
|
||||
|
||||
|
||||
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<PinpointSocketReconnectEventListener> getPinpointSocketReconnectEventListener() {
|
||||
List<PinpointSocketReconnectEventListener> result = new ArrayList<PinpointSocketReconnectEventListener>();
|
||||
synchronized (this) {
|
||||
for (PinpointSocketReconnectEventListener eventListener : this.reconnectEventListeners) {
|
||||
result.add(eventListener);
|
||||
}
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
private void notifyReconnectEvent() {
|
||||
List<PinpointSocketReconnectEventListener> reconnectEventListeners = getPinpointSocketReconnectEventListener();
|
||||
|
||||
for (PinpointSocketReconnectEventListener eachListener : reconnectEventListeners) {
|
||||
eachListener.reconnectPerformed(this);
|
||||
}
|
||||
}
|
||||
|
||||
public void sendSync(byte[] bytes) {
|
||||
ensureOpen();
|
||||
socketHandler.sendSync(bytes);
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user