mirror of
https://github.com/wahyd4/pinpoint.git
synced 2026-08-18 01:06:03 +10:00
Merge branch 'master' of sunsh318/pinpoint-2
from pull-request 200 * refs/heads/master: #22 feature: Allow stream data transfer in pinpointsocket.
This commit is contained in:
@@ -29,7 +29,21 @@ public abstract class BasicStreamPacket implements StreamPacket {
|
||||
public static final short PACKET_UNSUPPORT = 142;
|
||||
|
||||
public static final short UNKNWON_ERROR = 200;
|
||||
|
||||
|
||||
public static final short ROUTE_TYPE_ERROR = 330;
|
||||
public static final short ROUTE_TYPE_SERVER_UNSUPPORT = 331;
|
||||
public static final short ROUTE_TYPE_CLIENT = 336;
|
||||
public static final short ROUTE_TYPE_UNKOWN = 339;
|
||||
|
||||
public static final short ROUTE_PACKET_ERROR = 340;
|
||||
public static final short ROUTE_PACKET_UNKNOWN = 341;
|
||||
public static final short ROUTE_PACKET_UNSUPPORT = 342;
|
||||
|
||||
public static final short ROUTE_NOT_FOUND = 350;
|
||||
|
||||
public static final short ROUTE_CONNECTION_ERROR = 360;
|
||||
|
||||
|
||||
private static final byte[] EMPTY_PAYLOAD = new byte[0];
|
||||
|
||||
|
||||
@@ -2,6 +2,7 @@ package com.nhn.pinpoint.rpc.server;
|
||||
|
||||
import java.util.Collections;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
@@ -23,7 +24,7 @@ public class ChannelContext {
|
||||
|
||||
private final SocketChannelStateChangeEventListener stateChangeEventListener;
|
||||
|
||||
private volatile Map<Object, Object> channelProperties = Collections.emptyMap();
|
||||
private final AtomicReference<Map<Object, Object>> properties = new AtomicReference<Map<Object,Object>>();
|
||||
|
||||
public ChannelContext(SocketChannel socketChannel, StreamChannelManager streamChannelManager) {
|
||||
this(socketChannel, streamChannelManager, DoNothingChannelStateEventListener.INSTANCE);
|
||||
@@ -101,21 +102,16 @@ public class ChannelContext {
|
||||
}
|
||||
|
||||
public Map<Object, Object> getChannelProperties() {
|
||||
return channelProperties;
|
||||
Map<Object, Object> properties = this.properties.get();
|
||||
return properties == null ? Collections.emptyMap() : properties;
|
||||
}
|
||||
|
||||
public boolean setChannelProperties(Map<Object, Object> properties) {
|
||||
if (properties == null) {
|
||||
public boolean setChannelProperties(Map<Object, Object> value) {
|
||||
if (value == null) {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (this.channelProperties != Collections.emptyMap()) {
|
||||
logger.warn("Already Register ChannelProperties.({}).", this.channelProperties);
|
||||
return false;
|
||||
}
|
||||
|
||||
this.channelProperties = Collections.unmodifiableMap(properties);
|
||||
return true;
|
||||
return this.properties.compareAndSet(null, Collections.unmodifiableMap(value));
|
||||
}
|
||||
|
||||
public StreamChannelManager getStreamChannelManager() {
|
||||
|
||||
@@ -312,7 +312,7 @@ public class PinpointServerSocket extends SimpleChannelHandler {
|
||||
|
||||
private void sendHandshakeResponseMessage(int requestId, HandshakeResponseCode handShakeResponseCode, Channel channel) {
|
||||
try {
|
||||
logger.info("write HandshakeResponsePakcet channel:{}, HandshakeResponseCode:{}.", channel, handShakeResponseCode);
|
||||
logger.info("write HandshakeResponsePakcet. channel:{}, HandshakeResponseCode:{}.", channel, handShakeResponseCode);
|
||||
|
||||
Map<String, Object> result = new HashMap<String, Object>();
|
||||
result.put(ControlHandshakeResponsePacket.CODE, handShakeResponseCode.getCode());
|
||||
|
||||
+1
-1
@@ -10,6 +10,6 @@ public interface ClientStreamChannelMessageListener {
|
||||
|
||||
void handleStreamData(ClientStreamChannelContext streamChannelContext, StreamResponsePacket packet);
|
||||
|
||||
void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet);
|
||||
void handleStreamClose(ClientStreamChannelContext streamChannelContext, StreamClosePacket packet);
|
||||
|
||||
}
|
||||
|
||||
+1
-1
@@ -20,7 +20,7 @@ public class DisabledServerStreamChannelMessageListener implements ServerStreamC
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
public void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
logger.info("{} handleStreamClose unsupported operation. StreamChannel:{}, Packet:{}", this.getClass().getSimpleName(), streamChannelContext, packet);
|
||||
}
|
||||
|
||||
|
||||
+2
-2
@@ -26,7 +26,7 @@ public class LoggingStreamChannelMessageListener {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
public void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
LOGGER.info("handleStreamClose StreamChannel:{}, Packet:{}", streamChannelContext, packet);
|
||||
}
|
||||
|
||||
@@ -40,7 +40,7 @@ public class LoggingStreamChannelMessageListener {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
public void handleStreamClose(ClientStreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
LOGGER.info("handleStreamClose StreamChannel:{}, Packet:{}", streamChannelContext, packet);
|
||||
}
|
||||
|
||||
|
||||
+1
-1
@@ -10,6 +10,6 @@ public interface ServerStreamChannelMessageListener {
|
||||
|
||||
short handleStreamCreate(ServerStreamChannelContext streamChannelContext, StreamCreatePacket packet);
|
||||
|
||||
void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet);
|
||||
void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet);
|
||||
|
||||
}
|
||||
|
||||
@@ -100,11 +100,11 @@ public abstract class StreamChannel {
|
||||
this.streamChannelManager.clearResourceAndSendClose(getStreamId(), BasicStreamPacket.CHANNEL_CLOSE);
|
||||
}
|
||||
|
||||
protected Channel getChannel() {
|
||||
public Channel getChannel() {
|
||||
return channel;
|
||||
}
|
||||
|
||||
protected int getStreamId() {
|
||||
public int getStreamId() {
|
||||
return streamChannelId;
|
||||
}
|
||||
|
||||
|
||||
+1
-2
@@ -12,7 +12,6 @@ import com.nhn.pinpoint.rpc.packet.stream.StreamClosePacket;
|
||||
import com.nhn.pinpoint.rpc.packet.stream.StreamResponsePacket;
|
||||
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelContext;
|
||||
import com.nhn.pinpoint.rpc.stream.ClientStreamChannelMessageListener;
|
||||
import com.nhn.pinpoint.rpc.stream.StreamChannelContext;
|
||||
|
||||
/**
|
||||
* @author emeroad
|
||||
@@ -38,7 +37,7 @@ public class RecordedStreamChannelMessageListener implements ClientStreamChannel
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
public void handleStreamClose(ClientStreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
logger.info("handleClose {}, {}", streamChannelContext, packet);
|
||||
receivedMessageList.add(packet.getPayload());
|
||||
latch.countDown();
|
||||
|
||||
@@ -308,8 +308,8 @@ public class StreamChannelManagerTest {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void handleStreamClose(StreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
bo.removeServerStreamChannelContext((ServerStreamChannelContext) streamChannelContext);
|
||||
public void handleStreamClose(ServerStreamChannelContext streamChannelContext, StreamClosePacket packet) {
|
||||
bo.removeServerStreamChannelContext(streamChannelContext);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user