Improve state synchronization between server and client. #136

WritablePinpointServer -> PinpointServer
PinpointServer -> DefaultPinpointServer
This commit is contained in:
kr14910
2015-02-16 14:37:00 +09:00
parent 47bc41663e
commit ef45fa59cf
25 changed files with 574 additions and 536 deletions
@@ -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
@@ -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<String> getClusterData() {
@@ -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");
}
@@ -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) {
@@ -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);
}
@@ -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);
}
@@ -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());
}
@@ -116,10 +116,10 @@ public class RequestManager {
final int requestId = responsePacket.getRequestId();
final DefaultFuture<ResponseMessage> 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();
@@ -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<ChannelStateChangeEventHandler> stateChangeEventListeners;
private final StreamChannelManager streamChannelManager;
private final AtomicReference<Map<Object, Object>> properties = new AtomicReference<Map<Object, Object>>();
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<ChannelStateChangeEventHandler>(1);
} else {
this.stateChangeEventListeners = new ArrayList<ChannelStateChangeEventHandler>(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<ResponseMessage> 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<Object, Object> getChannelProperties() {
Map<Object, Object> properties = this.properties.get();
return properties == null ? Collections.emptyMap() : properties;
}
public boolean setChannelProperties(Map<Object, Object> 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<Object, Object> 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<String, Object> responseData = createHandshakeResponse(responseCode, isFirst);
sendHandshakeResponse0(requestId, responseData);
}
private void handleClosePacket(Channel channel) {
logger.debug("handleClosePacket channel:{}", channel);
changeStateBeingShutdown(PinpointServerStateCode.BEING_SHUTDOWN);
}
private Map<String, Object> 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<String, Object> result = new HashMap<String, Object>();
result.put(ControlHandshakeResponsePacket.CODE, createdCode.getCode());
result.put(ControlHandshakeResponsePacket.SUB_CODE, createdCode.getSubCode());
return result;
}
private void sendHandshakeResponse0(int requestId, Map<String, Object> 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<Object, Object> decodeHandshakePacket(ControlHandshakePacket message) {
try {
byte[] payload = message.getPayload();
Map<Object, Object> 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();
}
}
@@ -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<ChannelStateChangeEventHandler> stateChangeEventListeners;
private final StreamChannelManager streamChannelManager;
private final AtomicReference<Map<Object, Object>> properties = new AtomicReference<Map<Object, Object>>();
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<ChannelStateChangeEventHandler>(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<ResponseMessage> 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<Object, Object> getChannelProperties() {
Map<Object, Object> properties = this.properties.get();
return properties == null ? Collections.emptyMap() : properties;
}
public boolean setChannelProperties(Map<Object, Object> 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<Object, Object> 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<String, Object> responseData = createHandshakeResponse(responseCode, isFirst);
sendHandshakeResponse0(requestId, responseData);
}
private void handleClosePacket(Channel channel) {
logger.debug("handleClosePacket channel:{}", channel);
changeStateBeingShutdown(PinpointServerStateCode.BEING_SHUTDOWN);
}
private Map<String, Object> 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<String, Object> result = new HashMap<String, Object>();
result.put(ControlHandshakeResponsePacket.CODE, createdCode.getCode());
result.put(ControlHandshakeResponsePacket.SUB_CODE, createdCode.getSubCode());
return result;
}
private void sendHandshakeResponse0(int requestId, Map<String, Object> 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<Object, Object> decodeHandshakePacket(ControlHandshakePacket message) {
try {
byte[] payload = message.getPayload();
Map<Object, Object> 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<Object, Object> getChannelProperties();
}
@@ -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<WritablePinpointServer> getWritableServerList() {
List<WritablePinpointServer> pinpointServerList = new ArrayList<WritablePinpointServer>();
public List<PinpointServer> getWritableServerList() {
List<PinpointServer> pinpointServerList = new ArrayList<PinpointServer>();
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();
@@ -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);
}
@@ -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);
}
@@ -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<Object, Object> getChannelProperties();
}
@@ -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);
}
}
@@ -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<WritablePinpointServer> writableServerList = serverAcceptor.getWritableServerList();
List<PinpointServer> 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<WritablePinpointServer> writableServerList = serverAcceptor.getWritableServerList();
List<PinpointServer> 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));
@@ -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());
@@ -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());
@@ -72,7 +72,7 @@ public class HandshakeTest {
Thread.sleep(500);
List<WritablePinpointServer> writableServerList = serverAcceptor.getWritableServerList();
List<PinpointServer> 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<WritablePinpointServer> writableServerList) {
private PinpointServer getWritableServer(String applicationName, String agentId, long startTimeMillis, List<PinpointServer> writableServerList) {
if (applicationName == null) {
return null;
}
@@ -148,9 +148,9 @@ public class HandshakeTest {
return null;
}
List<WritablePinpointServer> result = new ArrayList<WritablePinpointServer>();
List<PinpointServer> result = new ArrayList<PinpointServer>();
for (WritablePinpointServer writableServer : writableServerList) {
for (PinpointServer writableServer : writableServerList) {
Map agentProperties = writableServer.getChannelProperties();
if (!applicationName.equals(agentProperties.get(AgentHandshakePropertyType.APPLICATION_NAME.getName()))) {
@@ -39,12 +39,12 @@ public class TestSeverMessageListener implements ServerMessageListener {
private List<byte[]> sendMessageList = new ArrayList<byte[]>();
@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());
@@ -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<WritablePinpointServer> writableServerList = serverAcceptor.getWritableServerList();
List<PinpointServer> 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<WritablePinpointServer> writableServerList = serverAcceptor.getWritableServerList();
List<PinpointServer> writableServerList = serverAcceptor.getWritableServerList();
Assert.assertEquals(1, writableServerList.size());
WritablePinpointServer writableServer = writableServerList.get(0);
PinpointServer writableServer = writableServerList.get(0);
RecordedStreamChannelMessageListener clientListener = new RecordedStreamChannelMessageListener(4);
@@ -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<ResponseMessage> future = writableServer.request(message);
future.await();
return future.getResult().getMessage();
@@ -187,12 +187,12 @@ public final class PinpointRPCTestUtils {
private final List<RequestPacket> requestPacketRepository = new ArrayList<RequestPacket>();
@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);
@@ -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");
@@ -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));
@@ -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<WritablePinpointServer> getCollectorList() {
public List<PinpointServer> getCollectorList() {
return serverAcceptor.getWritableServerList();
}
public WritablePinpointServer getCollector(String applicationName, String agentId, long startTimeStamp) {
public PinpointServer getCollector(String applicationName, String agentId, long startTimeStamp) {
List<String> 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<WritablePinpointServer> collectorList = getCollectorList();
List<PinpointServer> 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);
}