From ef45fa59cfb0d7f37e952775015c807eb08078d6 Mon Sep 17 00:00:00 2001 From: kr14910 Date: Mon, 16 Feb 2015 14:16:30 +0900 Subject: [PATCH] Improve state synchronization between server and client. #136 WritablePinpointServer -> PinpointServer PinpointServer -> DefaultPinpointServer --- .../cluster/PinpointServerClusterPoint.java | 16 +- .../ZookeeperProfilerClusterManager.java | 4 +- .../collector/receiver/tcp/TCPReceiver.java | 14 +- .../profiler/AgentInfoSenderTest.java | 6 +- .../sender/TcpDataSenderReconnectTest.java | 6 +- .../profiler/sender/TcpDataSenderTest.java | 6 +- .../RequestResponseServerMessageListener.java | 8 +- .../pinpoint/rpc/client/RequestManager.java | 4 +- .../rpc/server/DefaultPinpointServer.java | 470 ++++++++++++++++++ .../pinpoint/rpc/server/PinpointServer.java | 419 +--------------- .../rpc/server/PinpointServerAcceptor.java | 20 +- .../rpc/server/ServerMessageListener.java | 4 +- .../SimpleLoggingServerMessageListener.java | 4 +- .../rpc/server/WritablePinpointServer.java | 45 -- .../DoNothingChannelStateEventHandler.java | 6 +- .../rpc/client/ClientMessageListenerTest.java | 16 +- .../rpc/server/ControlPacketServerTest.java | 4 +- .../pinpoint/rpc/server/EventHandlerTest.java | 4 +- .../pinpoint/rpc/server/HandshakeTest.java | 10 +- .../rpc/server/TestSeverMessageListener.java | 4 +- .../rpc/stream/StreamChannelManagerTest.java | 10 +- .../rpc/util/PinpointRPCTestUtils.java | 8 +- .../thrift/io/TCommandTypeVersion.java | 2 + .../web/controller/CommandController.java | 6 +- .../web/server/PinpointSocketManager.java | 14 +- 25 files changed, 574 insertions(+), 536 deletions(-) create mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/server/DefaultPinpointServer.java delete mode 100644 rpc/src/main/java/com/navercorp/pinpoint/rpc/server/WritablePinpointServer.java diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/PinpointServerClusterPoint.java b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/PinpointServerClusterPoint.java index 1689d4319..5da5c103f 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/PinpointServerClusterPoint.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/PinpointServerClusterPoint.java @@ -92,7 +92,21 @@ public class PinpointServerClusterPoint implements TargetClusterPoint { @Override public String toString() { - return pinpointServer.toString(); + StringBuilder log = new StringBuilder(32); + log.append(this.getClass().getSimpleName()); + log.append("("); + log.append(applicationName); + log.append("/"); + log.append(agentId); + log.append("/"); + log.append(startTimeStamp); + log.append(")"); + log.append(", version:"); + log.append(version); + log.append(", pinpointServer:"); + log.append(pinpointServer); + + return log.toString(); } @Override diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java index 54f045a0a..e58cf12d3 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/cluster/zookeeper/ZookeeperProfilerClusterManager.java @@ -144,9 +144,7 @@ public class ZookeeperProfilerClusterManager implements ChannelStateChangeEventH @Override public void exceptionCaught(PinpointServer pinpointServer, PinpointServerStateCode stateCode, Throwable e) { - if (logger.isWarnEnabled()) { - logger.warn(this.getClass().getSimpleName() + " exception occured. Error: " + e.getMessage() + "." , e); - } + logger.warn("ZookeeperProfilerClusterManager exceptionCaught() (pinpointServer:{}, PinpointServerStateCode:{}). Error: {}.", pinpointServer, stateCode, e.getMessage(), e); } public List getClusterData() { diff --git a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java index 7da62e9e8..dd2f04cbf 100644 --- a/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java +++ b/collector/src/main/java/com/navercorp/pinpoint/collector/receiver/tcp/TCPReceiver.java @@ -47,7 +47,7 @@ import com.navercorp.pinpoint.rpc.packet.HandshakeResponseType; import com.navercorp.pinpoint.rpc.packet.RequestPacket; import com.navercorp.pinpoint.rpc.packet.SendPacket; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; import com.navercorp.pinpoint.rpc.util.MapUtils; import com.navercorp.pinpoint.thrift.io.DeserializerFactory; @@ -145,12 +145,12 @@ public class TCPReceiver { // pass them to a separate queue and handle them in a different thread. this.serverAcceptor.setMessageListener(new ServerMessageListener() { @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { receive(sendPacket, pinpointServer); } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { requestResponse(requestPacket, pinpointServer); } @@ -178,7 +178,7 @@ public class TCPReceiver { } - private void receive(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + private void receive(SendPacket sendPacket, PinpointServer pinpointServer) { try { worker.execute(new Dispatch(sendPacket.getPayload(), pinpointServer.getRemoteAddress())); } catch (RejectedExecutionException e) { @@ -187,7 +187,7 @@ public class TCPReceiver { } } - private void requestResponse(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + private void requestResponse(RequestPacket requestPacket, PinpointServer pinpointServer) { try { worker.execute(new RequestResponseDispatch(requestPacket, pinpointServer)); } catch (RejectedExecutionException e) { @@ -235,10 +235,10 @@ public class TCPReceiver { private class RequestResponseDispatch implements Runnable { private final RequestPacket requestPacket; - private final WritablePinpointServer pinpointServer; + private final PinpointServer pinpointServer; - private RequestResponseDispatch(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + private RequestResponseDispatch(RequestPacket requestPacket, PinpointServer pinpointServer) { if (requestPacket == null) { throw new NullPointerException("requestPacket"); } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java index 553342686..72a5abd44 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/AgentInfoSenderTest.java @@ -42,7 +42,7 @@ import com.navercorp.pinpoint.rpc.packet.RequestPacket; import com.navercorp.pinpoint.rpc.packet.SendPacket; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.thrift.dto.TResult; import com.navercorp.pinpoint.thrift.io.HeaderTBaseSerializer; import com.navercorp.pinpoint.thrift.io.HeaderTBaseSerializerFactory; @@ -264,13 +264,13 @@ public class AgentInfoSenderTest { } @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { logger.info("handleSend:{}", sendPacket); } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { int requestCount = this.requestCount.incrementAndGet(); if (requestCount < successCondition) { diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java index 66530d440..47c4dbcc3 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java @@ -32,7 +32,7 @@ import com.navercorp.pinpoint.rpc.packet.RequestPacket; import com.navercorp.pinpoint.rpc.packet.SendPacket; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.thrift.dto.TApiMetaData; /** @@ -52,13 +52,13 @@ public class TcpDataSenderReconnectTest { serverAcceptor.setMessageListener(new ServerMessageListener() { @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { logger.info("handleSend:{}", sendPacket); send++; } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { logger.info("handleRequest:{}", requestPacket); } diff --git a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java index 3b6077947..206a07435 100644 --- a/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java +++ b/profiler/src/test/java/com/navercorp/pinpoint/profiler/sender/TcpDataSenderTest.java @@ -38,7 +38,7 @@ import com.navercorp.pinpoint.rpc.packet.RequestPacket; import com.navercorp.pinpoint.rpc.packet.SendPacket; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.thrift.dto.TApiMetaData; /** @@ -60,7 +60,7 @@ public class TcpDataSenderTest { serverAcceptor.setMessageListener(new ServerMessageListener() { @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { logger.info("handleSend:{}", sendPacket); if (sendLatch != null) { sendLatch.countDown(); @@ -68,7 +68,7 @@ public class TcpDataSenderTest { } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { logger.info("handleRequest:{}", requestPacket); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java index cac51192b..2b0b0b3ac 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/RequestResponseServerMessageListener.java @@ -25,7 +25,7 @@ import com.navercorp.pinpoint.rpc.packet.HandshakeResponseCode; import com.navercorp.pinpoint.rpc.packet.HandshakeResponseType; import com.navercorp.pinpoint.rpc.packet.RequestPacket; import com.navercorp.pinpoint.rpc.packet.SendPacket; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; /** @@ -38,14 +38,14 @@ public class RequestResponseServerMessageListener implements ServerMessageListen public static final RequestResponseServerMessageListener LISTENER = new RequestResponseServerMessageListener(); @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { logger.info("handlerSend {} {}", sendPacket, pinpointServer); } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { - logger.info("handlerRequest {}", requestPacket, pinpointServer); + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { + logger.info("handlerRequest {} {}", requestPacket, pinpointServer); pinpointServer.response(requestPacket, requestPacket.getPayload()); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/RequestManager.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/RequestManager.java index a6ee659e3..2a84dfa38 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/RequestManager.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/client/RequestManager.java @@ -116,10 +116,10 @@ public class RequestManager { final int requestId = responsePacket.getRequestId(); final DefaultFuture future = removeMessageFuture(requestId); if (future == null) { - logger.warn("future not found:{}, channel:{}", responsePacket, pinpointServer); + logger.warn("future not found:{}, pinpointServer:{}", responsePacket, pinpointServer); return; } else { - logger.debug("responsePacket arrived packet:{}, channel:{}", responsePacket, pinpointServer); + logger.debug("responsePacket arrived packet:{}, pinpointServer:{}", responsePacket, pinpointServer); } ResponseMessage response = new ResponseMessage(); diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/DefaultPinpointServer.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/DefaultPinpointServer.java new file mode 100644 index 000000000..2ba0b6f97 --- /dev/null +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/DefaultPinpointServer.java @@ -0,0 +1,470 @@ +/* + * Copyright 2014 NAVER Corp. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.navercorp.pinpoint.rpc.server; + +import java.lang.reflect.Array; +import java.net.SocketAddress; +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.atomic.AtomicReference; + +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.navercorp.pinpoint.rpc.ChannelWriteFailListenableFuture; +import com.navercorp.pinpoint.rpc.Future; +import com.navercorp.pinpoint.rpc.ResponseMessage; +import com.navercorp.pinpoint.rpc.client.RequestManager; +import com.navercorp.pinpoint.rpc.client.WriteFailFutureListener; +import com.navercorp.pinpoint.rpc.control.ProtocolException; +import com.navercorp.pinpoint.rpc.packet.ControlHandshakePacket; +import com.navercorp.pinpoint.rpc.packet.ControlHandshakeResponsePacket; +import com.navercorp.pinpoint.rpc.packet.HandshakeResponseCode; +import com.navercorp.pinpoint.rpc.packet.Packet; +import com.navercorp.pinpoint.rpc.packet.PacketType; +import com.navercorp.pinpoint.rpc.packet.RequestPacket; +import com.navercorp.pinpoint.rpc.packet.ResponsePacket; +import com.navercorp.pinpoint.rpc.packet.SendPacket; +import com.navercorp.pinpoint.rpc.packet.ServerClosePacket; +import com.navercorp.pinpoint.rpc.packet.stream.StreamPacket; +import com.navercorp.pinpoint.rpc.server.handler.ChannelStateChangeEventHandler; +import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelContext; +import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelMessageListener; +import com.navercorp.pinpoint.rpc.stream.StreamChannelContext; +import com.navercorp.pinpoint.rpc.stream.StreamChannelManager; +import com.navercorp.pinpoint.rpc.util.AssertUtils; +import com.navercorp.pinpoint.rpc.util.ClassUtils; +import com.navercorp.pinpoint.rpc.util.ControlMessageEncodingUtils; +import com.navercorp.pinpoint.rpc.util.IDGenerator; +import com.navercorp.pinpoint.rpc.util.ListUtils; + +/** + * @author Taejin Koo + */ +public class DefaultPinpointServer implements PinpointServer { + + private final Logger logger = LoggerFactory.getLogger(this.getClass()); + + private final Channel channel; + private final RequestManager requestManager; + + private final PinpointServerState state; + + private final ServerMessageListener messageListener; + + private final List stateChangeEventListeners; + + private final StreamChannelManager streamChannelManager; + + private final AtomicReference> properties = new AtomicReference>(); + + private final String objectUniqName; + + private final ChannelFutureListener serverCloseWriteListener; + private final ChannelFutureListener responseWriteFailListener; + + public DefaultPinpointServer(Channel channel, PinpointServerConfig serverConfig) { + this(channel, serverConfig, null); + } + + public DefaultPinpointServer(Channel channel, PinpointServerConfig serverConfig, ChannelStateChangeEventHandler... stateChangeEventListeners) { + this.channel = channel; + + this.messageListener = serverConfig.getMessageListener(); + + StreamChannelManager streamChannelManager = new StreamChannelManager(channel, IDGenerator.createEvenIdGenerator(), serverConfig.getStreamMessageListener()); + this.streamChannelManager = streamChannelManager; + + if (stateChangeEventListeners == null) { + this.stateChangeEventListeners = new ArrayList(1); + } else { + this.stateChangeEventListeners = new ArrayList(Array.getLength(stateChangeEventListeners) + 1); + } + + ListUtils.addIfValueNotNull(this.stateChangeEventListeners, serverConfig.getStateChangeEventHandler()); + ListUtils.addAllExceptNullValue(this.stateChangeEventListeners, stateChangeEventListeners); + + RequestManager requestManager = new RequestManager(serverConfig.getRequestManagerTimer(), serverConfig.getDefaultRequestTimeout()); + this.requestManager = requestManager; + + this.state = new PinpointServerState(); + + this.objectUniqName = ClassUtils.simpleClassNameAndHashCodeString(this); + + this.serverCloseWriteListener = new WriteFailFutureListener(logger, objectUniqName + " sendClosePacket() write fail.", "serverClosePacket write success"); + this.responseWriteFailListener = new WriteFailFutureListener(logger, objectUniqName + " response() write fail."); + } + + public void start() { + changeStateToRunWithoutHandshake(); + } + + public void stop() { + if (PinpointServerStateCode.BEING_SHUTDOWN == getCurrentStateCode()) { + changeStateToShutdown(); + } else { + changeStateToUnexpectedShutdown(); + } + + if (this.channel.isConnected()) { + channel.close(); + } + + streamChannelManager.close(); + } + + @Override + public void send(byte[] payload) { + AssertUtils.assertNotNull(payload, "payload may not be null."); + if (!isEnableDuplexCommunication()) { + throw new IllegalStateException("Send fail. Error: Illegal State. pinpointServer:" + toString()); + } + + SendPacket send = new SendPacket(payload); + write0(send); + } + + @Override + public Future request(byte[] payload) { + AssertUtils.assertNotNull(payload, "payload may not be null."); + if (!isEnableDuplexCommunication()) { + throw new IllegalStateException("Request fail. Error: Illegal State. pinpointServer:" + toString()); + } + + RequestPacket requestPacket = new RequestPacket(payload); + ChannelWriteFailListenableFuture messageFuture = this.requestManager.register(requestPacket); + write0(requestPacket, messageFuture); + return messageFuture; + } + + @Override + public void response(RequestPacket requestPacket, byte[] payload) { + response(requestPacket.getRequestId(), payload); + } + + @Override + public void response(int requestId, byte[] payload) { + AssertUtils.assertNotNull(payload, "payload may not be null."); + if (!isEnableCommunication()) { + throw new IllegalStateException("Response fail. Error: Illegal State. pinpointServer:" + toString()); + } + + ResponsePacket responsePacket = new ResponsePacket(requestId, payload); + write0(responsePacket, responseWriteFailListener); + } + + private ChannelFuture write0(Object message) { + return write0(message, null); + } + + private ChannelFuture write0(Object message, ChannelFutureListener futureListener) { + ChannelFuture future = channel.write(message); + if (futureListener != null) { + future.addListener(futureListener);; + } + return future; + } + + public StreamChannelContext getStreamChannel(int channelId) { + return streamChannelManager.findStreamChannel(channelId); + } + + @Override + public ClientStreamChannelContext createStream(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) { + return streamChannelManager.openStreamChannel(payload, clientStreamChannelMessageListener); + } + + public void closeAllStreamChannel() { + streamChannelManager.close(); + } + + @Override + public Map getChannelProperties() { + Map properties = this.properties.get(); + return properties == null ? Collections.emptyMap() : properties; + } + + public boolean setChannelProperties(Map value) { + if (value == null) { + return false; + } + + return this.properties.compareAndSet(null, Collections.unmodifiableMap(value)); + } + + @Override + public SocketAddress getRemoteAddress() { + return channel.getRemoteAddress(); + } + + public ChannelFuture sendClosePacket() { + logger.info("sendServerClosedPacket start"); + + PinpointServerStateCode errorCode = changeStateBeingShutdown(); + + if (errorCode == null) { + final ChannelFuture writeFuture = this.channel.write(ServerClosePacket.DEFAULT_SERVER_CLOSE_PACKET); + writeFuture.addListener(serverCloseWriteListener); + + logger.info("sendServerClosedPacket end"); + return writeFuture; + } else { + logger.info("sendServerClosedPacket fail. Error: change state failed."); + return null; + } + } + + @Override + public void messageReceived(Object message) { + // TODO Auto-generated method stub + if (!PinpointServerStateCode.isRun(getCurrentStateCode())) { + // FIXME need change rules. + // as-is : do nothing when state is not run. + // candidate : close channel when state is not run. + logger.warn("{} messageReceived:{} from IllegalState this message will be ignore.", this, message); + return; + } + + final short packetType = getPacketType(message); + switch (packetType) { + case PacketType.APPLICATION_SEND: { + handleSend((SendPacket) message); + return; + } + case PacketType.APPLICATION_REQUEST: { + handleRequest((RequestPacket) message); + return; + } + case PacketType.APPLICATION_RESPONSE: { + handleResponse((ResponsePacket) message); + return; + } + case PacketType.APPLICATION_STREAM_CREATE: + case PacketType.APPLICATION_STREAM_CLOSE: + 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: + handleStreamEvent((StreamPacket) message); + return; + case PacketType.CONTROL_HANDSHAKE: + handleHandshake((ControlHandshakePacket) message); + return; + case PacketType.CONTROL_CLIENT_CLOSE: { + handleClosePacket(channel); + return; + } + default: { + logger.warn("invalid messageReceived msg:{}, connection:{}", message, channel); + } + } + } + + private short getPacketType(Object packet) { + if (packet == null) { + return PacketType.UNKNOWN; + } + + if (packet instanceof Packet) { + return ((Packet) packet).getPacketType(); + } + + return PacketType.UNKNOWN; + } + + private void handleSend(SendPacket sendPacket) { + messageListener.handleSend(sendPacket, this); + } + + private void handleRequest(RequestPacket requestPacket) { + messageListener.handleRequest(requestPacket, this); + } + + private void handleResponse(ResponsePacket responsePacket) { + this.requestManager.messageReceived(responsePacket, this); + } + + private void handleStreamEvent(StreamPacket streamPacket) { + streamChannelManager.messageReceived(streamPacket); + } + + private void handleHandshake(ControlHandshakePacket handshakepacket) { + int requestId = handshakepacket.getRequestId(); + Map handshakeData = decodeHandshakePacket(handshakepacket); + HandshakeResponseCode responseCode = messageListener.handleHandshake(handshakeData); + boolean isFirst = setChannelProperties(handshakeData); + if (isFirst) { + if (HandshakeResponseCode.DUPLEX_COMMUNICATION == responseCode) { + changeStateToRunDuplex(PinpointServerStateCode.RUN_DUPLEX); + } else if (HandshakeResponseCode.SIMPLEX_COMMUNICATION == responseCode) { + changeStateToRunSimplex(PinpointServerStateCode.RUN_SIMPLEX); + } + } + + Map responseData = createHandshakeResponse(responseCode, isFirst); + sendHandshakeResponse0(requestId, responseData); + } + + private void handleClosePacket(Channel channel) { + logger.debug("handleClosePacket channel:{}", channel); + changeStateBeingShutdown(PinpointServerStateCode.BEING_SHUTDOWN); + } + + private Map createHandshakeResponse(HandshakeResponseCode responseCode, boolean isFirst) { + HandshakeResponseCode createdCode = null; + if (isFirst) { + createdCode = responseCode; + } else { + if (HandshakeResponseCode.DUPLEX_COMMUNICATION == responseCode) { + createdCode = HandshakeResponseCode.ALREADY_DUPLEX_COMMUNICATION; + } else if (HandshakeResponseCode.SIMPLEX_COMMUNICATION == responseCode) { + createdCode = HandshakeResponseCode.ALREADY_SIMPLEX_COMMUNICATION; + } else { + createdCode = responseCode; + } + } + + Map result = new HashMap(); + result.put(ControlHandshakeResponsePacket.CODE, createdCode.getCode()); + result.put(ControlHandshakeResponsePacket.SUB_CODE, createdCode.getSubCode()); + + return result; + } + + private void sendHandshakeResponse0(int requestId, Map data) { + try { + logger.info("write HandshakeResponsePakcet. channel:{}, HandshakeResponseCode:{}.", channel, data); + + byte[] resultPayload = ControlMessageEncodingUtils.encode(data); + ControlHandshakeResponsePacket packet = new ControlHandshakeResponsePacket(requestId, resultPayload); + + channel.write(packet); + } catch (ProtocolException e) { + logger.warn(e.getMessage(), e); + } + } + + private Map decodeHandshakePacket(ControlHandshakePacket message) { + try { + byte[] payload = message.getPayload(); + Map properties = (Map) ControlMessageEncodingUtils.decode(payload); + return properties; + } catch (ProtocolException e) { + logger.warn(e.getMessage(), e); + } + + return null; + } + + @Override + public PinpointServerStateCode getCurrentStateCode() { + return state.getCurrentState(); + } + + private PinpointServerStateCode changeStateToRunWithoutHandshake(PinpointServerStateCode... skipLogicStateList) { + PinpointServerStateCode nextState = PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE; + return change0(nextState, skipLogicStateList); + } + + private PinpointServerStateCode changeStateToRunSimplex(PinpointServerStateCode... skipLogicStateList) { + PinpointServerStateCode nextState = PinpointServerStateCode.RUN_SIMPLEX; + return change0(nextState, skipLogicStateList); + } + + private PinpointServerStateCode changeStateToRunDuplex(PinpointServerStateCode... skipLogicStateList) { + PinpointServerStateCode nextState = PinpointServerStateCode.RUN_DUPLEX; + return change0(nextState, skipLogicStateList); + } + + private PinpointServerStateCode changeStateBeingShutdown(PinpointServerStateCode... skipLogicStateList) { + PinpointServerStateCode nextState = PinpointServerStateCode.BEING_SHUTDOWN; + return change0(nextState, skipLogicStateList); + } + + private PinpointServerStateCode changeStateToShutdown(PinpointServerStateCode... skipLogicStateList) { + PinpointServerStateCode nextState = PinpointServerStateCode.SHUTDOWN; + return change0(nextState, skipLogicStateList); + } + + private PinpointServerStateCode changeStateToUnexpectedShutdown(PinpointServerStateCode... skipLogicStateList) { + PinpointServerStateCode nextState = PinpointServerStateCode.UNEXPECTED_SHUTDOWN; + return change0(nextState, skipLogicStateList); + } + + private PinpointServerStateCode change0(PinpointServerStateCode nextState, PinpointServerStateCode... skipLogicStateList) { + logger.debug("Channel({}) state will be changed {}.", channel, nextState); + PinpointServerStateCode beforeState = state.changeState(nextState, skipLogicStateList); + if (beforeState == null) { + executeChangeEventHandler(this, nextState); + } else if (!isSkipedChange0(beforeState, skipLogicStateList)) { + executeChangeEventHandler(this, state.getCurrentState()); + } + + return beforeState; + } + + private boolean isSkipedChange0(PinpointServerStateCode beforeState, PinpointServerStateCode... skipLogicStateList) { + if (skipLogicStateList != null) { + for (PinpointServerStateCode skipLogicState : skipLogicStateList) { + if (beforeState == skipLogicState) { + return true; + } + } + } + return false; + } + + private void executeChangeEventHandler(DefaultPinpointServer pinpointServer, PinpointServerStateCode stateCode) { + for (ChannelStateChangeEventHandler eachListener : this.stateChangeEventListeners) { + try { + eachListener.eventPerformed(this, stateCode); + } catch (Exception e) { + eachListener.exceptionCaught(this, stateCode, e); + } + } + } + + public boolean isEnableCommunication() { + return PinpointServerStateCode.isRun(getCurrentStateCode()); + } + + public boolean isEnableDuplexCommunication() { + return PinpointServerStateCode.isRunDuplexCommunication(getCurrentStateCode()); + } + + @Override + public String toString() { + StringBuilder log = new StringBuilder(32); + log.append(objectUniqName); + log.append("("); + log.append("remote:"); + log.append(getRemoteAddress()); + log.append(", state:"); + log.append(getCurrentStateCode()); + log.append(")"); + + return log.toString(); + } + +} diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServer.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServer.java index 2faa59bfc..2b57b0796 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServer.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServer.java @@ -16,432 +16,33 @@ package com.navercorp.pinpoint.rpc.server; -import java.lang.reflect.Array; import java.net.SocketAddress; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashMap; -import java.util.List; import java.util.Map; -import java.util.concurrent.atomic.AtomicReference; -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.navercorp.pinpoint.rpc.ChannelWriteFailListenableFuture; import com.navercorp.pinpoint.rpc.Future; -import com.navercorp.pinpoint.rpc.ResponseMessage; -import com.navercorp.pinpoint.rpc.client.RequestManager; -import com.navercorp.pinpoint.rpc.client.WriteFailFutureListener; -import com.navercorp.pinpoint.rpc.control.ProtocolException; -import com.navercorp.pinpoint.rpc.packet.ControlHandshakePacket; -import com.navercorp.pinpoint.rpc.packet.ControlHandshakeResponsePacket; -import com.navercorp.pinpoint.rpc.packet.HandshakeResponseCode; -import com.navercorp.pinpoint.rpc.packet.Packet; -import com.navercorp.pinpoint.rpc.packet.PacketType; import com.navercorp.pinpoint.rpc.packet.RequestPacket; -import com.navercorp.pinpoint.rpc.packet.ResponsePacket; -import com.navercorp.pinpoint.rpc.packet.SendPacket; -import com.navercorp.pinpoint.rpc.packet.ServerClosePacket; -import com.navercorp.pinpoint.rpc.packet.stream.StreamPacket; -import com.navercorp.pinpoint.rpc.server.handler.ChannelStateChangeEventHandler; -import com.navercorp.pinpoint.rpc.server.handler.DoNothingChannelStateEventHandler; import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelContext; import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelMessageListener; -import com.navercorp.pinpoint.rpc.stream.StreamChannelContext; -import com.navercorp.pinpoint.rpc.stream.StreamChannelManager; -import com.navercorp.pinpoint.rpc.util.ClassUtils; -import com.navercorp.pinpoint.rpc.util.ControlMessageEncodingUtils; -import com.navercorp.pinpoint.rpc.util.IDGenerator; -import com.navercorp.pinpoint.rpc.util.ListUtils; /** * @author Taejin Koo */ -public class PinpointServer implements WritablePinpointServer { +public interface PinpointServer { + void send(byte[] payload); - private final Logger logger = LoggerFactory.getLogger(this.getClass()); - - private final Channel channel; - private final RequestManager requestManager; - - private final PinpointServerState state; - - private final ServerMessageListener messageListener; - - private final List stateChangeEventListeners; - - private final StreamChannelManager streamChannelManager; - - private final AtomicReference> properties = new AtomicReference>(); - - private final String objectUniqName; + Future request(byte[] payload); - private final ChannelFutureListener serverCloseWriteListener; - private final ChannelFutureListener responseWriteFailListener; - - public PinpointServer(Channel channel, PinpointServerConfig serverConfig) { - this(channel, serverConfig, DoNothingChannelStateEventHandler.INSTANCE); - } + void response(RequestPacket requestPacket, byte[] payload); + void response(int requestId, byte[] payload); - public PinpointServer(Channel channel, PinpointServerConfig serverConfig, ChannelStateChangeEventHandler... stateChangeEventListeners) { - this.channel = channel; + ClientStreamChannelContext createStream(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener); - this.messageListener = serverConfig.getMessageListener(); + void messageReceived(Object message); - StreamChannelManager streamChannelManager = new StreamChannelManager(channel, IDGenerator.createEvenIdGenerator(), serverConfig.getStreamMessageListener()); - this.streamChannelManager = streamChannelManager; + PinpointServerStateCode getCurrentStateCode(); - this.stateChangeEventListeners = new ArrayList(Array.getLength(stateChangeEventListeners) + 1); - ListUtils.addIfValueNotNull(this.stateChangeEventListeners, serverConfig.getStateChangeEventHandler()); - ListUtils.addAllExceptNullValue(this.stateChangeEventListeners, stateChangeEventListeners); + SocketAddress getRemoteAddress(); - RequestManager requestManager = new RequestManager(serverConfig.getRequestManagerTimer(), serverConfig.getDefaultRequestTimeout()); - this.requestManager = requestManager; - - this.state = new PinpointServerState(); - - this.objectUniqName = ClassUtils.simpleClassNameAndHashCodeString(this); - - this.serverCloseWriteListener = new WriteFailFutureListener(logger, objectUniqName + " sendClosePacket() write fail.", "serverClosePacket write success"); - this.responseWriteFailListener = new WriteFailFutureListener(logger, objectUniqName + " response() write fail."); - } - - public void start() { - changeStateToRunWithoutHandshake(); - } - - public void stop() { - if (PinpointServerStateCode.BEING_SHUTDOWN == getCurrentStateCode()) { - changeStateToShutdown(); - } else { - changeStateToUnexpectedShutdown(); - } - - if (this.channel.isConnected()) { - channel.close(); - } - - streamChannelManager.close(); - } - - @Override - public void send(byte[] payload) { - SendPacket send = new SendPacket(payload); - this.channel.write(send); - } - - @Override - public Future request(byte[] payload) { - if (payload == null) { - throw new NullPointerException("requestMessage must not be null"); - } - RequestPacket requestPacket = new RequestPacket(payload); - - ChannelWriteFailListenableFuture messageFuture = this.requestManager.register(requestPacket); - - ChannelFuture write = this.channel.write(requestPacket); - write.addListener(messageFuture); - - return messageFuture; - } - - @Override - public void response(RequestPacket requestPacket, byte[] payload) { - if (requestPacket == null) { - throw new NullPointerException("requestPacket must not be null"); - } - - response(requestPacket.getRequestId(), payload); - } - - @Override - public void response(int requestId, byte[] payload) { - ResponsePacket responsePacket = new ResponsePacket(requestId, payload); - ChannelFuture write = this.channel.write(responsePacket); - write.addListener(responseWriteFailListener); - } - - public void receiveResponsePacket(ResponsePacket packet) { - this.requestManager.messageReceived(packet, this); - } - - public StreamChannelContext getStreamChannel(int channelId) { - return streamChannelManager.findStreamChannel(channelId); - } - - @Override - public ClientStreamChannelContext createStream(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener) { - return streamChannelManager.openStreamChannel(payload, clientStreamChannelMessageListener); - } - - public void closeAllStreamChannel() { - streamChannelManager.close(); - } - - public Map getChannelProperties() { - Map properties = this.properties.get(); - return properties == null ? Collections.emptyMap() : properties; - } - - public boolean setChannelProperties(Map value) { - if (value == null) { - return false; - } - - return this.properties.compareAndSet(null, Collections.unmodifiableMap(value)); - } - - @Override - public SocketAddress getRemoteAddress() { - return channel.getRemoteAddress(); - } - - public ChannelFuture sendClosePacket() { - logger.info("sendServerClosedPacket start"); - - PinpointServerStateCode errorCode = changeStateBeingShutdown(); - - if (errorCode == null) { - final ChannelFuture writeFuture = this.channel.write(ServerClosePacket.DEFAULT_SERVER_CLOSE_PACKET); - writeFuture.addListener(serverCloseWriteListener); - - logger.info("sendServerClosedPacket end"); - return writeFuture; - } else { - logger.info("sendServerClosedPacket fail. Error: change state failed."); - return null; - } - } - - public void messageReceived(Object message) { - // TODO Auto-generated method stub - if (!PinpointServerStateCode.isRun(getCurrentStateCode())) { - // FIXME need change rules. - // as-is : do nothing when state is not run. - // candidate : close channel when state is not run. - logger.warn("{} messageReceived:{} from IllegalState this message will be ignore.", this, message); - return; - } - - final short packetType = getPacketType(message); - switch (packetType) { - case PacketType.APPLICATION_SEND: { - handleSend((SendPacket) message); - return; - } - case PacketType.APPLICATION_REQUEST: { - handleRequest((RequestPacket) message); - return; - } - case PacketType.APPLICATION_RESPONSE: { - handleResponse((ResponsePacket) message); - return; - } - case PacketType.APPLICATION_STREAM_CREATE: - case PacketType.APPLICATION_STREAM_CLOSE: - 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: - handleStreamEvent((StreamPacket) message); - return; - case PacketType.CONTROL_HANDSHAKE: - handleHandshake((ControlHandshakePacket) message); - return; - case PacketType.CONTROL_CLIENT_CLOSE: { - handleClosePacket(channel); - return; - } - default: { - logger.warn("invalid messageReceived msg:{}, connection:{}", message, channel); - } - } - } - - private short getPacketType(Object packet) { - if (packet == null) { - return PacketType.UNKNOWN; - } - - if (packet instanceof Packet) { - return ((Packet) packet).getPacketType(); - } - - return PacketType.UNKNOWN; - } - - private void handleSend(SendPacket sendPacket) { - messageListener.handleSend(sendPacket, this); - } - - private void handleRequest(RequestPacket requestPacket) { - messageListener.handleRequest(requestPacket, this); - } - - private void handleResponse(ResponsePacket responsePacket) { - this.requestManager.messageReceived(responsePacket, this); - } - - private void handleStreamEvent(StreamPacket streamPacket) { - streamChannelManager.messageReceived(streamPacket); - } - - private void handleHandshake(ControlHandshakePacket handshakepacket) { - int requestId = handshakepacket.getRequestId(); - Map handshakeData = decodeHandshakePacket(handshakepacket); - HandshakeResponseCode responseCode = messageListener.handleHandshake(handshakeData); - boolean isFirst = setChannelProperties(handshakeData); - if (isFirst) { - if (HandshakeResponseCode.DUPLEX_COMMUNICATION == responseCode) { - changeStateToRunDuplex(PinpointServerStateCode.RUN_DUPLEX); - } else if (HandshakeResponseCode.SIMPLEX_COMMUNICATION == responseCode) { - changeStateToRunSimplex(PinpointServerStateCode.RUN_SIMPLEX); - } - } - - Map responseData = createHandshakeResponse(responseCode, isFirst); - sendHandshakeResponse0(requestId, responseData); - } - - private void handleClosePacket(Channel channel) { - logger.debug("handleClosePacket channel:{}", channel); - changeStateBeingShutdown(PinpointServerStateCode.BEING_SHUTDOWN); - } - - private Map createHandshakeResponse(HandshakeResponseCode responseCode, boolean isFirst) { - HandshakeResponseCode createdCode = null; - if (isFirst) { - createdCode = responseCode; - } else { - if (HandshakeResponseCode.DUPLEX_COMMUNICATION == responseCode) { - createdCode = HandshakeResponseCode.ALREADY_DUPLEX_COMMUNICATION; - } else if (HandshakeResponseCode.SIMPLEX_COMMUNICATION == responseCode) { - createdCode = HandshakeResponseCode.ALREADY_SIMPLEX_COMMUNICATION; - } else { - createdCode = responseCode; - } - } - - Map result = new HashMap(); - result.put(ControlHandshakeResponsePacket.CODE, createdCode.getCode()); - result.put(ControlHandshakeResponsePacket.SUB_CODE, createdCode.getSubCode()); - - return result; - } - - private void sendHandshakeResponse0(int requestId, Map data) { - try { - logger.info("write HandshakeResponsePakcet. channel:{}, HandshakeResponseCode:{}.", channel, data); - - byte[] resultPayload = ControlMessageEncodingUtils.encode(data); - ControlHandshakeResponsePacket packet = new ControlHandshakeResponsePacket(requestId, resultPayload); - - channel.write(packet); - } catch (ProtocolException e) { - logger.warn(e.getMessage(), e); - } - } - - private Map decodeHandshakePacket(ControlHandshakePacket message) { - try { - byte[] payload = message.getPayload(); - Map properties = (Map) ControlMessageEncodingUtils.decode(payload); - return properties; - } catch (ProtocolException e) { - logger.warn(e.getMessage(), e); - } - - return null; - } - - public PinpointServerStateCode getCurrentStateCode() { - return state.getCurrentState(); - } - - private PinpointServerStateCode changeStateToRunWithoutHandshake(PinpointServerStateCode... skipLogicStateList) { - PinpointServerStateCode nextState = PinpointServerStateCode.RUN_WITHOUT_HANDSHAKE; - return change0(nextState, skipLogicStateList); - } - - private PinpointServerStateCode changeStateToRunSimplex(PinpointServerStateCode... skipLogicStateList) { - PinpointServerStateCode nextState = PinpointServerStateCode.RUN_SIMPLEX; - return change0(nextState, skipLogicStateList); - } - - private PinpointServerStateCode changeStateToRunDuplex(PinpointServerStateCode... skipLogicStateList) { - PinpointServerStateCode nextState = PinpointServerStateCode.RUN_DUPLEX; - return change0(nextState, skipLogicStateList); - } - - private PinpointServerStateCode changeStateBeingShutdown(PinpointServerStateCode... skipLogicStateList) { - PinpointServerStateCode nextState = PinpointServerStateCode.BEING_SHUTDOWN; - return change0(nextState, skipLogicStateList); - } - - private PinpointServerStateCode changeStateToShutdown(PinpointServerStateCode... skipLogicStateList) { - PinpointServerStateCode nextState = PinpointServerStateCode.SHUTDOWN; - return change0(nextState, skipLogicStateList); - } - - private PinpointServerStateCode changeStateToUnexpectedShutdown(PinpointServerStateCode... skipLogicStateList) { - PinpointServerStateCode nextState = PinpointServerStateCode.UNEXPECTED_SHUTDOWN; - return change0(nextState, skipLogicStateList); - } - - private PinpointServerStateCode change0(PinpointServerStateCode nextState, PinpointServerStateCode... skipLogicStateList) { - logger.debug("Channel({}) state will be changed {}.", channel, nextState); - PinpointServerStateCode beforeState = state.changeState(nextState, skipLogicStateList); - if (beforeState == null) { - executeChangeEventHandler(this, nextState); - } else if (!isSkipedChange0(beforeState, skipLogicStateList)) { - executeChangeEventHandler(this, state.getCurrentState()); - } - - return beforeState; - } - - private boolean isSkipedChange0(PinpointServerStateCode beforeState, PinpointServerStateCode... skipLogicStateList) { - if (skipLogicStateList != null) { - for (PinpointServerStateCode skipLogicState : skipLogicStateList) { - if (beforeState == skipLogicState) { - return true; - } - } - } - return false; - } - - private void executeChangeEventHandler(PinpointServer pinpointServer, PinpointServerStateCode stateCode) { - for (ChannelStateChangeEventHandler eachListener : this.stateChangeEventListeners) { - try { - eachListener.eventPerformed(this, stateCode); - } catch (Exception e) { - eachListener.exceptionCaught(this, stateCode, e); - } - } - } - - public boolean isEnableDuplexCommunication() { - return PinpointServerStateCode.isRunDuplexCommunication(getCurrentStateCode()); - } - - @Override - public String toString() { - StringBuilder log = new StringBuilder(32); - log.append(objectUniqName); - log.append("("); - log.append("remote:"); - log.append(getRemoteAddress()); - log.append(", state:"); - log.append(getCurrentStateCode()); - log.append(")"); - - return log.toString(); - } + Map getChannelProperties(); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerAcceptor.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerAcceptor.java index 32d0d3cf5..be0eb2cb2 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerAcceptor.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/PinpointServerAcceptor.java @@ -154,8 +154,8 @@ public class PinpointServerAcceptor implements PinpointServerConfig { sendPing(); } - private PinpointServer createPinpointServer(Channel channel) { - PinpointServer pinpointServer = new PinpointServer(channel, this); + private DefaultPinpointServer createPinpointServer(Channel channel) { + DefaultPinpointServer pinpointServer = new DefaultPinpointServer(channel, this); return pinpointServer; } @@ -304,7 +304,7 @@ public class PinpointServerAcceptor implements PinpointServerConfig { private void sendServerClosePacket() { for (Channel channel : channelGroup) { - PinpointServer pinpointServer = (PinpointServer) channel.getAttachment(); + DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment(); if (pinpointServer != null) { pinpointServer.sendClosePacket(); @@ -314,17 +314,17 @@ public class PinpointServerAcceptor implements PinpointServerConfig { private void closePinpointServer() { for (Channel channel : channelGroup) { - PinpointServer pinpointServer = (PinpointServer) channel.getAttachment(); + DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment(); pinpointServer.sendClosePacket(); } } - public List getWritableServerList() { - List pinpointServerList = new ArrayList(); + public List getWritableServerList() { + List pinpointServerList = new ArrayList(); for (Channel channel : channelGroup) { - PinpointServer pinpointServer = (PinpointServer) channel.getAttachment(); + DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment(); if (pinpointServer != null) { if (pinpointServer.isEnableDuplexCommunication()) { pinpointServerList.add(pinpointServer); @@ -358,7 +358,7 @@ public class PinpointServerAcceptor implements PinpointServerConfig { return; } - PinpointServer pinpointServer = createPinpointServer(channel); + DefaultPinpointServer pinpointServer = createPinpointServer(channel); channel.setAttachment(pinpointServer); channelGroup.add(channel); @@ -372,7 +372,7 @@ public class PinpointServerAcceptor implements PinpointServerConfig { public void channelDisconnected(ChannelHandlerContext ctx, ChannelStateEvent e) throws Exception { final Channel channel = e.getChannel(); - PinpointServer pinpointServer = (PinpointServer) channel.getAttachment(); + DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment(); if (pinpointServer != null) { pinpointServer.stop(); } @@ -396,7 +396,7 @@ public class PinpointServerAcceptor implements PinpointServerConfig { public void messageReceived(ChannelHandlerContext ctx, MessageEvent e) throws Exception { final Channel channel = e.getChannel(); - PinpointServer pinpointServer = (PinpointServer) channel.getAttachment(); + DefaultPinpointServer pinpointServer = (DefaultPinpointServer) channel.getAttachment(); if (pinpointServer != null) { Object message = e.getMessage(); diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java index e1440a7dc..2a490570d 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/ServerMessageListener.java @@ -25,9 +25,9 @@ import com.navercorp.pinpoint.rpc.server.handler.HandshakerHandler; */ public interface ServerMessageListener extends HandshakerHandler { - void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer); + void handleSend(SendPacket sendPacket, PinpointServer pinpointServer); // TODO make another tcp channel in case of exposed channel. - void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer); + void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java index 049a33eba..ebb84171f 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/SimpleLoggingServerMessageListener.java @@ -36,12 +36,12 @@ public class SimpleLoggingServerMessageListener implements ServerMessageListener public static final SimpleLoggingServerMessageListener LISTENER = new SimpleLoggingServerMessageListener(); @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { logger.info("handlerSend {} {}", sendPacket, pinpointServer); } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { logger.info("handlerRequest {} {}", requestPacket, pinpointServer); } diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/WritablePinpointServer.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/WritablePinpointServer.java deleted file mode 100644 index cccbb84d0..000000000 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/WritablePinpointServer.java +++ /dev/null @@ -1,45 +0,0 @@ -/* - * Copyright 2014 NAVER Corp. - * - * Licensed under the Apache License, Version 2.0 (the "License"); - * you may not use this file except in compliance with the License. - * You may obtain a copy of the License at - * - * http://www.apache.org/licenses/LICENSE-2.0 - * - * Unless required by applicable law or agreed to in writing, software - * distributed under the License is distributed on an "AS IS" BASIS, - * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. - * See the License for the specific language governing permissions and - * limitations under the License. - */ - -package com.navercorp.pinpoint.rpc.server; - -import java.net.SocketAddress; -import java.util.Map; - -import com.navercorp.pinpoint.rpc.Future; -import com.navercorp.pinpoint.rpc.packet.RequestPacket; -import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelContext; -import com.navercorp.pinpoint.rpc.stream.ClientStreamChannelMessageListener; - -/** - * @author Taejin Koo - */ -public interface WritablePinpointServer { - - void send(byte[] payload); - - Future request(byte[] payload); - - void response(RequestPacket requestPacket, byte[] payload); - void response(int requestId, byte[] payload); - - ClientStreamChannelContext createStream(byte[] payload, ClientStreamChannelMessageListener clientStreamChannelMessageListener); - - SocketAddress getRemoteAddress(); - - Map getChannelProperties(); - -} diff --git a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/handler/DoNothingChannelStateEventHandler.java b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/handler/DoNothingChannelStateEventHandler.java index d684ddb6e..bfa81c68c 100644 --- a/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/handler/DoNothingChannelStateEventHandler.java +++ b/rpc/src/main/java/com/navercorp/pinpoint/rpc/server/handler/DoNothingChannelStateEventHandler.java @@ -33,14 +33,12 @@ public class DoNothingChannelStateEventHandler implements ChannelStateChangeEven @Override public void eventPerformed(PinpointServer pinpointServer, PinpointServerStateCode stateCode) { - logger.info("{} eventPerformed {}:{}", this.getClass().getSimpleName(), pinpointServer, stateCode); + logger.info("{} eventPerformed(). pinpointServer:{}, code:{}", this.getClass().getSimpleName(), pinpointServer, stateCode); } @Override public void exceptionCaught(PinpointServer pinpointServer, PinpointServerStateCode stateCode, Throwable e) { - if (logger.isWarnEnabled()) { - logger.warn(this.getClass().getSimpleName() + " exception occured. Error: " + e.getMessage() + "." , e); - } + logger.warn("{} exceptionCaught(). pinpointServer:{}, code:{}. Error: {}.", this.getClass().getSimpleName(), pinpointServer, stateCode, e.getMessage(), e); } } diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/client/ClientMessageListenerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/client/ClientMessageListenerTest.java index 0af8382dc..45bb13057 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/client/ClientMessageListenerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/client/ClientMessageListenerTest.java @@ -29,7 +29,7 @@ import org.slf4j.LoggerFactory; import com.navercorp.pinpoint.rpc.packet.HandshakeResponseCode; import com.navercorp.pinpoint.rpc.packet.HandshakeResponseType; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.rpc.server.SimpleLoggingServerMessageListener; import com.navercorp.pinpoint.rpc.util.PinpointRPCTestUtils; import com.navercorp.pinpoint.rpc.util.PinpointRPCTestUtils.EchoClientListener; @@ -56,12 +56,12 @@ public class ClientMessageListenerTest { PinpointSocket socket = clientSocketFactory.connect("127.0.0.1", bindPort); Thread.sleep(500); - List writableServerList = serverAcceptor.getWritableServerList(); + List writableServerList = serverAcceptor.getWritableServerList(); if (writableServerList.size() != 1) { Assert.fail(); } - WritablePinpointServer writableServer = writableServerList.get(0); + PinpointServer writableServer = writableServerList.get(0); assertSendMessage(writableServer, "simple", echoMessageListener); assertRequestMessage(writableServer, "request", echoMessageListener); @@ -88,15 +88,15 @@ public class ClientMessageListenerTest { Thread.sleep(500); - List writableServerList = serverAcceptor.getWritableServerList(); + List writableServerList = serverAcceptor.getWritableServerList(); if (writableServerList.size() != 2) { Assert.fail(); } - WritablePinpointServer writableServer = writableServerList.get(0); + PinpointServer writableServer = writableServerList.get(0); assertRequestMessage(writableServer, "socket1", null); - WritablePinpointServer writableServer2 = writableServerList.get(1); + PinpointServer writableServer2 = writableServerList.get(1); assertRequestMessage(writableServer2, "socket2", null); Assert.assertEquals(1, echoMessageListener1.getRequestPacketRepository().size()); @@ -111,14 +111,14 @@ public class ClientMessageListenerTest { } } - private void assertSendMessage(WritablePinpointServer writableServer, String message, EchoClientListener echoMessageListener) throws InterruptedException { + private void assertSendMessage(PinpointServer writableServer, String message, EchoClientListener echoMessageListener) throws InterruptedException { writableServer.send(message.getBytes()); Thread.sleep(100); Assert.assertEquals(message, new String(echoMessageListener.getSendPacketRepository().get(0).getPayload())); } - private void assertRequestMessage(WritablePinpointServer writableServer, String message, EchoClientListener echoMessageListener) throws InterruptedException { + private void assertRequestMessage(PinpointServer writableServer, String message, EchoClientListener echoMessageListener) throws InterruptedException { byte[] response = PinpointRPCTestUtils.request(writableServer, message.getBytes()); Assert.assertEquals(message, new String(response)); diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java index c00f13731..2fd1840ea 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/ControlPacketServerTest.java @@ -234,12 +234,12 @@ public class ControlPacketServerTest { class SimpleListener implements ServerMessageListener { @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { logger.info("handlerRequest {} {}", requestPacket, pinpointServer); pinpointServer.response(requestPacket, requestPacket.getPayload()); diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventHandlerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventHandlerTest.java index ec6662436..c2a6610ae 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventHandlerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/EventHandlerTest.java @@ -193,12 +193,12 @@ public class EventHandlerTest { class SimpleListener implements ServerMessageListener { @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { logger.info("handlerRequest {}", requestPacket); pinpointServer.response(requestPacket, requestPacket.getPayload()); diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/HandshakeTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/HandshakeTest.java index 5f1741438..d46229db7 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/HandshakeTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/HandshakeTest.java @@ -72,7 +72,7 @@ public class HandshakeTest { Thread.sleep(500); - List writableServerList = serverAcceptor.getWritableServerList(); + List writableServerList = serverAcceptor.getWritableServerList(); if (writableServerList.size() != 2) { Assert.fail(); } @@ -98,7 +98,7 @@ public class HandshakeTest { PinpointSocket socket = clientSocketFactory1.connect("127.0.0.1", bindPort); Thread.sleep(500); - WritablePinpointServer writableServer = getWritableServer("application", "agent", (Long) params.get(AgentHandshakePropertyType.START_TIMESTAMP.getName()), serverAcceptor.getWritableServerList()); + PinpointServer writableServer = getWritableServer("application", "agent", (Long) params.get(AgentHandshakePropertyType.START_TIMESTAMP.getName()), serverAcceptor.getWritableServerList()); Assert.assertNotNull(writableServer); writableServer = getWritableServer("application", "agent", (Long) params.get(AgentHandshakePropertyType.START_TIMESTAMP.getName()) + 1, serverAcceptor.getWritableServerList()); @@ -135,7 +135,7 @@ public class HandshakeTest { Assert.assertTrue(handshaker.isFinished()); } - private WritablePinpointServer getWritableServer(String applicationName, String agentId, long startTimeMillis, List writableServerList) { + private PinpointServer getWritableServer(String applicationName, String agentId, long startTimeMillis, List writableServerList) { if (applicationName == null) { return null; } @@ -148,9 +148,9 @@ public class HandshakeTest { return null; } - List result = new ArrayList(); + List result = new ArrayList(); - for (WritablePinpointServer writableServer : writableServerList) { + for (PinpointServer writableServer : writableServerList) { Map agentProperties = writableServer.getChannelProperties(); if (!applicationName.equals(agentProperties.get(AgentHandshakePropertyType.APPLICATION_NAME.getName()))) { diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java index 463c04861..7cb4ab564 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/server/TestSeverMessageListener.java @@ -39,12 +39,12 @@ public class TestSeverMessageListener implements ServerMessageListener { private List sendMessageList = new ArrayList(); @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { logger.debug("sendPacket:{} channel:{}", sendPacket, pinpointServer); } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { logger.debug("requestPacket:{} channel:{}", requestPacket, pinpointServer); pinpointServer.response(requestPacket, requestPacket.getPayload()); diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java index 5ee590ae0..3a78f982b 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/stream/StreamChannelManagerTest.java @@ -35,7 +35,7 @@ import com.navercorp.pinpoint.rpc.client.SimpleLoggingMessageListener; import com.navercorp.pinpoint.rpc.packet.stream.StreamClosePacket; import com.navercorp.pinpoint.rpc.packet.stream.StreamCreatePacket; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; import com.navercorp.pinpoint.rpc.server.TestSeverMessageListener; import com.navercorp.pinpoint.rpc.util.PinpointRPCTestUtils; @@ -151,10 +151,10 @@ public class StreamChannelManagerTest { Thread.sleep(100); - List writableServerList = serverAcceptor.getWritableServerList(); + List writableServerList = serverAcceptor.getWritableServerList(); Assert.assertEquals(1, writableServerList.size()); - WritablePinpointServer writableServer = writableServerList.get(0); + PinpointServer writableServer = writableServerList.get(0); RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4); @@ -253,10 +253,10 @@ public class StreamChannelManagerTest { Thread.sleep(100); - List writableServerList = serverAcceptor.getWritableServerList(); + List writableServerList = serverAcceptor.getWritableServerList(); Assert.assertEquals(1, writableServerList.size()); - WritablePinpointServer writableServer = writableServerList.get(0); + PinpointServer writableServer = writableServerList.get(0); RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4); diff --git a/rpc/src/test/java/com/navercorp/pinpoint/rpc/util/PinpointRPCTestUtils.java b/rpc/src/test/java/com/navercorp/pinpoint/rpc/util/PinpointRPCTestUtils.java index 41c0bdfc9..a63432018 100644 --- a/rpc/src/test/java/com/navercorp/pinpoint/rpc/util/PinpointRPCTestUtils.java +++ b/rpc/src/test/java/com/navercorp/pinpoint/rpc/util/PinpointRPCTestUtils.java @@ -40,7 +40,7 @@ import com.navercorp.pinpoint.rpc.packet.ResponsePacket; import com.navercorp.pinpoint.rpc.packet.SendPacket; import com.navercorp.pinpoint.rpc.server.AgentHandshakePropertyType; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; public final class PinpointRPCTestUtils { @@ -119,7 +119,7 @@ public final class PinpointRPCTestUtils { return socketFactory; } - public static byte[] request(WritablePinpointServer writableServer, byte[] message) { + public static byte[] request(PinpointServer writableServer, byte[] message) { Future future = writableServer.request(message); future.await(); return future.getResult().getMessage(); @@ -187,12 +187,12 @@ public final class PinpointRPCTestUtils { private final List requestPacketRepository = new ArrayList(); @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { sendPacketRepository.add(sendPacket); } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { requestPacketRepository.add(requestPacket); logger.info("handlerRequest {}", requestPacket); diff --git a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/TCommandTypeVersion.java b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/TCommandTypeVersion.java index 676800999..19b3d55ea 100644 --- a/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/TCommandTypeVersion.java +++ b/thrift/src/main/java/com/navercorp/pinpoint/thrift/io/TCommandTypeVersion.java @@ -33,6 +33,8 @@ public enum TCommandTypeVersion { V_1_0_3("1.0.3", V_1_0_3_SNAPSHOT), V_1_0_4_SNAPSHOT("1.0.4-SNAPSHOT", V_1_0_3), V_1_0_4("1.0.4", V_1_0_4_SNAPSHOT), + V_1_1_0("1.1.0-SNAPSHOT", V_1_0_4_SNAPSHOT), + UNKNOWN("UNKNOWN"); diff --git a/web/src/main/java/com/navercorp/pinpoint/web/controller/CommandController.java b/web/src/main/java/com/navercorp/pinpoint/web/controller/CommandController.java index 93627bccd..e14506960 100644 --- a/web/src/main/java/com/navercorp/pinpoint/web/controller/CommandController.java +++ b/web/src/main/java/com/navercorp/pinpoint/web/controller/CommandController.java @@ -34,7 +34,7 @@ import org.springframework.web.servlet.ModelAndView; import com.navercorp.pinpoint.rpc.Future; import com.navercorp.pinpoint.rpc.ResponseMessage; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.thrift.dto.TResult; import com.navercorp.pinpoint.thrift.dto.command.TCommandEcho; import com.navercorp.pinpoint.thrift.dto.command.TCommandThreadDump; @@ -72,7 +72,7 @@ public class CommandController { public ModelAndView echo(@RequestParam("application") String applicationName, @RequestParam("agent") String agentId, @RequestParam("startTimeStamp") long startTimeStamp, @RequestParam("message") String message) throws TException { - WritablePinpointServer collector = socketManager.getCollector(applicationName, agentId, startTimeStamp); + PinpointServer collector = socketManager.getCollector(applicationName, agentId, startTimeStamp); if (collector == null) { return createResponse(false, String.format("Can't find suitable PinpointServer(%s/%s/%d).", applicationName, agentId, startTimeStamp)); @@ -119,7 +119,7 @@ public class CommandController { public ModelAndView echo(@RequestParam("application") String applicationName, @RequestParam("agent") String agentId, @RequestParam("startTimeStamp") long startTimeStamp) throws TException { - WritablePinpointServer collector = socketManager.getCollector(applicationName, agentId, startTimeStamp); + PinpointServer collector = socketManager.getCollector(applicationName, agentId, startTimeStamp); if (collector == null) { return createResponse(false, String.format("Can't find suitable PinpointServer(%s/%s/%d).", applicationName, agentId, startTimeStamp)); diff --git a/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java b/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java index fd56de950..99213595d 100644 --- a/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java +++ b/web/src/main/java/com/navercorp/pinpoint/web/server/PinpointSocketManager.java @@ -37,7 +37,7 @@ import com.navercorp.pinpoint.rpc.packet.RequestPacket; import com.navercorp.pinpoint.rpc.packet.SendPacket; import com.navercorp.pinpoint.rpc.server.PinpointServerAcceptor; import com.navercorp.pinpoint.rpc.server.ServerMessageListener; -import com.navercorp.pinpoint.rpc.server.WritablePinpointServer; +import com.navercorp.pinpoint.rpc.server.PinpointServer; import com.navercorp.pinpoint.web.cluster.ClusterManager; import com.navercorp.pinpoint.web.cluster.zookeeper.ZookeeperClusterManager; import com.navercorp.pinpoint.web.config.WebConfig; @@ -108,11 +108,11 @@ public class PinpointSocketManager { } } - public List getCollectorList() { + public List getCollectorList() { return serverAcceptor.getWritableServerList(); } - public WritablePinpointServer getCollector(String applicationName, String agentId, long startTimeStamp) { + public PinpointServer getCollector(String applicationName, String agentId, long startTimeStamp) { List agentNameList = clusterManager.getRegisteredAgentList(applicationName, agentId, startTimeStamp); // having duplicate AgentName registered is an exceptional case @@ -126,9 +126,9 @@ public class PinpointSocketManager { String agentName = agentNameList.get(0); - List collectorList = getCollectorList(); + List collectorList = getCollectorList(); - for (WritablePinpointServer collector : collectorList) { + for (PinpointServer collector : collectorList) { String id = (String) collector.getChannelProperties().get("id"); if (agentName.startsWith(id)) { return collector; @@ -172,12 +172,12 @@ public class PinpointSocketManager { private class PinpointSocketManagerHandler implements ServerMessageListener { @Override - public void handleSend(SendPacket sendPacket, WritablePinpointServer pinpointServer) { + public void handleSend(SendPacket sendPacket, PinpointServer pinpointServer) { logger.warn("Unsupport send received {} {}", sendPacket, pinpointServer); } @Override - public void handleRequest(RequestPacket requestPacket, WritablePinpointServer pinpointServer) { + public void handleRequest(RequestPacket requestPacket, PinpointServer pinpointServer) { logger.warn("Unsupport request received {} {}", requestPacket, pinpointServer); }