#12 thrift command 객체 추가

1. version 등록 (1.0.3-SNAPSHOT)
2. command 기능 등록
 echo (echo)
 commandTransfer (단순 req/res 방식 형태의 command message를 전달하는 기능)
 (stream은 이후 추가 예정)
This commit is contained in:
koo-taejin
2014-09-16 14:06:02 +09:00
parent 2365bbef8b
commit f3a92ee327
6 changed files with 589 additions and 135 deletions
@@ -1,124 +1,127 @@
package com.nhn.pinpoint.profiler.receiver;
import org.apache.thrift.TBase;
import org.apache.thrift.TException;
import org.apache.thrift.protocol.TCompactProtocol;
import org.apache.thrift.protocol.TProtocolFactory;
import org.jboss.netty.channel.Channel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.common.Version;
import com.nhn.pinpoint.profiler.receiver.bo.ThreadDumpBO;
import com.nhn.pinpoint.rpc.client.MessageListener;
import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.ResponsePacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.thrift.dto.TResult;
import com.nhn.pinpoint.thrift.dto.command.TCommandThreadDump;
import com.nhn.pinpoint.thrift.io.HeaderTBaseDeserializer;
import com.nhn.pinpoint.thrift.io.HeaderTBaseDeserializerFactory;
import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializer;
import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializerFactory;
import com.nhn.pinpoint.thrift.io.SerializerFactory;
import com.nhn.pinpoint.thrift.io.TBaseLocator;
import com.nhn.pinpoint.thrift.io.TCommandRegistry;
import com.nhn.pinpoint.thrift.io.TCommandTypeVersion;
/**
* @author koo.taejin
*/
public class CommandDispatcher implements MessageListener {
// 일단은 현재 스레드가 워커스레드로 되는 것을 등록 (이후에 변경하자.)
private static final TProtocolFactory DEFAULT_PROTOCOL_FACTORY = new TCompactProtocol.Factory();
private static final TBaseLocator commandTbaseLocator = new TCommandRegistry(TCommandTypeVersion.getVersion(Version.VERSION));
private final Logger logger = LoggerFactory.getLogger(this.getClass());
// 여기만 따로 TBaseLocator를 상속받아 만들어주는 것이 좋을듯
private final TBaseBOLocator locator;
private final SerializerFactory serializerFactory = new HeaderTBaseSerializerFactory(true, HeaderTBaseSerializerFactory.DEFAULT_UDP_STREAM_MAX_SIZE, DEFAULT_PROTOCOL_FACTORY, commandTbaseLocator);
private final HeaderTBaseDeserializerFactory deserializerFactory = new HeaderTBaseDeserializerFactory(DEFAULT_PROTOCOL_FACTORY, commandTbaseLocator);
public CommandDispatcher() {
TBaseBORegistry registry = new TBaseBORegistry();
registry.addBO(TCommandThreadDump.class, new ThreadDumpBO());
this.locator = registry;
}
@Override
public void handleRequest(RequestPacket packet, Channel channel) {
logger.info("MessageReceive {} {}", packet, channel);
TBase<?, ?> request = deserialize(packet.getPayload());
TBase response = null;
if (request == null) {
TResult tResult = new TResult(false);
tResult.setMessage("Unsupported Type.");
response = tResult;
} else {
TBaseRequestBO bo = locator.getRequestBO(request);
if (bo == null) {
TResult tResult = new TResult(false);
tResult.setMessage("Unsupported Listener.");
response = tResult;
} else {
response = bo.handleRequest(request);
}
}
byte[] payload = serialize(response);
if (payload != null) {
channel.write(new ResponsePacket(packet.getRequestId(), payload));
}
}
@Override
public void handleSend(SendPacket packet, Channel channel) {
logger.info("MessageReceive {} {}", packet, channel);
}
private TBase deserialize(byte[] payload) {
if (payload == null) {
logger.warn("Payload may not be null.");
return null;
}
try {
final HeaderTBaseDeserializer deserializer = deserializerFactory.createDeserializer();
TBase<?, ?> tBase = deserializer.deserialize(payload);
return tBase;
} catch (TException e) {
logger.warn(e.getMessage(), e);
}
return null;
}
private byte[] serialize(TBase result) {
if (result == null) {
logger.warn("tBase may not be null.");
return null;
}
try {
HeaderTBaseSerializer serializer = serializerFactory.createSerializer();
byte[] payload = serializer.serialize(result);
return payload;
} catch (TException e) {
logger.warn(e.getMessage(), e);
}
return null;
}
}
package com.nhn.pinpoint.profiler.receiver;
import org.apache.thrift.TBase;
import org.apache.thrift.TException;
import org.apache.thrift.protocol.TCompactProtocol;
import org.apache.thrift.protocol.TProtocolFactory;
import org.jboss.netty.channel.Channel;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.common.Version;
import com.nhn.pinpoint.profiler.receiver.bo.EchoBO;
import com.nhn.pinpoint.profiler.receiver.bo.ThreadDumpBO;
import com.nhn.pinpoint.rpc.client.MessageListener;
import com.nhn.pinpoint.rpc.packet.RequestPacket;
import com.nhn.pinpoint.rpc.packet.ResponsePacket;
import com.nhn.pinpoint.rpc.packet.SendPacket;
import com.nhn.pinpoint.thrift.dto.TResult;
import com.nhn.pinpoint.thrift.dto.command.TCommandEcho;
import com.nhn.pinpoint.thrift.dto.command.TCommandThreadDump;
import com.nhn.pinpoint.thrift.io.HeaderTBaseDeserializer;
import com.nhn.pinpoint.thrift.io.HeaderTBaseDeserializerFactory;
import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializer;
import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializerFactory;
import com.nhn.pinpoint.thrift.io.SerializerFactory;
import com.nhn.pinpoint.thrift.io.TBaseLocator;
import com.nhn.pinpoint.thrift.io.TCommandRegistry;
import com.nhn.pinpoint.thrift.io.TCommandTypeVersion;
/**
* @author koo.taejin
*/
public class CommandDispatcher implements MessageListener {
// 일단은 현재 스레드가 워커스레드로 되는 것을 등록 (이후에 변경하자.)
private static final TProtocolFactory DEFAULT_PROTOCOL_FACTORY = new TCompactProtocol.Factory();
private static final TBaseLocator commandTbaseLocator = new TCommandRegistry(TCommandTypeVersion.getVersion(Version.VERSION));
private final Logger logger = LoggerFactory.getLogger(this.getClass());
// 여기만 따로 TBaseLocator를 상속받아 만들어주는 것이 좋을듯
private final TBaseBOLocator locator;
private final SerializerFactory serializerFactory = new HeaderTBaseSerializerFactory(true, HeaderTBaseSerializerFactory.DEFAULT_UDP_STREAM_MAX_SIZE, DEFAULT_PROTOCOL_FACTORY, commandTbaseLocator);
private final HeaderTBaseDeserializerFactory deserializerFactory = new HeaderTBaseDeserializerFactory(DEFAULT_PROTOCOL_FACTORY, commandTbaseLocator);
public CommandDispatcher() {
TBaseBORegistry registry = new TBaseBORegistry();
registry.addBO(TCommandThreadDump.class, new ThreadDumpBO());
registry.addBO(TCommandEcho.class, new EchoBO());
this.locator = registry;
}
@Override
public void handleRequest(RequestPacket packet, Channel channel) {
logger.info("MessageReceive {} {}", packet, channel);
TBase<?, ?> request = deserialize(packet.getPayload());
TBase response = null;
if (request == null) {
TResult tResult = new TResult(false);
tResult.setMessage("Unsupported Type.");
response = tResult;
} else {
TBaseRequestBO bo = locator.getRequestBO(request);
if (bo == null) {
TResult tResult = new TResult(false);
tResult.setMessage("Unsupported Listener.");
response = tResult;
} else {
response = bo.handleRequest(request);
}
}
byte[] payload = serialize(response);
if (payload != null) {
channel.write(new ResponsePacket(packet.getRequestId(), payload));
}
}
@Override
public void handleSend(SendPacket packet, Channel channel) {
logger.info("MessageReceive {} {}", packet, channel);
}
private TBase deserialize(byte[] payload) {
if (payload == null) {
logger.warn("Payload may not be null.");
return null;
}
try {
final HeaderTBaseDeserializer deserializer = deserializerFactory.createDeserializer();
TBase<?, ?> tBase = deserializer.deserialize(payload);
return tBase;
} catch (TException e) {
logger.warn(e.getMessage(), e);
}
return null;
}
private byte[] serialize(TBase result) {
if (result == null) {
logger.warn("tBase may not be null.");
return null;
}
try {
HeaderTBaseSerializer serializer = serializerFactory.createSerializer();
byte[] payload = serializer.serialize(result);
return payload;
} catch (TException e) {
logger.warn(e.getMessage(), e);
}
return null;
}
}
@@ -0,0 +1,24 @@
package com.nhn.pinpoint.profiler.receiver.bo;
import org.apache.thrift.TBase;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.nhn.pinpoint.profiler.receiver.TBaseRequestBO;
import com.nhn.pinpoint.thrift.dto.command.TCommandEcho;
/**
* @author koo.taejin <kr14910>
*/
public class EchoBO implements TBaseRequestBO {
private final Logger logger = LoggerFactory.getLogger(this.getClass());
@Override
public TBase<?, ?> handleRequest(TBase tbase) {
logger.info("{} execute {}.", this, tbase);
TCommandEcho param = (TCommandEcho) tbase;
return param;
}
}
@@ -0,0 +1,385 @@
/**
* Autogenerated by Thrift Compiler (0.9.1)
*
* DO NOT EDIT UNLESS YOU ARE SURE THAT YOU KNOW WHAT YOU ARE DOING
* @generated
*/
package com.nhn.pinpoint.thrift.dto.command;
import org.apache.thrift.scheme.IScheme;
import org.apache.thrift.scheme.SchemeFactory;
import org.apache.thrift.scheme.StandardScheme;
import org.apache.thrift.scheme.TupleScheme;
import org.apache.thrift.protocol.TTupleProtocol;
import org.apache.thrift.protocol.TProtocolException;
import org.apache.thrift.EncodingUtils;
import org.apache.thrift.TException;
import org.apache.thrift.async.AsyncMethodCallback;
import org.apache.thrift.server.AbstractNonblockingServer.*;
import java.util.List;
import java.util.ArrayList;
import java.util.Map;
import java.util.HashMap;
import java.util.EnumMap;
import java.util.Set;
import java.util.HashSet;
import java.util.EnumSet;
import java.util.Collections;
import java.util.BitSet;
import java.nio.ByteBuffer;
import java.util.Arrays;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
public class TCommandEcho implements org.apache.thrift.TBase<TCommandEcho, TCommandEcho._Fields>, java.io.Serializable, Cloneable, Comparable<TCommandEcho> {
private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("TCommandEcho");
private static final org.apache.thrift.protocol.TField MESSAGE_FIELD_DESC = new org.apache.thrift.protocol.TField("message", org.apache.thrift.protocol.TType.STRING, (short)1);
private static final Map<Class<? extends IScheme>, SchemeFactory> schemes = new HashMap<Class<? extends IScheme>, SchemeFactory>();
static {
schemes.put(StandardScheme.class, new TCommandEchoStandardSchemeFactory());
schemes.put(TupleScheme.class, new TCommandEchoTupleSchemeFactory());
}
private String message; // required
/** The set of fields this struct contains, along with convenience methods for finding and manipulating them. */
public enum _Fields implements org.apache.thrift.TFieldIdEnum {
MESSAGE((short)1, "message");
private static final Map<String, _Fields> byName = new HashMap<String, _Fields>();
static {
for (_Fields field : EnumSet.allOf(_Fields.class)) {
byName.put(field.getFieldName(), field);
}
}
/**
* Find the _Fields constant that matches fieldId, or null if its not found.
*/
public static _Fields findByThriftId(int fieldId) {
switch(fieldId) {
case 1: // MESSAGE
return MESSAGE;
default:
return null;
}
}
/**
* Find the _Fields constant that matches fieldId, throwing an exception
* if it is not found.
*/
public static _Fields findByThriftIdOrThrow(int fieldId) {
_Fields fields = findByThriftId(fieldId);
if (fields == null) throw new IllegalArgumentException("Field " + fieldId + " doesn't exist!");
return fields;
}
/**
* Find the _Fields constant that matches name, or null if its not found.
*/
public static _Fields findByName(String name) {
return byName.get(name);
}
private final short _thriftId;
private final String _fieldName;
_Fields(short thriftId, String fieldName) {
_thriftId = thriftId;
_fieldName = fieldName;
}
public short getThriftFieldId() {
return _thriftId;
}
public String getFieldName() {
return _fieldName;
}
}
// isset id assignments
public static final Map<_Fields, org.apache.thrift.meta_data.FieldMetaData> metaDataMap;
static {
Map<_Fields, org.apache.thrift.meta_data.FieldMetaData> tmpMap = new EnumMap<_Fields, org.apache.thrift.meta_data.FieldMetaData>(_Fields.class);
tmpMap.put(_Fields.MESSAGE, new org.apache.thrift.meta_data.FieldMetaData("message", org.apache.thrift.TFieldRequirementType.DEFAULT,
new org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
metaDataMap = Collections.unmodifiableMap(tmpMap);
org.apache.thrift.meta_data.FieldMetaData.addStructMetaDataMap(TCommandEcho.class, metaDataMap);
}
public TCommandEcho() {
}
public TCommandEcho(
String message)
{
this();
this.message = message;
}
/**
* Performs a deep copy on <i>other</i>.
*/
public TCommandEcho(TCommandEcho other) {
if (other.isSetMessage()) {
this.message = other.message;
}
}
public TCommandEcho deepCopy() {
return new TCommandEcho(this);
}
@Override
public void clear() {
this.message = null;
}
public String getMessage() {
return this.message;
}
public void setMessage(String message) {
this.message = message;
}
public void unsetMessage() {
this.message = null;
}
/** Returns true if field message is set (has been assigned a value) and false otherwise */
public boolean isSetMessage() {
return this.message != null;
}
public void setMessageIsSet(boolean value) {
if (!value) {
this.message = null;
}
}
public void setFieldValue(_Fields field, Object value) {
switch (field) {
case MESSAGE:
if (value == null) {
unsetMessage();
} else {
setMessage((String)value);
}
break;
}
}
public Object getFieldValue(_Fields field) {
switch (field) {
case MESSAGE:
return getMessage();
}
throw new IllegalStateException();
}
/** Returns true if field corresponding to fieldID is set (has been assigned a value) and false otherwise */
public boolean isSet(_Fields field) {
if (field == null) {
throw new IllegalArgumentException();
}
switch (field) {
case MESSAGE:
return isSetMessage();
}
throw new IllegalStateException();
}
@Override
public boolean equals(Object that) {
if (that == null)
return false;
if (that instanceof TCommandEcho)
return this.equals((TCommandEcho)that);
return false;
}
public boolean equals(TCommandEcho that) {
if (that == null)
return false;
boolean this_present_message = true && this.isSetMessage();
boolean that_present_message = true && that.isSetMessage();
if (this_present_message || that_present_message) {
if (!(this_present_message && that_present_message))
return false;
if (!this.message.equals(that.message))
return false;
}
return true;
}
@Override
public int hashCode() {
return 0;
}
@Override
public int compareTo(TCommandEcho other) {
if (!getClass().equals(other.getClass())) {
return getClass().getName().compareTo(other.getClass().getName());
}
int lastComparison = 0;
lastComparison = Boolean.valueOf(isSetMessage()).compareTo(other.isSetMessage());
if (lastComparison != 0) {
return lastComparison;
}
if (isSetMessage()) {
lastComparison = org.apache.thrift.TBaseHelper.compareTo(this.message, other.message);
if (lastComparison != 0) {
return lastComparison;
}
}
return 0;
}
public _Fields fieldForId(int fieldId) {
return _Fields.findByThriftId(fieldId);
}
public void read(org.apache.thrift.protocol.TProtocol iprot) throws org.apache.thrift.TException {
schemes.get(iprot.getScheme()).getScheme().read(iprot, this);
}
public void write(org.apache.thrift.protocol.TProtocol oprot) throws org.apache.thrift.TException {
schemes.get(oprot.getScheme()).getScheme().write(oprot, this);
}
@Override
public String toString() {
StringBuilder sb = new StringBuilder("TCommandEcho(");
boolean first = true;
sb.append("message:");
if (this.message == null) {
sb.append("null");
} else {
sb.append(this.message);
}
first = false;
sb.append(")");
return sb.toString();
}
public void validate() throws org.apache.thrift.TException {
// check for required fields
// check for sub-struct validity
}
private void writeObject(java.io.ObjectOutputStream out) throws java.io.IOException {
try {
write(new org.apache.thrift.protocol.TCompactProtocol(new org.apache.thrift.transport.TIOStreamTransport(out)));
} catch (org.apache.thrift.TException te) {
throw new java.io.IOException(te);
}
}
private void readObject(java.io.ObjectInputStream in) throws java.io.IOException, ClassNotFoundException {
try {
read(new org.apache.thrift.protocol.TCompactProtocol(new org.apache.thrift.transport.TIOStreamTransport(in)));
} catch (org.apache.thrift.TException te) {
throw new java.io.IOException(te);
}
}
private static class TCommandEchoStandardSchemeFactory implements SchemeFactory {
public TCommandEchoStandardScheme getScheme() {
return new TCommandEchoStandardScheme();
}
}
private static class TCommandEchoStandardScheme extends StandardScheme<TCommandEcho> {
public void read(org.apache.thrift.protocol.TProtocol iprot, TCommandEcho struct) throws org.apache.thrift.TException {
org.apache.thrift.protocol.TField schemeField;
iprot.readStructBegin();
while (true)
{
schemeField = iprot.readFieldBegin();
if (schemeField.type == org.apache.thrift.protocol.TType.STOP) {
break;
}
switch (schemeField.id) {
case 1: // MESSAGE
if (schemeField.type == org.apache.thrift.protocol.TType.STRING) {
struct.message = iprot.readString();
struct.setMessageIsSet(true);
} else {
org.apache.thrift.protocol.TProtocolUtil.skip(iprot, schemeField.type);
}
break;
default:
org.apache.thrift.protocol.TProtocolUtil.skip(iprot, schemeField.type);
}
iprot.readFieldEnd();
}
iprot.readStructEnd();
struct.validate();
}
public void write(org.apache.thrift.protocol.TProtocol oprot, TCommandEcho struct) throws org.apache.thrift.TException {
struct.validate();
oprot.writeStructBegin(STRUCT_DESC);
if (struct.message != null) {
oprot.writeFieldBegin(MESSAGE_FIELD_DESC);
oprot.writeString(struct.message);
oprot.writeFieldEnd();
}
oprot.writeFieldStop();
oprot.writeStructEnd();
}
}
private static class TCommandEchoTupleSchemeFactory implements SchemeFactory {
public TCommandEchoTupleScheme getScheme() {
return new TCommandEchoTupleScheme();
}
}
private static class TCommandEchoTupleScheme extends TupleScheme<TCommandEcho> {
@Override
public void write(org.apache.thrift.protocol.TProtocol prot, TCommandEcho struct) throws org.apache.thrift.TException {
TTupleProtocol oprot = (TTupleProtocol) prot;
BitSet optionals = new BitSet();
if (struct.isSetMessage()) {
optionals.set(0);
}
oprot.writeBitSet(optionals, 1);
if (struct.isSetMessage()) {
oprot.writeString(struct.message);
}
}
@Override
public void read(org.apache.thrift.protocol.TProtocol prot, TCommandEcho struct) throws org.apache.thrift.TException {
TTupleProtocol iprot = (TTupleProtocol) prot;
BitSet incoming = iprot.readBitSet(1);
if (incoming.get(0)) {
struct.message = iprot.readString();
struct.setMessageIsSet(true);
}
}
}
}
@@ -3,7 +3,9 @@ package com.nhn.pinpoint.thrift.io;
import org.apache.thrift.TBase;
import com.nhn.pinpoint.thrift.dto.TResult;
import com.nhn.pinpoint.thrift.dto.command.TCommandEcho;
import com.nhn.pinpoint.thrift.dto.command.TCommandThreadDump;
import com.nhn.pinpoint.thrift.dto.command.TCommandTransfer;
/**
* @author koo.taejin
@@ -19,6 +21,18 @@ public enum TCommandType {
return new TResult();
}
},
TRANSFER((short) 700, TCommandTransfer.class) {
@Override
public TBase newObject() {
return new TCommandTransfer();
}
},
ECHO((short) 710, TCommandEcho.class) {
@Override
public TBase newObject() {
return new TCommandEcho();
}
},
THREAD_DUMP((short) 720, TCommandThreadDump.class) {
@Override
public TBase newObject() {
@@ -3,6 +3,8 @@ package com.nhn.pinpoint.thrift.io;
import java.util.ArrayList;
import java.util.List;
import org.apache.thrift.TBase;
/**
* @author koo.taejin
*/
@@ -10,6 +12,9 @@ public enum TCommandTypeVersion {
// Agent 버젼과 맞추면 좋을듯 일단은 Agent 버전과 맞춰놓음
V_1_0_2_SNAPSHOT("1.0.2-SNAPSHOT", TCommandType.RESULT, TCommandType.THREAD_DUMP),
V_1_0_2("1.0.2", V_1_0_2_SNAPSHOT),
V_1_0_3_SNAPSHOT("1.0.3-SNAPSHOT", V_1_0_2, TCommandType.ECHO, TCommandType.TRANSFER),
UNKNOWN("UNKNOWN");
private final String versionName;
@@ -38,6 +43,20 @@ public enum TCommandTypeVersion {
public List<TCommandType> getSupportCommandList() {
return supportCommandList;
}
public boolean isSupportCommand(TBase command) {
if (command == null) {
return false;
}
for (TCommandType eachCommand : supportCommandList) {
if (command.getClass().isInstance(eachCommand)) {
return true;
}
}
return false;
}
public String getVersionName() {
return versionName;
+20 -11
View File
@@ -1,11 +1,20 @@
namespace java com.nhn.pinpoint.thrift.dto.command
enum TThreadDumpType {
TARGET,
PENDING
}
struct TCommandThreadDump {
1: TThreadDumpType type = TThreadDumpType.TARGET
2: optional string name
3: optional i64 pendingTimeMillis
}
namespace java com.nhn.pinpoint.thrift.dto.command
enum TThreadDumpType {
TARGET,
PENDING
}
struct TCommandThreadDump {
1: TThreadDumpType type = TThreadDumpType.TARGET
2: optional string name
3: optional i64 pendingTimeMillis
}
struct TCommandEcho {
1: string message
}
struct TCommandTransfer {
1: string applicationName
2: string agentId
3: optional i64 startTime
4: binary payload
}