#22 Cluster간 Stream 기능 고안 및 개발

1. 이전에 사용하던 StreamChannel 관련 코드 제거 및 신규 StreamChannel로 변경
2. MessageListener에서 StreamChnnel 관련 처리를 분리 
3. PinpointSocketFactory의 인터페이스 변경 
4. 단순 테스트 코드 수정
This commit is contained in:
koo-taejin
2014-10-17 11:56:31 +09:00
parent a37ae06eba
commit 284bc433d1
37 changed files with 511 additions and 929 deletions
@@ -22,16 +22,12 @@ public class WebClusterPoint implements ClusterPoint {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private final PinpointSocketFactory factory;
private final MessageListener messageListener;
// InetSocketAddress List로 전달 하는게 좋을거 같은데 이걸 Key로만들기가 쉽지 않네;
private final Map<InetSocketAddress, PinpointSocket> clusterRepository = new HashMap<InetSocketAddress, PinpointSocket>();
public WebClusterPoint(String id, MessageListener messageListener) {
this.messageListener = messageListener;
this.factory = new PinpointSocketFactory();
this.factory.setTimeoutMillis(1000 * 5);
this.factory.setMessageListener(messageListener);
Map<String, Object> properties = new HashMap<String, Object>();
properties.put("id", id);
@@ -74,7 +70,7 @@ public class WebClusterPoint implements ClusterPoint {
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
logger.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -82,7 +78,7 @@ public class WebClusterPoint implements ClusterPoint {
}
}
logger.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -29,10 +29,8 @@ import com.nhn.pinpoint.common.util.PinpointThreadFactory;
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.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
import com.nhn.pinpoint.thrift.io.DeserializerFactory;
import com.nhn.pinpoint.thrift.io.Header;
@@ -136,11 +134,6 @@ public class TCPReceiver {
requestResponse(requestPacket, channel);
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.warn("unsupported streamPacket received {}", streamPacket);
}
@Override
public int handleEnableWorker(Map properties) {
if (properties == null) {
@@ -18,11 +18,9 @@ import com.nhn.pinpoint.collector.receiver.tcp.AgentProperties;
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.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.ChannelContext;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
@RunWith(SpringJUnit4ClassRunner.class)
@@ -88,11 +86,6 @@ public class ClusterPointRouterTest {
logger.warn("Unsupport request received {} {}", requestPacket, channel);
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.warn("unsupported streamPacket received {}", streamPacket);
}
@Override
public int handleEnableWorker(Map properties) {
logger.warn("do handleEnableWorker {}", properties);
@@ -216,6 +216,7 @@ public class ZookeeperProfilerClusterStressTest {
this.factory = new PinpointSocketFactory();
this.factory.setProperties(properties);
this.factory.setMessageListener(messageListener);
}
private void connect(InetSocketAddress address) {
@@ -249,7 +250,7 @@ public class ZookeeperProfilerClusterStressTest {
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
logger.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -257,7 +258,7 @@ public class ZookeeperProfilerClusterStressTest {
}
}
logger.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -13,9 +13,7 @@ import com.nhn.pinpoint.rpc.client.MessageListener;
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.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
final class ZookeeperTestUtils {
@@ -86,11 +84,6 @@ final class ZookeeperTestUtils {
LOGGER.warn("Unsupport request received {} {}", requestPacket, channel);
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
LOGGER.warn("unsupported streamPacket received {}", streamPacket);
}
@Override
public int handleEnableWorker(Map properties) {
LOGGER.warn("do handleEnableWorker {}", properties);
@@ -137,9 +137,8 @@ public class DefaultAgent implements Agent {
this.tAgentInfo = createTAgentInfo();
this.factory = createPinpointSocketFactory();
this.socket = createPinpointSocket(this.profilerConfig.getCollectorServerIp(), this.profilerConfig.getCollectorTcpServerPort(), factory,
this.profilerConfig.isTcpDataSenderCommandAcceptEnable());
this.factory = createPinpointSocketFactory(this.profilerConfig.isTcpDataSenderCommandAcceptEnable());
this.socket = createPinpointSocket(this.profilerConfig.getCollectorServerIp(), this.profilerConfig.getCollectorTcpServerPort(), factory);
this.tcpDataSender = createTcpDataSender(socket);
@@ -285,7 +284,7 @@ public class DefaultAgent implements Agent {
return serverMetaDataHolder;
}
protected PinpointSocketFactory createPinpointSocketFactory() {
protected PinpointSocketFactory createPinpointSocketFactory(boolean isSupportServerMode) {
Map<String, Object> properties = this.agentInformation.toMap();
properties.put(AgentPropertiesType.IP.getName(), serverInfo.getHostip());
@@ -293,28 +292,22 @@ public class DefaultAgent implements Agent {
pinpointSocketFactory.setTimeoutMillis(1000 * 5);
pinpointSocketFactory.setProperties(properties);
if (isSupportServerMode) {
pinpointSocketFactory.setMessageListener(new CommandDispatcher());
}
return pinpointSocketFactory;
}
protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) {
return createPinpointSocket(host, port, factory, false);
}
protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory, boolean useMessageListener) {
// 1.2 버전이 Tcp Data Command 허용하는 버전이 아니기 떄문에 true이던 false이던 무조건 SimpleLoggingMessageListener를 이용하게 함
// SimpleLoggingMessageListener.LISTENER 는 서로 통신을 하지 않게 설정되어 있음 (테스트코드는 pinpoint-rpc에 존재)
// 1.3 버전으로 할 경우 아래 분기에서 MessageListener 변경 필요
MessageListener messageListener = null;
if (useMessageListener) {
messageListener = new CommandDispatcher();
} else {
messageListener = SimpleLoggingMessageListener.LISTENER;
}
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
logger.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -322,7 +315,7 @@ public class DefaultAgent implements Agent {
}
}
logger.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -93,18 +93,17 @@ public class NetworkAvailabilityChecker implements PinpointTools {
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
pinpointSocketFactory.setTimeoutMillis(1000 * 5);
pinpointSocketFactory.setProperties(Collections.<String, Object>emptyMap());
pinpointSocketFactory.setMessageListener(new CommandDispatcher());
return pinpointSocketFactory;
}
private static PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) {
MessageListener messageListener = new CommandDispatcher();
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
LOGGER.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -112,7 +111,7 @@ public class NetworkAvailabilityChecker implements PinpointTools {
}
}
LOGGER.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -15,15 +15,12 @@ import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.profiler.receiver.CommandDispatcher;
import com.nhn.pinpoint.profiler.sender.TcpDataSender;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.client.MessageListener;
import com.nhn.pinpoint.rpc.client.PinpointSocket;
import com.nhn.pinpoint.rpc.client.PinpointSocketFactory;
import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
import com.nhn.pinpoint.thrift.dto.TAgentInfo;
import com.nhn.pinpoint.thrift.dto.TResult;
@@ -210,11 +207,6 @@ public class HeartBeatCheckerTest {
}
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.info("handleStreamPacket:{}", streamPacket);
}
@Override
public int handleEnableWorker(Map arg0) {
return 0;
@@ -225,18 +217,17 @@ public class HeartBeatCheckerTest {
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
pinpointSocketFactory.setTimeoutMillis(1000 * 5);
pinpointSocketFactory.setProperties(Collections.EMPTY_MAP);
pinpointSocketFactory.setMessageListener(new CommandDispatcher());
return pinpointSocketFactory;
}
private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) {
MessageListener messageListener = new CommandDispatcher();
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
logger.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -244,7 +235,7 @@ public class HeartBeatCheckerTest {
}
}
logger.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -12,15 +12,12 @@ import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.profiler.receiver.CommandDispatcher;
import com.nhn.pinpoint.profiler.sender.TcpDataSender;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.client.MessageListener;
import com.nhn.pinpoint.rpc.client.PinpointSocket;
import com.nhn.pinpoint.rpc.client.PinpointSocketFactory;
import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
import com.nhn.pinpoint.thrift.dto.TAgentInfo;
import com.nhn.pinpoint.thrift.dto.TResult;
@@ -160,11 +157,6 @@ public class HeartBitCheckerStressTest {
e.printStackTrace();
}
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.info("handleStreamPacket:{}", streamPacket);
}
@Override
public int handleEnableWorker(Map arg0) {
@@ -176,18 +168,17 @@ public class HeartBitCheckerStressTest {
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
pinpointSocketFactory.setTimeoutMillis(1000 * 5);
pinpointSocketFactory.setProperties(Collections.EMPTY_MAP);
pinpointSocketFactory.setMessageListener(new CommandDispatcher());
return pinpointSocketFactory;
}
private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) {
MessageListener messageListener = new CommandDispatcher();
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
logger.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -195,7 +186,7 @@ public class HeartBitCheckerStressTest {
}
}
logger.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -9,15 +9,12 @@ import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.profiler.receiver.CommandDispatcher;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.client.MessageListener;
import com.nhn.pinpoint.rpc.client.PinpointSocket;
import com.nhn.pinpoint.rpc.client.PinpointSocketFactory;
import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
import com.nhn.pinpoint.thrift.dto.TApiMetaData;
@@ -48,11 +45,6 @@ public class TcpDataSenderReconnectTest {
logger.info("handleRequest:{}", requestPacket);
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.info("handleStreamPacket:{}", streamPacket);
}
@Override
public int handleEnableWorker(Map properties) {
return 0;
@@ -95,18 +87,17 @@ public class TcpDataSenderReconnectTest {
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
pinpointSocketFactory.setTimeoutMillis(1000 * 5);
pinpointSocketFactory.setProperties(Collections.EMPTY_MAP);
pinpointSocketFactory.setMessageListener(new CommandDispatcher());
return pinpointSocketFactory;
}
private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) {
MessageListener messageListener = new CommandDispatcher();
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
logger.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -114,7 +105,7 @@ public class TcpDataSenderReconnectTest {
}
}
logger.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -15,15 +15,12 @@ import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.profiler.receiver.CommandDispatcher;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.client.MessageListener;
import com.nhn.pinpoint.rpc.client.PinpointSocket;
import com.nhn.pinpoint.rpc.client.PinpointSocketFactory;
import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
import com.nhn.pinpoint.thrift.dto.TApiMetaData;
@@ -57,11 +54,6 @@ public class TcpDataSenderTest {
public void handleRequest(RequestPacket requestPacket, SocketChannel channel) {
logger.info("handleRequest:{}", requestPacket);
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.info("handleStreamPacket:{}", streamPacket);
}
@Override
public int handleEnableWorker(Map arg0) {
@@ -83,6 +75,8 @@ public class TcpDataSenderTest {
this.sendLatch = new CountDownLatch(2);
PinpointSocketFactory socketFactory = createPinpointSocketFactory();
socketFactory.setMessageListener(new CommandDispatcher());
PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory);
TcpDataSender sender = new TcpDataSender(socket);
@@ -113,15 +107,12 @@ public class TcpDataSenderTest {
return pinpointSocketFactory;
}
private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) {
MessageListener messageListener = new CommandDispatcher();
PinpointSocket socket = null;
for (int i = 0; i < 3; i++) {
try {
socket = factory.connect(host, port, messageListener);
socket = factory.connect(host, port);
logger.info("tcp connect success:{}/{}", host, port);
return socket;
} catch (PinpointSocketException e) {
@@ -129,7 +120,7 @@ public class TcpDataSenderTest {
}
}
logger.warn("change background tcp connect mode {}/{} ", host, port);
socket = factory.scheduledConnect(host, port, messageListener);
socket = factory.scheduledConnect(host, port);
return socket;
}
@@ -64,21 +64,11 @@ public class MockAgent extends DefaultAgent {
return new HoldingSpanStorageFactory(getSpanDataSender());
}
@Override
protected PinpointSocketFactory createPinpointSocketFactory() {
return null;
}
@Override
protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) {
return null;
}
@Override
protected PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory, boolean useMessageListener) {
return null;
}
@Override
protected EnhancedDataSender createTcpDataSender(PinpointSocket socket) {
return new LoggingDataSender();
@@ -8,9 +8,7 @@ import org.slf4j.LoggerFactory;
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.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
/**
@@ -34,12 +32,6 @@ public class RequestResponseServerMessageListener implements ServerMessageListen
channel.sendResponseMessage(requestPacket, requestPacket.getPayload());
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.info("handlerStream {} {}", streamChannel, streamChannel);
}
@Override
public int handleEnableWorker(Map properties) {
logger.info("handleEnableWorker {}", properties);
@@ -10,6 +10,8 @@ import com.nhn.pinpoint.rpc.DefaultFuture;
import com.nhn.pinpoint.rpc.Future;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.ResponseMessage;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.util.AssertUtils;
@@ -21,7 +23,6 @@ import com.nhn.pinpoint.rpc.util.AssertUtils;
public class PinpointSocket {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private final MessageListener messageListener;
private volatile SocketHandler socketHandler;
@@ -32,28 +33,19 @@ public class PinpointSocket {
public PinpointSocket() {
this(new ReconnectStateSocketHandler());
}
public PinpointSocket(MessageListener messageListener) {
this(new ReconnectStateSocketHandler(), messageListener);
}
public PinpointSocket(SocketHandler socketHandler) {
this(socketHandler, SimpleLoggingMessageListener.LISTENER);
}
public PinpointSocket(SocketHandler socketHandler, MessageListener messageListener) {
AssertUtils.assertNotNull(socketHandler, "socketHandler");
AssertUtils.assertNotNull(messageListener, "messageListener");
this.messageListener = messageListener;
socketHandler.setMessageListener(this.messageListener);
if (socketHandler.isSupportServerMode()) {
socketHandler.turnOnServerMode();
}
this.socketHandler = socketHandler;
socketHandler.setPinpointSocket(this);
}
void reconnectSocketHandler(SocketHandler socketHandler) {
AssertUtils.assertNotNull(socketHandler, "socketHandler");
@@ -64,8 +56,11 @@ public class PinpointSocket {
}
logger.warn("reconnectSocketHandler:{}", socketHandler);
// Pinpoint 소켓 내부 객체가 되기전에 listener를 먼저 등록
socketHandler.setMessageListener(messageListener);
// Pinpoint 소켓 내부 객체가 되기전에 listener를 먼저 등록
if (socketHandler.isSupportServerMode()) {
socketHandler.turnOnServerMode();
}
this.socketHandler = socketHandler;
notifyReconnectEvent();
@@ -118,12 +113,11 @@ public class PinpointSocket {
return socketHandler.request(bytes);
}
public StreamChannel createStreamChannel() {
public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) {
// 실패를 리턴하는 StreamChannel을 던져야 되는데. StreamChannel을 interface로 변경해야 됨.
// 일단 그냥 ex를 던지도록 하겠음.
ensureOpen();
return socketHandler.createStreamChannel();
return socketHandler.createStreamChannel(payload, clientStreamChannelMessageListener);
}
private Future<ResponseMessage> returnFailureFuture() {
@@ -29,6 +29,8 @@ import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.common.util.PinpointThreadFactory;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.stream.DisabledServerStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.ServerStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.util.AssertUtils;
import com.nhn.pinpoint.rpc.util.LoggerFactorySetup;
import com.nhn.pinpoint.rpc.util.TimerFactory;
@@ -61,6 +63,8 @@ public class PinpointSocketFactory {
private long enableWorkerPacketDelay = DEFAULT_ENABLE_WORKER_PACKET_DELAY;
private long timeoutMillis = DEFAULT_TIMEOUTMILLIS;
private MessageListener messageListener = SimpleLoggingMessageListener.LISTENER;
private ServerStreamChannelMessageListener serverStreamChannelMessageListener = DisabledServerStreamChannelMessageListener.INSTANCE;
static {
LoggerFactorySetup.setupSlf4jLoggerFactory();
@@ -181,33 +185,21 @@ public class PinpointSocketFactory {
}
public PinpointSocket connect(String host, int port) throws PinpointSocketException {
return connect(host, port, SimpleLoggingMessageListener.LISTENER);
}
public PinpointSocket connect(String host, int port, MessageListener messageListener) throws PinpointSocketException {
AssertUtils.assertNotNull(messageListener);
SocketAddress address = new InetSocketAddress(host, port);
ChannelFuture connectFuture = bootstrap.connect(address);
SocketHandler socketHandler = getSocketHandler(connectFuture, address);
PinpointSocket pinpointSocket = new PinpointSocket(socketHandler, messageListener);
PinpointSocket pinpointSocket = new PinpointSocket(socketHandler);
traceSocket(pinpointSocket);
return pinpointSocket;
}
public PinpointSocket reconnect(String host, int port) throws PinpointSocketException {
return reconnect(host, port, SimpleLoggingMessageListener.LISTENER);
}
public PinpointSocket reconnect(String host, int port, MessageListener messageListener) throws PinpointSocketException {
AssertUtils.assertNotNull(messageListener);
SocketAddress address = new InetSocketAddress(host, port);
ChannelFuture connectFuture = bootstrap.connect(address);
SocketHandler socketHandler = getSocketHandler(connectFuture, address);
PinpointSocket pinpointSocket = new PinpointSocket(socketHandler, messageListener);
PinpointSocket pinpointSocket = new PinpointSocket(socketHandler);
traceSocket(pinpointSocket);
return pinpointSocket;
}
@@ -218,13 +210,7 @@ public class PinpointSocketFactory {
}
public PinpointSocket scheduledConnect(String host, int port) {
return scheduledConnect(host, port, SimpleLoggingMessageListener.LISTENER);
}
public PinpointSocket scheduledConnect(String host, int port, MessageListener messageListener) {
AssertUtils.assertNotNull(messageListener);
PinpointSocket pinpointSocket = new PinpointSocket(new ReconnectStateSocketHandler(), messageListener);
PinpointSocket pinpointSocket = new PinpointSocket(new ReconnectStateSocketHandler());
SocketAddress address = new InetSocketAddress(host, port);
reconnect(pinpointSocket, address);
return pinpointSocket;
@@ -383,4 +369,24 @@ public class PinpointSocketFactory {
this.properties = Collections.unmodifiableMap(agentProperties);
}
public MessageListener getMessageListener() {
return messageListener;
}
public void setMessageListener(MessageListener messageListener) {
AssertUtils.assertNotNull(messageListener, "messageListener must not be null");
this.messageListener = messageListener;
}
public ServerStreamChannelMessageListener getServerStreamChannelMessageListener() {
return serverStreamChannelMessageListener;
}
public void setServerStreamChannelMessageListener(ServerStreamChannelMessageListener serverStreamChannelMessageListener) {
AssertUtils.assertNotNull(messageListener, "messageListener must not be null");
this.serverStreamChannelMessageListener = serverStreamChannelMessageListener;
}
}
@@ -37,8 +37,13 @@ import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.ResponsePacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.util.AssertUtils;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.DisabledServerStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.ServerStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.StreamChannelManager;
import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils;
import com.nhn.pinpoint.rpc.util.IDGenerator;
import com.nhn.pinpoint.rpc.util.MapUtils;
import com.nhn.pinpoint.rpc.util.TimerFactory;
@@ -60,7 +65,6 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
private final State state = new State();
private volatile Channel channel;
private volatile MessageListener messageListener = SimpleLoggingMessageListener.LISTENER;
private long timeoutMillis = DEFAULT_TIMEOUTMILLIS;
private long pingDelay = DEFAULT_PING_DELAY;
@@ -74,8 +78,10 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
private SocketAddress connectSocketAddress;
private volatile PinpointSocket pinpointSocket;
private final MessageListener messageListener;
private final ServerStreamChannelMessageListener serverStreamChannelMessageListener;
private final RequestManager requestManager;
private final StreamChannelManager streamChannelManager;
private final ChannelFutureListener pingWriteFailFutureListener = new WriteFailFutureListener(this.logger, "ping write fail.", "ping write success.");
private final ChannelFutureListener sendWriteFailFutureListener = new WriteFailFutureListener(this.logger, "send() write fail.", "send() write fail.");
@@ -95,10 +101,25 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
this.channelTimer = timer;
this.pinpointSocketFactory = pinpointSocketFactory;
this.requestManager = new RequestManager(timer);
this.streamChannelManager = new StreamChannelManager();
this.pingDelay = pingDelay;
this.enableWorkerPacketDelay = enableWorkerPacketDelay;
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;
}
pinpointSocketFactory.getServerStreamChannelMessageListener();
}
public Timer getChannelTimer() {
@@ -134,24 +155,25 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
if (!state.changeRun()) {
throw new IllegalStateException("invalid open state:" + state.getString());
}
Channel channel = this.channel;
if (channel != null) {
prepareChannel(channel);
}
}
@Override
public void setMessageListener(MessageListener messageListener) {
AssertUtils.assertNotNull(messageListener, "messageListener");
logger.info("{} registered Listner({}).", toString(), messageListener);
if (messageListener != SimpleLoggingMessageListener.LISTENER) {
this.messageListener = messageListener;
// MessageListener 등록시 EnableWorkerPacket전달
sendEnableWorkerPacket();
RegisterEnableWorkerPacketJob job = new RegisterEnableWorkerPacketJob(enableWorkerPacketRetryCount);
reservationEnableWorkerPacketJob(job);
}
}
private void prepareChannel(Channel channel) {
ServerStreamChannelMessageListener serverStreamChannelMessageListener = this.serverStreamChannelMessageListener;
StreamChannelManager streamChannelManager = new StreamChannelManager(channel, IDGenerator.createOddIdGenerator(), serverStreamChannelMessageListener);
SocketHandlerContext context = new SocketHandlerContext(channel, streamChannelManager);
channel.setAttachment(context);
}
private SocketHandlerContext getChannelContext(Channel channel) {
return (SocketHandlerContext) channel.getAttachment();
}
@Override
public void initReconnect() {
@@ -373,13 +395,14 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
return messageFuture;
}
public StreamChannel createStreamChannel() {
@Override
public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) {
ensureOpen();
final Channel channel = this.channel;
return this.streamChannelManager.createStreamChannel(channel);
SocketHandlerContext context = getChannelContext(channel);
return context.getStreamChannelManager().openStreamChannel(payload, clientStreamChannelMessageListener);
}
@@ -405,7 +428,10 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
case PacketType.APPLICATION_STREAM_CREATE_SUCCESS:
case PacketType.APPLICATION_STREAM_CREATE_FAIL:
case PacketType.APPLICATION_STREAM_RESPONSE:
this.streamChannelManager.messageReceived((StreamPacket) message, e.getChannel());
case PacketType.APPLICATION_STREAM_PING:
case PacketType.APPLICATION_STREAM_PONG:
SocketHandlerContext context = getChannelContext(channel);
context.getStreamChannelManager().messageReceived((StreamPacket) message);
return;
case PacketType.CONTROL_SERVER_CLOSE:
messageReceivedServerClosed(e.getChannel());
@@ -550,7 +576,12 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
private void releaseResource() {
logger.debug("releaseResource()");
this.requestManager.close();
this.streamChannelManager.close();
if (this.channel != null) {
SocketHandlerContext context = getChannelContext(channel);
context.getStreamChannelManager().close();
}
this.channelTimer.stop();
}
@@ -588,5 +619,37 @@ public class PinpointSocketHandler extends SimpleChannelHandler implements Socke
public boolean isConnected() {
return this.state.isRun();
}
@Override
public boolean isSupportServerMode() {
return messageListener != SimpleLoggingMessageListener.LISTENER;
}
@Override
public void turnOnServerMode() {
// MessageListener 등록시 EnableWorkerPacket전달
sendEnableWorkerPacket();
RegisterEnableWorkerPacketJob job = new RegisterEnableWorkerPacketJob(enableWorkerPacketRetryCount);
reservationEnableWorkerPacketJob(job);
}
class SocketHandlerContext {
private final Channel channel;
private final StreamChannelManager streamChannelManager;
public SocketHandlerContext(Channel channel, StreamChannelManager streamChannelManager) {
this.channel = channel;
this.streamChannelManager = streamChannelManager;
}
public Channel getChannel() {
return channel;
}
public StreamChannelManager getStreamChannelManager() {
return streamChannelManager;
}
}
}
@@ -4,6 +4,8 @@ import com.nhn.pinpoint.rpc.DefaultFuture;
import com.nhn.pinpoint.rpc.Future;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.ResponseMessage;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
import java.net.SocketAddress;
@@ -22,10 +24,6 @@ public class ReconnectStateSocketHandler implements SocketHandler {
public void open() {
throw new IllegalStateException();
}
@Override
public void setMessageListener(MessageListener messageListener) {
}
@Override
public void initReconnect() {
@@ -70,10 +68,10 @@ public class ReconnectStateSocketHandler implements SocketHandler {
}
@Override
public StreamChannel createStreamChannel() {
throw new UnsupportedOperationException();
public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) {
throw new UnsupportedOperationException();
}
@Override
public void sendPing() {
}
@@ -82,4 +80,15 @@ public class ReconnectStateSocketHandler implements SocketHandler {
public boolean isConnected() {
return false;
}
@Override
public boolean isSupportServerMode() {
return false;
}
@Override
public void turnOnServerMode() {
throw new UnsupportedOperationException();
}
}
@@ -31,12 +31,15 @@ public class SocketClientPipelineFactory implements ChannelPipelineFactory {
ChannelPipeline pipeline = Channels.pipeline();
pipeline.addLast("encoder", new PacketEncoder());
pipeline.addLast("decoder", new PacketDecoder());
long pingDelay = pinpointSocketFactory.getPingDelay();
long enableWorkerPacketDelay = pinpointSocketFactory.getEnableWorkerPacketDelay();
long timeoutMillis = pinpointSocketFactory.getTimeoutMillis();
PinpointSocketHandler pinpointSocketHandler = new PinpointSocketHandler(pinpointSocketFactory, pingDelay, enableWorkerPacketDelay, timeoutMillis);
pipeline.addLast("writeTimeout", new WriteTimeoutHandler(pinpointSocketHandler.getChannelTimer(), 3000, TimeUnit.MILLISECONDS));
pipeline.addLast("socketHandler", pinpointSocketHandler);
return pipeline;
}
}
@@ -1,9 +1,11 @@
package com.nhn.pinpoint.rpc.client;
import java.net.SocketAddress;
import com.nhn.pinpoint.rpc.Future;
import com.nhn.pinpoint.rpc.ResponseMessage;
import java.net.SocketAddress;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
/**
* @author emeroad
@@ -29,11 +31,14 @@ public interface SocketHandler {
Future<ResponseMessage> request(byte[] bytes);
StreamChannel createStreamChannel();
ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener);
void sendPing();
boolean isConnected();
void setMessageListener(MessageListener messageListener);
boolean isSupportServerMode();
void turnOnServerMode();
}
@@ -1,225 +0,0 @@
package com.nhn.pinpoint.rpc.client;
import java.util.concurrent.atomic.AtomicInteger;
import org.jboss.netty.channel.Channel;
import org.jboss.netty.channel.ChannelFuture;
import org.jboss.netty.channel.ChannelFutureListener;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.rpc.DefaultFuture;
import com.nhn.pinpoint.rpc.FailureEventHandler;
import com.nhn.pinpoint.rpc.Future;
import com.nhn.pinpoint.rpc.StreamCreateResponse;
import com.nhn.pinpoint.rpc.packet.PacketType;
import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.stream.StreamChannelMessageListener;
/**
* @author emeroad
*/
public class StreamChannel {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private static final int NONE = 0;
// OPEN 호출
private static final int OPEN = 1;
// OPEN 결과 대기
private static final int OPEN_AWAIT = 2;
// 동작중
private static final int RUN = 3;
// 닫힘
private static final int CLOSED = 4;
private final AtomicInteger state = new AtomicInteger(NONE);
private final int channelId;
private StreamChannelManager streamChannelManager;
private StreamChannelMessageListener streamChannelMessageListener;
private DefaultFuture<StreamCreateResponse> openLatch;
private Channel channel;
public StreamChannel(int channelId) {
this.channelId = channelId;
}
public int getChannelId() {
return channelId;
}
public void setChannel(Channel channel) {
this.channel = channel;
}
public Future<StreamCreateResponse> open(byte[] bytes) {
if (!state.compareAndSet(NONE, OPEN)) {
throw new IllegalStateException("invalid state:" + state.get());
}
StreamCreatePacket streamCreatePacket = new StreamCreatePacket(channelId, bytes);
this.openLatch = new DefaultFuture<StreamCreateResponse>();
openLatch.setFailureEventHandler(new FailureEventHandler() {
@Override
public boolean fireFailure() {
streamChannelManager.closeChannel(channelId);
return false;
}
});
ChannelFuture channelFuture = this.channel.write(streamCreatePacket);
channelFuture.addListener(new ChannelFutureListener() {
@Override
public void operationComplete(ChannelFuture future) throws Exception {
if (!future.isSuccess()) {
future.setFailure(future.getCause());
}
}
});
if (!state.compareAndSet(OPEN, OPEN_AWAIT)) {
throw new IllegalStateException("invalid state");
}
return openLatch;
}
public boolean receiveStreamPacket(StreamPacket packet) {
final short packetType = packet.getPacketType();
switch (packetType) {
case PacketType.APPLICATION_STREAM_CREATE_SUCCESS:
logger.debug("APPLICATION_STREAM_CREATE_SUCCESS {}", channel);
StreamCreateResponse success = new StreamCreateResponse(true);
success.setMessage(packet.getPayload());
return openChannel(RUN, success);
case PacketType.APPLICATION_STREAM_CREATE_FAIL:
logger.debug("APPLICATION_STREAM_CREATE_FAIL {}", channel);
StreamCreateResponse failResult = new StreamCreateResponse(false);
failResult.setMessage(packet.getPayload());
return openChannel(CLOSED, failResult);
case PacketType.APPLICATION_STREAM_RESPONSE: {
logger.debug("APPLICATION_STREAM_RESPONSE {}", channel);
StreamResponsePacket streamResponsePacket = (StreamResponsePacket) packet;
StreamChannelMessageListener streamChannelMessageListener = this.streamChannelMessageListener;
if (streamChannelMessageListener != null) {
streamChannelMessageListener.handleStreamData(this, streamResponsePacket);
}
return true;
}
case PacketType.APPLICATION_STREAM_CLOSE: {
logger.debug("APPLICATION_STREAM_CLOSE {}", channel);
this.closeInternal();
StreamClosePacket streamClosePacket = (StreamClosePacket) packet;
StreamChannelMessageListener streamChannelMessageListener = this.streamChannelMessageListener;
if (streamChannelMessageListener != null) {
streamChannelMessageListener.handleStreamClose(this, streamClosePacket);
}
return true;
}
}
return false;
}
private boolean openChannel(int channelState, StreamCreateResponse streamCreateResponse) {
if (state.compareAndSet(OPEN_AWAIT, channelState)) {
notifyOpenResult(streamCreateResponse);
return true;
} else {
logger.info("invalid stream channel state:{}", state.get());
return false;
}
}
private boolean notifyOpenResult(StreamCreateResponse failResult) {
DefaultFuture<StreamCreateResponse> openLatch = this.openLatch;
if (openLatch != null) {
return openLatch.setResult(failResult);
}
return false;
}
public boolean close() {
return close0(true);
}
boolean closeInternal() {
return close0(false);
}
private boolean close0(boolean safeClose) {
if (!state.compareAndSet(RUN, CLOSED)) {
return false;
}
if (safeClose) {
StreamClosePacket closePacket = new StreamClosePacket(this.channelId, StreamClosePacket.SUCCESS);
this.channel.write(closePacket);
StreamChannelManager streamChannelManager = this.streamChannelManager;
if (streamChannelManager != null) {
streamChannelManager.closeChannel(channelId);
this.streamChannelManager = null;
}
}
return true;
}
public void setStreamChannelManager(StreamChannelManager streamChannelManager) {
this.streamChannelManager = streamChannelManager;
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
StreamChannel that = (StreamChannel) o;
if (channelId != that.channelId) return false;
if (channel != null ? !channel.equals(that.channel) : that.channel != null) return false;
return true;
}
@Override
public int hashCode() {
int result = channelId;
result = 31 * result + (channel != null ? channel.hashCode() : 0);
return result;
}
public void setStreamChannelMessageListener(StreamChannelMessageListener streamChannelMessageListener) {
this.streamChannelMessageListener = streamChannelMessageListener;
}
@Override
public String toString() {
final StringBuilder sb = new StringBuilder();
sb.append("StreamChannel");
sb.append("{channelId=").append(channelId);
sb.append(", channel=").append(channel);
sb.append('}');
return sb.toString();
}
}
@@ -1,81 +0,0 @@
package com.nhn.pinpoint.rpc.client;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.atomic.AtomicInteger;
import org.jboss.netty.channel.Channel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
/**
* @author emeroad
*/
public class StreamChannelManager {
private Logger logger = LoggerFactory.getLogger(this.getClass());
private final AtomicInteger idAllocator = new AtomicInteger(0);
private final ConcurrentMap<Integer, StreamChannel> channelMap = new ConcurrentHashMap<Integer, StreamChannel>();
public StreamChannel createStreamChannel(Channel channel) {
final int channelId = allocateChannelId();
StreamChannel streamChannel = new StreamChannel(channelId);
streamChannel.setChannel(channel);
StreamChannel old = channelMap.put(channelId, streamChannel);
if (old != null) {
throw new PinpointSocketException("already channelId exist:" + channelId + " streamChannel:" + old);
}
// handle을 붙여서 리턴.
streamChannel.setStreamChannelManager(this);
return streamChannel;
}
private int allocateChannelId() {
return idAllocator.get();
}
public StreamChannel findStreamChannel(int channelId) {
return this.channelMap.get(channelId);
}
public boolean closeChannel(int channelId) {
StreamChannel remove = this.channelMap.remove(channelId);
return remove != null;
}
public void close() {
logger.debug("close()");
final ConcurrentMap<Integer, StreamChannel> channelMap = this.channelMap;
int forceCloseChannel = 0;
for (Map.Entry<Integer, StreamChannel> entry : channelMap.entrySet()) {
if(entry.getValue().closeInternal()) {
forceCloseChannel++;
}
}
channelMap.clear();
if(forceCloseChannel > 0) {
logger.info("streamChannelManager forceCloseChannel {}", forceCloseChannel);
}
}
public boolean messageReceived(StreamPacket streamPacket, Channel channel) {
final int channelId = streamPacket.getStreamChannelId();
final StreamChannel streamChannel = findStreamChannel(channelId);
if (streamChannel == null) {
logger.warn("streamChannel not found. channelId:{} ", channelId, channel);
return false;
}
return streamChannel.receiveStreamPacket(streamPacket);
}
}
@@ -6,11 +6,16 @@ import java.util.Map;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.StreamChannelContext;
import com.nhn.pinpoint.rpc.stream.StreamChannelManager;
public class ChannelContext {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private final ServerStreamChannelManager streamChannelManager;
private final StreamChannelManager streamChannelManager;
private final SocketChannel socketChannel;
@@ -20,11 +25,11 @@ public class ChannelContext {
private volatile Map<Object, Object> channelProperties = Collections.emptyMap();
public ChannelContext(SocketChannel socketChannel, ServerStreamChannelManager streamChannelManager) {
public ChannelContext(SocketChannel socketChannel, StreamChannelManager streamChannelManager) {
this(socketChannel, streamChannelManager, DoNothingChannelStateEventListener.INSTANCE);
}
public ChannelContext(SocketChannel socketChannel, ServerStreamChannelManager streamChannelManager, SocketChannelStateChangeEventListener stateChangeEventListener) {
public ChannelContext(SocketChannel socketChannel, StreamChannelManager streamChannelManager, SocketChannelStateChangeEventListener stateChangeEventListener) {
this.socketChannel = socketChannel;
this.streamChannelManager = streamChannelManager;
@@ -33,16 +38,16 @@ public class ChannelContext {
this.state = new PinpointServerSocketState();
}
public ServerStreamChannel getStreamChannel(int channelId) {
public StreamChannelContext getStreamChannel(int channelId) {
return streamChannelManager.findStreamChannel(channelId);
}
public ServerStreamChannel createStreamChannel(int channelId) {
return streamChannelManager.createStreamChannel(channelId);
public ClientStreamChannelContext createStreamChannel(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) {
return streamChannelManager.openStreamChannel(payload, clientStreamChannelMessageListener);
}
public void closeAllStreamChannel() {
streamChannelManager.closeInternal();
streamChannelManager.close();
}
public SocketChannel getSocketChannel() {
@@ -112,5 +117,9 @@ public class ChannelContext {
this.channelProperties = Collections.unmodifiableMap(properties);
return true;
}
public StreamChannelManager getStreamChannelManager() {
return streamChannelManager;
}
}
@@ -47,11 +47,14 @@ import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.ResponsePacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.rpc.packet.ServerClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.stream.DisabledServerStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.ServerStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.StreamChannelManager;
import com.nhn.pinpoint.rpc.util.AssertUtils;
import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils;
import com.nhn.pinpoint.rpc.util.CpuUtils;
import com.nhn.pinpoint.rpc.util.IDGenerator;
import com.nhn.pinpoint.rpc.util.LoggerFactorySetup;
import com.nhn.pinpoint.rpc.util.TimerFactory;
@@ -76,6 +79,8 @@ public class PinpointServerSocket extends SimpleChannelHandler {
private final Timer requestManagerTimer;
private ServerMessageListener messageListener = SimpleLoggingServerMessageListener.LISTENER;
private ServerStreamChannelMessageListener serverStreamChannelMessageListener = DisabledServerStreamChannelMessageListener.INSTANCE;
private WriteFailFutureListener traceSendAckWriteFailFutureListener = new WriteFailFutureListener(logger, "TraceSendAckPacket send fail.", "TraceSendAckPacket send() success.");
private InetAddress[] ignoreAddressList;
@@ -126,6 +131,12 @@ public class PinpointServerSocket extends SimpleChannelHandler {
}
this.messageListener = messageListener;
}
public void setServerStreamChannelMessageListener(ServerStreamChannelMessageListener serverStreamChannelMessageListener) {
AssertUtils.assertNotNull(serverStreamChannelMessageListener, "serverStreamChannelMessageListener must not be null");
this.serverStreamChannelMessageListener = serverStreamChannelMessageListener;
}
private void setOptions(ServerBootstrap bootstrap) {
// read write timeout이 있어야 되나? nio라서 없어도 되던가?
@@ -214,6 +225,8 @@ public class PinpointServerSocket extends SimpleChannelHandler {
case PacketType.APPLICATION_STREAM_CREATE_SUCCESS:
case PacketType.APPLICATION_STREAM_CREATE_FAIL:
case PacketType.APPLICATION_STREAM_RESPONSE:
case PacketType.APPLICATION_STREAM_PING:
case PacketType.APPLICATION_STREAM_PONG:
handleStreamPacket((StreamPacket) message, channel);
return;
case PacketType.CONTROL_ENABLE_WORKER:
@@ -258,31 +271,7 @@ public class PinpointServerSocket extends SimpleChannelHandler {
private void handleStreamPacket(StreamPacket packet, Channel channel) {
ChannelContext context = getChannelContext(channel);
if (packet instanceof StreamCreatePacket) {
logger.debug("StreamCreate {}, streamId:{}", channel, packet.getStreamChannelId());
try {
ServerStreamChannel streamChannel = context.createStreamChannel(packet.getStreamChannelId());
boolean success = streamChannel.receiveChannelCreate((StreamCreatePacket) packet);
if (success) {
messageListener.handleStream(packet, streamChannel);
}
} catch (PinpointSocketException e) {
logger.warn("channel create fail. channel:{} Caused:{}", channel, e);
}
} else if (packet instanceof StreamClosePacket) {
logger.debug("StreamDestroy {}, streamId:{}", channel, packet.getStreamChannelId());
ServerStreamChannel streamChannel = context.getStreamChannel(packet.getStreamChannelId());
// null이 나올수 있음.
boolean close = streamChannel.close();
if (close) {
messageListener.handleStream(packet, streamChannel);
} else {
logger.warn("invalid streamClosePacket. already close. channel:{} Caused:{}", channel);
}
} else {
logger.warn("invalid streamPacket. channel:{}", channel);
}
context.getStreamChannelManager().messageReceived(packet);
}
private Map<Object, Object> decodeSocketProperties(ControlEnableWorkerPacket message) {
@@ -439,7 +428,7 @@ public class PinpointServerSocket extends SimpleChannelHandler {
private void prepareChannel(Channel channel) {
SocketChannel socketChannel = new SocketChannel(channel, DEFAULT_TIMEOUTMILLIS, requestManagerTimer);
ServerStreamChannelManager streamChannelManager = new ServerStreamChannelManager(channel);
StreamChannelManager streamChannelManager = new StreamChannelManager(channel, IDGenerator.createEvenIdGenerator(), serverStreamChannelMessageListener);
ChannelContext channelContext = new ChannelContext(socketChannel, streamChannelManager, channelStateChangeEventListener);
@@ -4,7 +4,6 @@ import java.util.Map;
import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
/**
* @author emeroad
@@ -15,8 +14,6 @@ public interface ServerMessageListener {
// 외부 노출 Channel은 별도의 Tcp Channel로 감싸는걸로 변경할 것.
void handleRequest(RequestPacket requestPacket, SocketChannel channel);
void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel);
int handleEnableWorker(Map properties);
}
@@ -1,159 +0,0 @@
package com.nhn.pinpoint.rpc.server;
import java.util.concurrent.atomic.AtomicInteger;
import org.jboss.netty.channel.Channel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreateFailPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreateSuccessPacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket;
/**
* @author emeroad
*/
public class ServerStreamChannel {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private static final int NONE = 0;
// OPEN이 도착함.
private static final int OPEN_ARRIVED = 1;
// create success 던짐. 동작중
private static final int RUN = 2;
// 닫힘
private static final int CLOSED = 2;
private final AtomicInteger state = new AtomicInteger(NONE);
private final int channelId;
private ServerStreamChannelManager serverStreamChannelManager;
private Channel channel;
public ServerStreamChannel(int channelId) {
this.channelId = channelId;
}
public int getChannelId() {
return channelId;
}
public void setChannel(Channel channel) {
this.channel = channel;
}
// public boolean receiveStreamPacket(StreamPacket packet) {
// final short packetType = packet.getPacketType();
// switch (packetType) {
// case PacketType.APPLICATION_STREAM_CREATE:
// logger.info("APPLICATION_STREAM_CREATE_SUCCESS");
// return receiveChannelCreate((StreamCreatePacket) packet);
// }
// return false;
// }
public boolean receiveChannelCreate(StreamCreatePacket streamCreateResponse) {
if (state.compareAndSet(NONE, OPEN_ARRIVED)) {
return true;
} else {
logger.info("invalid state:{}", state.get());
return false;
}
}
public boolean sendOpenResult(boolean success, byte[] bytes) {
if(success ) {
if(!state.compareAndSet(OPEN_ARRIVED, RUN)) {
return false;
}
StreamCreateSuccessPacket streamCreateSuccessPacket = new StreamCreateSuccessPacket(channelId);
this.channel.write(streamCreateSuccessPacket);
return true;
} else {
if(!state.compareAndSet(OPEN_ARRIVED, CLOSED)) {
return false;
}
StreamCreateFailPacket streamCreateFailPacket = new StreamCreateFailPacket(channelId, StreamCreateFailPacket.UNKNWON_ERROR);
this.channel.write(streamCreateFailPacket);
return true;
}
}
public boolean sendStreamMessage(byte[] bytes) {
if (state.get() != RUN) {
return false;
}
StreamResponsePacket response = new StreamResponsePacket(channelId, bytes);
this.channel.write(response);
return true;
}
public boolean close() {
return close0(true);
}
boolean closeInternal() {
return close0(false);
}
private boolean close0(boolean safeClose) {
if (!state.compareAndSet(RUN, CLOSED)) {
return false;
}
if (safeClose) {
StreamClosePacket streamClosePacket = new StreamClosePacket(channelId, StreamClosePacket.SUCCESS);
this.channel.write(streamClosePacket);
ServerStreamChannelManager serverStreamChannelManager = this.serverStreamChannelManager;
if (serverStreamChannelManager != null) {
serverStreamChannelManager.closeChannel(channelId);
this.serverStreamChannelManager = null;
}
}
return true;
}
public void setServerStreamChannelManager(ServerStreamChannelManager serverStreamChannelManager) {
this.serverStreamChannelManager = serverStreamChannelManager;
}
@Override
public boolean equals(Object o) {
if (this == o) return true;
if (o == null || getClass() != o.getClass()) return false;
ServerStreamChannel that = (ServerStreamChannel) o;
if (channelId != that.channelId) return false;
if (channel != null ? !channel.equals(that.channel) : that.channel != null) return false;
return true;
}
@Override
public int hashCode() {
int result = channelId;
result = 31 * result + (channel != null ? channel.hashCode() : 0);
return result;
}
@Override
public String toString() {
final StringBuilder sb = new StringBuilder();
sb.append("ServerStreamChannel");
sb.append("{state=").append(state);
sb.append(", channelId=").append(channelId);
sb.append(", channel=").append(channel);
sb.append('}');
return sb.toString();
}
}
@@ -1,65 +0,0 @@
package com.nhn.pinpoint.rpc.server;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import org.jboss.netty.channel.Channel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import java.util.Map;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentMap;
/**
* @author emeroad
*/
public class ServerStreamChannelManager {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
private final Channel channel;
private final ConcurrentMap<Integer, ServerStreamChannel> channelMap = new ConcurrentHashMap<Integer, ServerStreamChannel>();
public ServerStreamChannelManager(Channel channel) {
if (channel == null) {
throw new NullPointerException("channel");
}
this.channel = channel;
}
public ServerStreamChannel createStreamChannel(int channelId) {
ServerStreamChannel streamChannel = new ServerStreamChannel(channelId);
streamChannel.setChannel(channel);
ServerStreamChannel old = channelMap.put(channelId, streamChannel);
if (old != null) {
throw new PinpointSocketException("already channelId exist:" + channelId + " streamChannel:" + old);
}
// handle을 붙여서 리턴.
streamChannel.setServerStreamChannelManager(this);
return streamChannel;
}
public ServerStreamChannel findStreamChannel(int channelId) {
return this.channelMap.get(channelId);
}
public boolean closeChannel(int channelId) {
ServerStreamChannel remove = this.channelMap.remove(channelId);
return remove != null;
}
public void closeInternal() {
final boolean debugEnabled = logger.isDebugEnabled();
for (Map.Entry<Integer, ServerStreamChannel> streamChannel : this.channelMap.entrySet()) {
streamChannel.getValue().closeInternal();
if (debugEnabled) {
logger.debug("ServerStreamChannel.closeInternal() id:{}, {}", streamChannel.getKey(), channel);
}
}
}
}
@@ -8,7 +8,6 @@ import org.slf4j.LoggerFactory;
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.stream.StreamPacket;
/**
* @author emeroad
@@ -29,12 +28,6 @@ public class SimpleLoggingServerMessageListener implements ServerMessageListener
logger.info("handlerRequest {} {}", requestPacket, channel);
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.info("handlerStream {} {}", streamChannel, streamChannel);
}
@Override
public int handleEnableWorker(Map properties) {
logger.info("handleEnableWorker {}", properties);
@@ -1,19 +0,0 @@
package com.nhn.pinpoint.rpc.stream;
import com.nhn.pinpoint.rpc.client.StreamChannel;
import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket;
/**
* @author koo.taejin <kr14910>
*/
public interface StreamChannelMessageListener {
short handleStreamCreate(StreamChannel streamChannel, StreamCreatePacket packet);
void handleStreamData(StreamChannel streamChannel, StreamResponsePacket packet);
void handleStreamClose(StreamChannel streamChannel, StreamClosePacket packet);
}
@@ -8,17 +8,17 @@ import java.util.concurrent.CountDownLatch;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.rpc.client.StreamChannel;
import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket;
import com.nhn.pinpoint.rpc.stream.StreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext;
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.stream.StreamChannelContext;
/**
* @author emeroad
* @author koo.taejin <kr14910>
*/
public class RecordedStreamChannelMessageListener implements StreamChannelMessageListener {
public class RecordedStreamChannelMessageListener implements ClientStreamChannelMessageListener {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
@@ -29,24 +29,17 @@ public class RecordedStreamChannelMessageListener implements StreamChannelMessag
public RecordedStreamChannelMessageListener(int receiveMessageCount) {
this.latch = new CountDownLatch(receiveMessageCount);
}
@Override
public short handleStreamCreate(StreamChannel streamChannel, StreamCreatePacket packet) {
// TODO Auto-generated method stub
return 0;
}
@Override
public void handleStreamData(StreamChannel streamChannel, StreamResponsePacket packet) {
logger.info("handleStreamData {}, {}", streamChannel, packet);
public void handleStreamData(ClientStreamChannelContext streamChannelContext, StreamResponsePacket packet) {
logger.info("handleStreamData {}, {}", streamChannelContext, packet);
receivedMessageList.add(packet.getPayload());
latch.countDown();
}
@Override
public void handleStreamClose(StreamChannel streamChannel, StreamClosePacket packet) {
logger.info("handleClose {}, {}", streamChannel, packet);
public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) {
logger.info("handleClose {}, {}", streamChannelContext, packet);
receivedMessageList.add(packet.getPayload());
latch.countDown();
}
@@ -180,55 +180,6 @@ public class PinpointSocketFactoryTest {
}
@Test
public void stream() throws IOException, InterruptedException {
PinpointServerSocket ss = new PinpointServerSocket();
TestSeverMessageListener testSeverMessageListener = new TestSeverMessageListener();
ss.setMessageListener(testSeverMessageListener);
ss.bind("localhost", 10234);
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
try {
PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234);
StreamChannel streamChannel = socket.createStreamChannel();
byte[] openBytes = TestByteUtils.createRandomByte(30);
// 현재 서버에서 3번 보내게 되어 있음.
RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4);
streamChannel.setStreamChannelMessageListener(clientListener);
Future<StreamCreateResponse> open = streamChannel.open(openBytes);
open.await();
StreamCreateResponse response = open.getResult();
Assert.assertTrue(response.isSuccess());
// stream 메시지를 대기함.
clientListener.getLatch().await();
List<byte[]> receivedMessage = clientListener.getReceivedMessage();
List<byte[]> sendMessage = testSeverMessageListener.getSendMessage();
// 한개는 close 패킷임.
Assert.assertEquals(receivedMessage.size(), sendMessage.size());
for(int i =0; i<receivedMessage.size(); i++) {
Assert.assertArrayEquals(receivedMessage.get(i), sendMessage.get(i));
}
socket.close();
} finally {
pinpointSocketFactory.release();
ss.close();
}
}
@Test
public void connectTimeout() {
PinpointSocketFactory pinpointSocketFactory = null;
@@ -22,7 +22,6 @@ 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;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils;
import com.nhn.pinpoint.rpc.util.MapUtils;
@@ -246,12 +245,6 @@ public class ControlPacketServerTest {
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) {
@@ -21,7 +21,6 @@ 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;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.util.ControlMessageEnDeconderUtils;
import com.nhn.pinpoint.rpc.util.MapUtils;
@@ -168,11 +167,6 @@ public class EventListnerTest {
logger.info("handlerRequest {}", requestPacket, channel);
channel.sendResponseMessage(requestPacket, requestPacket.getPayload());
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
}
@Override
public int handleEnableWorker(Map properties) {
@@ -32,12 +32,16 @@ public class MessageListenerTest {
PinpointServerSocket ss = new PinpointServerSocket();
ss.bind("127.0.0.1", 10234);
PinpointSocketFactory socketFactory = createPinpointSocketFactory();
PinpointSocketFactory socketFactory1 = createPinpointSocketFactory();
socketFactory1.setMessageListener(new EchoMessageListener());
PinpointSocketFactory socketFactory2 = createPinpointSocketFactory();
try {
// 리스터를 등록한 것만 RegisterAgent 로 나옴
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, new EchoMessageListener());
PinpointSocket socket2 = socketFactory.connect("127.0.0.1", 10234);
PinpointSocket socket = socketFactory1.connect("127.0.0.1", 10234);
PinpointSocket socket2 = socketFactory2.connect("127.0.0.1", 10234);
Thread.sleep(500);
@@ -49,7 +53,9 @@ public class MessageListenerTest {
socket.close();
socket2.close();
} finally {
socketFactory.release();
socketFactory1.release();
socketFactory2.release();
ss.close();
}
}
@@ -59,14 +65,14 @@ public class MessageListenerTest {
PinpointServerSocket ss = new PinpointServerSocket();
ss.bind("127.0.0.1", 10234);
EchoMessageListener echoMessageListener = new EchoMessageListener();
PinpointSocketFactory socketFactory = createPinpointSocketFactory();
socketFactory.setMessageListener(echoMessageListener);
try {
EchoMessageListener echoMessageListener = new EchoMessageListener();
// 리스터를 등록한 것만 RegisterAgent 로 나옴
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, echoMessageListener);
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234);
Thread.sleep(500);
List<ChannelContext> channelContextList = ss.getDuplexCommunicationChannelContext();
@@ -99,17 +105,18 @@ public class MessageListenerTest {
PinpointServerSocket ss = new PinpointServerSocket();
ss.bind("127.0.0.1", 10234);
PinpointSocketFactory socketFactory = createPinpointSocketFactory();
PinpointSocketFactory socketFactory1 = createPinpointSocketFactory();
EchoMessageListener echoMessageListener1 = new EchoMessageListener();
socketFactory1.setMessageListener(echoMessageListener1);
PinpointSocketFactory socketFactory2 = createPinpointSocketFactory();
EchoMessageListener echoMessageListener2 = new EchoMessageListener();
socketFactory2.setMessageListener(echoMessageListener2);
try {
EchoMessageListener echoMessageListener1 = new EchoMessageListener();
EchoMessageListener echoMessageListener2 = new EchoMessageListener();
// 리스터를 등록한 것만 RegisterAgent 로 나옴
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, echoMessageListener1);
PinpointSocket socket2 = socketFactory.connect("127.0.0.1", 10234, echoMessageListener2);
PinpointSocket socket = socketFactory1.connect("127.0.0.1", 10234);
PinpointSocket socket2 = socketFactory2.connect("127.0.0.1", 10234);
Thread.sleep(500);
@@ -131,7 +138,9 @@ public class MessageListenerTest {
socket.close();
socket2.close();
} finally {
socketFactory.release();
socketFactory1.release();
socketFactory2.release();
ss.close();
}
}
@@ -143,12 +152,12 @@ public class MessageListenerTest {
Map params = getParams();
PinpointSocketFactory socketFactory = createPinpointSocketFactory(params);
socketFactory.setMessageListener(new EchoMessageListener());
try {
EchoMessageListener echoMessageListener1 = new EchoMessageListener();
// 리스터를 등록한 것만 RegisterAgent 로 나옴
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, echoMessageListener1);
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234);
Thread.sleep(500);
@@ -171,10 +180,11 @@ public class MessageListenerTest {
ss.bind("127.0.0.1", 10234);
PinpointSocketFactory socketFactory = createPinpointSocketFactory();
socketFactory.setMessageListener(SimpleLoggingMessageListener.LISTENER);
try {
// Listener가 없을때 디폴트로 등록하는 SimpleLoggingMessageListener.LISTENER인 경우 상호 연결이 불가능함
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, SimpleLoggingMessageListener.LISTENER);
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234);
Thread.sleep(500);
@@ -201,11 +211,12 @@ public class MessageListenerTest {
PinpointSocketFactory socketFactory = createPinpointSocketFactory();
socketFactory.setEnableWorkerPacketDelay(500);
socketFactory.setMessageListener(new EchoMessageListener());
try {
// Listener가 없을때 디폴트로 등록하는 SimpleLoggingMessageListener.LISTENER인 경우 상호 연결이 불가능함
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234, new EchoMessageListener());
PinpointSocket socket = socketFactory.connect("127.0.0.1", 10234);
Thread.sleep(5000);
List<ChannelContext> channelContextList = ss.getDuplexCommunicationChannelContext();
@@ -7,13 +7,9 @@ import java.util.Map;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.rpc.TestByteUtils;
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.stream.StreamClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
/**
* @author emeroad
@@ -37,42 +33,12 @@ public class TestSeverMessageListener implements ServerMessageListener {
channel.sendResponseMessage(requestPacket, requestPacket.getPayload());
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.debug("streamPacket:{} channel:{}", streamPacket, streamChannel);
if (streamPacket instanceof StreamCreatePacket) {
byte[] payload = streamPacket.getPayload();
this.open = payload;
streamChannel.sendOpenResult(true, payload);
sendStreamMessage(streamChannel);
sendStreamMessage(streamChannel);
sendStreamMessage(streamChannel);
sendClose(streamChannel);
} else if(streamPacket instanceof StreamClosePacket) {
// 채널 종료해야 함.
}
}
@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]);
streamChannel.close();
}
private void sendStreamMessage(ServerStreamChannel streamChannel) {
byte[] randomByte = TestByteUtils.createRandomByte(10);
streamChannel.sendStreamMessage(randomByte);
sendMessageList.add(randomByte);
}
public byte[] getOpen() {
return open;
}
@@ -0,0 +1,244 @@
package com.nhn.pinpoint.rpc.stream;
import java.io.IOException;
import java.util.List;
import java.util.concurrent.CopyOnWriteArrayList;
import junit.framework.Assert;
import org.junit.Test;
import com.nhn.pinpoint.rpc.PinpointSocketException;
import com.nhn.pinpoint.rpc.RecordedStreamChannelMessageListener;
import com.nhn.pinpoint.rpc.TestByteUtils;
import com.nhn.pinpoint.rpc.client.MessageListener;
import com.nhn.pinpoint.rpc.client.PinpointSocket;
import com.nhn.pinpoint.rpc.client.PinpointSocketFactory;
import com.nhn.pinpoint.rpc.client.SimpleLoggingMessageListener;
import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket;
import com.nhn.pinpoint.rpc.packet.stream.StreamCreatePacket;
import com.nhn.pinpoint.rpc.server.ChannelContext;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.TestSeverMessageListener;
public class StreamChannelManagerTest {
// Client to Server Stream
@Test
public void stream1() throws IOException, InterruptedException {
SimpleStreamBO bo = new SimpleStreamBO();
PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), new ServerListener(bo));
ss.bind("localhost", 10234);
PinpointSocketFactory pinpointSocketFactory = createSocketFactory();
try {
PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234);
RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4);
ClientStreamChannelContext clientContext = socket.createStreamChannel(new byte[0], clientListener);
int sendCount = 4;
for (int i = 0; i < 4; i++) {
sendRandomBytes(bo);
}
Thread.sleep(100);
Assert.assertEquals(sendCount, clientListener.getReceivedMessage().size());
clientContext.getStreamChannel().close();
socket.close();
} finally {
pinpointSocketFactory.release();
ss.close();
}
}
@Test(expected = PinpointSocketException.class)
public void stream2() throws IOException, InterruptedException {
PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), null);
ss.bind("localhost", 10234);
PinpointSocketFactory pinpointSocketFactory = createSocketFactory();
try {
PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234);
RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4);
ClientStreamChannelContext clientContext = socket.createStreamChannel(new byte[0], clientListener);
Thread.sleep(100);
clientContext.getStreamChannel().close();
socket.close();
} finally {
pinpointSocketFactory.release();
ss.close();
}
}
@Test(expected = PinpointSocketException.class)
public void stream3() throws IOException, InterruptedException {
SimpleStreamBO bo = new SimpleStreamBO();
PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), new ServerListener(bo));
ss.bind("localhost", 10234);
PinpointSocketFactory pinpointSocketFactory = createSocketFactory();
PinpointSocket socket = null;
try {
socket = pinpointSocketFactory.connect("127.0.0.1", 10234);
RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4);
ClientStreamChannelContext clientContext = socket.createStreamChannel(new byte[0], clientListener);
Thread.sleep(100);
clientContext.getStreamChannel().close();
Thread.sleep(100);
sendRandomBytes(bo);
} finally {
if (socket != null) {
socket.close();
}
pinpointSocketFactory.release();
ss.close();
}
}
// ServerSocket to Client Stream
@Test
public void stream4() throws IOException, InterruptedException {
PinpointServerSocket ss = createServerSocket(new TestSeverMessageListener(), null);
ss.bind("localhost", 10234);
SimpleStreamBO bo = new SimpleStreamBO();
PinpointSocketFactory pinpointSocketFactory = createSocketFactory(new TestListener(), new ServerListener(bo));
try {
PinpointSocket socket = pinpointSocketFactory.connect("127.0.0.1", 10234);
Thread.sleep(100);
List<ChannelContext> contextList = ss.getDuplexCommunicationChannelContext();
Assert.assertEquals(1, contextList.size());
ChannelContext context = contextList.get(0);
RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4);
ClientStreamChannelContext clientContext = context.createStreamChannel(new byte[0], clientListener);
int sendCount = 4;
for (int i = 0; i < 4; i++) {
sendRandomBytes(bo);
}
Thread.sleep(100);
Assert.assertEquals(sendCount, clientListener.getReceivedMessage().size());
clientContext.getStreamChannel().close();
socket.close();
} finally {
pinpointSocketFactory.release();
ss.close();
}
}
private PinpointServerSocket createServerSocket(ServerMessageListener severMessageListener,
ServerStreamChannelMessageListener serverStreamChannelMessageListener) {
PinpointServerSocket serverSocket = new PinpointServerSocket();
if (severMessageListener != null) {
serverSocket.setMessageListener(severMessageListener);
}
if (serverStreamChannelMessageListener != null) {
serverSocket.setServerStreamChannelMessageListener(serverStreamChannelMessageListener);
}
return serverSocket;
}
private PinpointSocketFactory createSocketFactory() {
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
return pinpointSocketFactory;
}
private PinpointSocketFactory createSocketFactory(MessageListener messageListener, ServerStreamChannelMessageListener serverStreamChannelMessageListener) {
PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory();
pinpointSocketFactory.setMessageListener(messageListener);
pinpointSocketFactory.setServerStreamChannelMessageListener(serverStreamChannelMessageListener);
return pinpointSocketFactory;
}
class TestListener extends SimpleLoggingMessageListener {
}
private void sendRandomBytes(SimpleStreamBO bo) {
byte[] openBytes = TestByteUtils.createRandomByte(30);
bo.sendResponse(openBytes);
}
class ServerListener implements ServerStreamChannelMessageListener {
private final SimpleStreamBO bo;
public ServerListener(SimpleStreamBO bo) {
this.bo = bo;
}
@Override
public short handleStreamCreate(ServerStreamChannelContext streamChannelContext, StreamCreatePacket packet) {
bo.addServerStreamChannelContext(streamChannelContext);
return 0;
}
@Override
public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) {
}
}
class SimpleStreamBO {
private final List<ServerStreamChannelContext> serverStreamChannelContextList;
public SimpleStreamBO() {
serverStreamChannelContextList = new CopyOnWriteArrayList<ServerStreamChannelContext>();
}
public void addServerStreamChannelContext(ServerStreamChannelContext context) {
serverStreamChannelContextList.add(context);
}
public void removeServerStreamChannelContext(ServerStreamChannelContext context) {
serverStreamChannelContextList.remove(context);
}
void sendResponse(byte[] data) {
for (ServerStreamChannelContext context : serverStreamChannelContextList) {
context.getStreamChannel().sendData(data);
}
}
}
}
@@ -22,8 +22,8 @@ import com.nhn.pinpoint.rpc.packet.stream.StreamPacket;
import com.nhn.pinpoint.rpc.server.ChannelContext;
import com.nhn.pinpoint.rpc.server.PinpointServerSocket;
import com.nhn.pinpoint.rpc.server.ServerMessageListener;
import com.nhn.pinpoint.rpc.server.ServerStreamChannel;
import com.nhn.pinpoint.rpc.server.SocketChannel;
import com.nhn.pinpoint.rpc.stream.ServerStreamChannel;
import com.nhn.pinpoint.web.cluster.ClusterManager;
import com.nhn.pinpoint.web.cluster.zookeeper.ZookeeperClusterManager;
import com.nhn.pinpoint.web.config.WebConfig;
@@ -168,11 +168,6 @@ public class PinpointSocketManager {
logger.warn("Unsupport request received {} {}", requestPacket, channel);
}
@Override
public void handleStream(StreamPacket streamPacket, ServerStreamChannel streamChannel) {
logger.warn("unsupported streamPacket received {}", streamPacket);
}
@Override
public int handleEnableWorker(Map properties) {
logger.warn("do handleEnableWorker {}", properties);
@@ -123,7 +123,9 @@ public class ClusterTest {
Assert.assertEquals(0, socketManager.getCollectorChannelContext().size());
factory = new PinpointSocketFactory();
socket = factory.connect(DEFAULT_IP, DEFAULT_ACCEPTOR_PORT, new SimpleListener());
factory.setMessageListener(new SimpleListener());
socket = factory.connect(DEFAULT_IP, DEFAULT_ACCEPTOR_PORT);
Thread.sleep(1000);