diff --git a/pom.xml b/pom.xml
index 15e607863..5e3a4c682 100644
--- a/pom.xml
+++ b/pom.xml
@@ -158,6 +158,13 @@
provided
+
+ arcus
+ arcus-client
+ 1.6.2.1
+ provided
+
+
diff --git a/src/main/java/com/profiler/context/AsyncTrace.java b/src/main/java/com/profiler/context/AsyncTrace.java
index 1b52d1609..d2965b305 100644
--- a/src/main/java/com/profiler/context/AsyncTrace.java
+++ b/src/main/java/com/profiler/context/AsyncTrace.java
@@ -23,6 +23,25 @@ public class AsyncTrace {
this.span = span;
}
+ public void setDataSender(DataSender dataSender) {
+ this.dataSender = dataSender;
+ }
+
+ public Span getSpan() {
+ return span;
+ }
+
+ private Object attachObject;
+
+ public Object getAttachObject() {
+ return attachObject;
+ }
+
+ public void setAttachObject(Object attachObject) {
+ this.attachObject = attachObject;
+ }
+
+
public void record(Annotation annotation) {
annotate(annotation.getCode(), null);
}
@@ -52,14 +71,8 @@ public class AsyncTrace {
public void recordRpcName(final String service, final String rpc) {
try {
- spanUpdate(new SpanUpdater() {
- @Override
- public Span updateSpan(Span span) {
- span.setServiceName(service);
- span.setName(rpc);
- return span;
- }
- });
+ this.span.setServiceName(service);
+ this.span.setName(rpc);
} catch (Exception e) {
logger.log(Level.SEVERE, e.getMessage(), e);
}
@@ -76,15 +89,8 @@ public class AsyncTrace {
// TODO: final String... endPoint로 받으면 합치는데 비용이 들어가 그냥 한번에 받는게 나을것 같음.
private void recordEndPoint(final String endPoint, final boolean isTerminal) {
try {
- spanUpdate(new SpanUpdater() {
- @Override
- public Span updateSpan(Span span) {
- // set endpoint to both span and annotations
- span.setEndPoint(endPoint);
- span.setTerminal(isTerminal);
- return span;
- }
- });
+ this.span.setEndPoint(endPoint);
+ this.span.setTerminal(isTerminal);
} catch (Exception e) {
logger.log(Level.SEVERE, e.getMessage(), e);
}
@@ -93,21 +99,19 @@ public class AsyncTrace {
private void annotate(final String key, final Long duration) {
try {
- spanUpdate(new SpanUpdater() {
- @Override
- public Span updateSpan(Span span) {
- span.addAnnotation(new HippoAnnotation(System.currentTimeMillis(), key, duration));
- return span;
- }
- });
+ this.span.addAnnotation(new HippoAnnotation(System.currentTimeMillis(), key, duration));
+ logSpan(key, span);
} catch (Exception e) {
logger.log(Level.SEVERE, e.getMessage(), e);
}
}
- private void spanUpdate(SpanUpdater spanUpdater) {
- if (span.isExistsAnnotationKey(Annotation.ClientRecv.getCode()) || span.isExistsAnnotationKey(Annotation.ServerSend.getCode())) {
+ private void logSpan(String key, Span span) {
+ if (key == null) {
+ return;
+ }
+ if (key.equals(Annotation.ClientRecv.getCode()) || key.equals(Annotation.ServerSend.getCode())) {
logSpan(span);
}
}
diff --git a/src/main/java/com/profiler/context/DeadlineSpanMap.java b/src/main/java/com/profiler/context/DeadlineSpanMap.java
index fdc33463d..462985812 100644
--- a/src/main/java/com/profiler/context/DeadlineSpanMap.java
+++ b/src/main/java/com/profiler/context/DeadlineSpanMap.java
@@ -22,7 +22,7 @@ public class DeadlineSpanMap {
map.put(traceIdKey, span);
TimerTask task = new FlushTimedoutSpanTask(span);
- span.setTimerTask(task);
+// span.setTimerTask(task);
timer.schedule(task, FLUSH_TIMEOUT);
}
diff --git a/src/main/java/com/profiler/context/GlobalCallTrace.java b/src/main/java/com/profiler/context/GlobalCallTrace.java
index 3a5711697..247a8d63c 100644
--- a/src/main/java/com/profiler/context/GlobalCallTrace.java
+++ b/src/main/java/com/profiler/context/GlobalCallTrace.java
@@ -10,22 +10,24 @@ import java.util.concurrent.atomic.AtomicInteger;
/**
*
*/
-public class GlobalCallTrace {
+public class GlobalCallTrace {
private static AtomicInteger timerId = new AtomicInteger(0);
- private ConcurrentMap trace = new ConcurrentHashMap(32);
+ private ConcurrentMap trace = new ConcurrentHashMap(32);
private AtomicInteger idGenerator = new AtomicInteger(0);
private DataSender dataSender;
private Timer timer = new Timer("GlobalCallTrace-Timer-" + timerId.getAndIncrement(), true);
- public int registerTraceObject(T target) {
+ public int registerTraceObject(AsyncTrace target) {
+ // datasender쪽 전달부분이 영 별로임.
+ target.setDataSender(this.dataSender);
int id = idGenerator.getAndIncrement();
trace.put(id, target);
return id;
}
- public T removeTraceObject(int id) {
+ public AsyncTrace removeTraceObject(int id) {
return trace.get(id);
}
diff --git a/src/main/java/com/profiler/context/Span.java b/src/main/java/com/profiler/context/Span.java
index 2466c325c..5ff2fa14d 100644
--- a/src/main/java/com/profiler/context/Span.java
+++ b/src/main/java/com/profiler/context/Span.java
@@ -10,157 +10,147 @@ import com.profiler.Agent;
/**
* Span represent RPC
- *
+ *
* @author netspider
- *
*/
public class Span {
- private final TraceID traceID;
- private final long createTime;
+ private final TraceID traceID;
+ private final long createTime;
- private String serviceName;
- private String name;
- private String endPoint;
- private boolean isTerminal = false;
+ private String serviceName;
+ private String name;
+ private String endPoint;
+ private boolean isTerminal = false;
- private final List annotations = new ArrayList(5);
- private final Set annotationKeys = new HashSet(5);
-
- private long rpcStartTime;
- private long rpcEndTime;
-
- /**
- * Cancel timer logic.
- * TODO: refactor this.
- */
- private TimerTask timerTask;
+ private final List annotations = new ArrayList(5);
+ private final Set annotationKeys = new HashSet(5);
- public void setTimerTask(TimerTask task) {
- this.timerTask = task;
- }
+ private long rpcStartTime;
+ private long rpcEndTime;
- public boolean cancelTimer() {
- return timerTask.cancel();
- }
+ public Span(TraceID traceId, String name, String endPoint) {
+ this.traceID = traceId;
+ this.name = name;
+ this.endPoint = endPoint;
+ this.createTime = System.currentTimeMillis();
+ }
- public Span(TraceID traceId, String name, String endPoint) {
- this.traceID = traceId;
- this.name = name;
- this.endPoint = endPoint;
- this.createTime = System.currentTimeMillis();
- }
+ public boolean addAnnotation(HippoAnnotation annotation) {
+ annotationKeys.add(annotation.getKey());
+ if (annotation.getKey().equals(Annotation.ClientSend.getCode()) || annotation.getKey().equals(Annotation.ServerRecv.getCode())) {
+ rpcStartTime = annotation.getTimestamp();
+ }
+ if (annotation.getKey().equals(Annotation.ClientRecv.getCode()) || annotation.getKey().equals(Annotation.ServerSend.getCode())) {
+ rpcEndTime = annotation.getTimestamp();
+ }
+ return annotations.add(annotation);
+ }
- public boolean addAnnotation(HippoAnnotation annotation) {
- annotationKeys.add(annotation.getKey());
- if (annotation.getKey().equals(Annotation.ClientSend.getCode()) || annotation.getKey().equals(Annotation.ServerRecv.getCode())) {
- rpcStartTime = annotation.getTimestamp();
- }
- if (annotation.getKey().equals(Annotation.ClientRecv.getCode()) || annotation.getKey().equals(Annotation.ServerSend.getCode())) {
- rpcEndTime = annotation.getTimestamp();
- }
- return annotations.add(annotation);
- }
+ public int getAnnotationSize() {
+ return annotations.size();
+ }
- public int getAnnotationSize() {
- return annotations.size();
- }
+ /**
+ * this method only works for Trace.mutate()
+ *
+ * @param value
+ * @return
+ */
+ public boolean isExistsAnnotationKey(String key) {
+ return annotationKeys.contains(key);
+ }
- /**
- * this method only works for Trace.mutate()
- *
- * @param value
- * @return
- */
- public boolean isExistsAnnotationKey(String key) {
- return annotationKeys.contains(key);
- }
+ public String getEndPoint() {
+ return this.endPoint;
+ }
- public String getEndPoint() {
- return this.endPoint;
- }
+ public String getServiceName() {
+ return serviceName;
+ }
- public String getServiceName() {
- return serviceName;
- }
+ public void setServiceName(String serviceName) {
+ this.serviceName = serviceName;
+ }
- public void setServiceName(String serviceName) {
- this.serviceName = serviceName;
- }
+ public String getName() {
+ return name;
+ }
- public String getName() {
- return name;
- }
+ public void setName(String name) {
+ this.name = name;
+ }
- public void setName(String name) {
- this.name = name;
- }
+ public void setEndPoint(String endPoint) {
+ this.endPoint = endPoint;
+ }
- public void setEndPoint(String endPoint) {
- this.endPoint = endPoint;
- }
-
- public boolean isTerminal() {
- return isTerminal;
- }
+ public boolean isTerminal() {
+ return isTerminal;
+ }
- public void setTerminal(boolean isTerminal) {
- this.isTerminal = isTerminal;
- }
-
- public long getRpcStartTime() {
- return rpcStartTime;
- }
+ public void setTerminal(boolean isTerminal) {
+ this.isTerminal = isTerminal;
+ }
- public long getRpcEndTime() {
- return rpcEndTime;
- }
+ public long getRpcStartTime() {
+ return rpcStartTime;
+ }
- public String toString() {
- StringBuilder sb = new StringBuilder();
+ public long getRpcEndTime() {
+ return rpcEndTime;
+ }
- sb.append("{");
- sb.append("\n\t TraceID = ").append(traceID);
- sb.append(",\n\t CreateTime = ").append(createTime);
- sb.append(",\n\t Name = ").append(name);
- sb.append(",\n\t ServiceName = ").append(serviceName);
- sb.append(",\n\t EndPoint = ").append(endPoint);
+ public long getCreateTime() {
+ return createTime;
+ }
- sb.append(",\n\t Annotations = {");
- for (HippoAnnotation a : annotations) {
- sb.append("\n\t\t").append(a);
- }
- sb.append("\n\t}");
- sb.append("}");
+ public String toString() {
+ StringBuilder sb = new StringBuilder();
- return sb.toString();
- }
-
- public com.profiler.common.dto.thrift.Span toThrift() {
- com.profiler.common.dto.thrift.Span span = new com.profiler.common.dto.thrift.Span();
+ sb.append("{");
+ sb.append("\n\t TraceID = ").append(traceID);
+ sb.append(",\n\t CreateTime = ").append(createTime);
+ sb.append(",\n\t Name = ").append(name);
+ sb.append(",\n\t ServiceName = ").append(serviceName);
+ sb.append(",\n\t EndPoint = ").append(endPoint);
- span.setAgentId(Agent.getInstance().getAgentId());
- span.setTimestamp(createTime);
- span.setMostTraceId(traceID.getId().getMostSignificantBits());
- span.setLeastTraceId(traceID.getId().getLeastSignificantBits());
- span.setName(name);
- span.setServiceName(serviceName);
- span.setSpanId(traceID.getSpanId());
- span.setParentSpanId(traceID.getParentSpanId());
- span.setEndPoint(endPoint);
- span.setTerminal(isTerminal);
-
- // TODO: set duration.
-
- List annotationList = new ArrayList(annotations.size());
- for (HippoAnnotation a : annotations) {
- annotationList.add(a.toThrift());
- }
- span.setAnnotations(annotationList);
+ sb.append(",\n\t Annotations = {");
+ for (HippoAnnotation a : annotations) {
+ sb.append("\n\t\t").append(a);
+ }
+ sb.append("\n\t}");
- span.setFlag(traceID.getFlags());
+ sb.append("}");
- return span;
- }
+ return sb.toString();
+ }
+
+ public com.profiler.common.dto.thrift.Span toThrift() {
+ com.profiler.common.dto.thrift.Span span = new com.profiler.common.dto.thrift.Span();
+
+ span.setAgentId(Agent.getInstance().getAgentId());
+ span.setTimestamp(createTime);
+ span.setMostTraceId(traceID.getId().getMostSignificantBits());
+ span.setLeastTraceId(traceID.getId().getLeastSignificantBits());
+ span.setName(name);
+ span.setServiceName(serviceName);
+ span.setSpanId(traceID.getSpanId());
+ span.setParentSpanId(traceID.getParentSpanId());
+ span.setEndPoint(endPoint);
+ span.setTerminal(isTerminal);
+
+ // TODO: set duration.
+
+ List annotationList = new ArrayList(annotations.size());
+ for (HippoAnnotation a : annotations) {
+ annotationList.add(a.toThrift());
+ }
+ span.setAnnotations(annotationList);
+
+ span.setFlag(traceID.getFlags());
+
+ return span;
+ }
}
diff --git a/src/main/java/com/profiler/context/Trace.java b/src/main/java/com/profiler/context/Trace.java
index 1bec1aeb6..be4c29ce6 100644
--- a/src/main/java/com/profiler/context/Trace.java
+++ b/src/main/java/com/profiler/context/Trace.java
@@ -150,14 +150,12 @@ public final class Trace {
}
- private void spanUpdate(SpanUpdater spanUpdater) {
- StackFrame currentStackFrame = getCurrentStackFrame();
- Span span = spanUpdater.updateSpan(currentStackFrame.getSpan());
- if (span.isExistsAnnotationKey(Annotation.ClientRecv.getCode()) || span.isExistsAnnotationKey(Annotation.ServerSend.getCode())) {
- // remove current context threadId from callStack
-// removeCurrentTraceIdFromStack();
+ private void logSpan(String key, Span span) {
+ if (key == null) {
+ return;
+ }
+ if (key.equals(Annotation.ClientRecv.getCode()) || key.equals(Annotation.ServerSend.getCode())) {
logSpan(span);
-
}
}
@@ -197,24 +195,17 @@ public final class Trace {
}
public void recordAttribute(final String key, final String value) {
- recordAttibute(key, (Object) value);
+ recordAttribute(key, (Object) value);
}
- public void recordAttibute(final String key, final Object value) {
+ public void recordAttribute(final String key, final Object value) {
if (!tracingEnabled)
return;
try {
-
- spanUpdate(new SpanUpdater() {
- @Override
- public Span updateSpan(Span span) {
- // TODO 사용자 thread에서 encoding을 하지 않도록 변경.
- Encoded enc = transcoder.encode(value);
- span.addAnnotation(new HippoAnnotation(System.currentTimeMillis(), key, enc.getValueType(), enc.getBytes(), null));
- return span;
- }
- });
+ Span span = getCurrentStackFrame().getSpan();
+ Encoded enc = transcoder.encode(value);
+ span.addAnnotation(new HippoAnnotation(System.currentTimeMillis(), key, enc.getValueType(), enc.getBytes(), null));
} catch (Exception e) {
logger.log(Level.SEVERE, e.getMessage(), e);
}
@@ -232,14 +223,9 @@ public final class Trace {
return;
try {
- spanUpdate(new SpanUpdater() {
- @Override
- public Span updateSpan(Span span) {
- span.setServiceName(service);
- span.setName(rpc);
- return span;
- }
- });
+ Span span = getCurrentStackFrame().getSpan();
+ span.setServiceName(service);
+ span.setName(rpc);
} catch (Exception e) {
logger.log(Level.SEVERE, e.getMessage(), e);
}
@@ -259,15 +245,9 @@ public final class Trace {
return;
try {
- spanUpdate(new SpanUpdater() {
- @Override
- public Span updateSpan(Span span) {
- // set endpoint to both span and annotations
- span.setEndPoint(endPoint);
- span.setTerminal(isTerminal);
- return span;
- }
- });
+ Span span = getCurrentStackFrame().getSpan();
+ span.setEndPoint(endPoint);
+ span.setTerminal(isTerminal);
} catch (Exception e) {
logger.log(Level.SEVERE, e.getMessage(), e);
}
@@ -278,13 +258,9 @@ public final class Trace {
return;
try {
- spanUpdate(new SpanUpdater() {
- @Override
- public Span updateSpan(Span span) {
- span.addAnnotation(new HippoAnnotation(System.currentTimeMillis(), key, duration));
- return span;
- }
- });
+ Span span = getCurrentStackFrame().getSpan();
+ span.addAnnotation(new HippoAnnotation(System.currentTimeMillis(), key, duration));
+ logSpan(key, span);
} catch (Exception e) {
logger.log(Level.SEVERE, e.getMessage(), e);
}
diff --git a/src/main/java/com/profiler/modifier/arcus/ArcusClientModifier.java b/src/main/java/com/profiler/modifier/arcus/ArcusClientModifier.java
index c6a440957..08d58491e 100644
--- a/src/main/java/com/profiler/modifier/arcus/ArcusClientModifier.java
+++ b/src/main/java/com/profiler/modifier/arcus/ArcusClientModifier.java
@@ -4,6 +4,7 @@ import java.security.ProtectionDomain;
import java.util.logging.Level;
import java.util.logging.Logger;
+import com.profiler.interceptor.Interceptor;
import com.profiler.interceptor.bci.ByteCodeInstrumentor;
import com.profiler.interceptor.bci.InstrumentClass;
import com.profiler.modifier.AbstractModifier;
@@ -35,29 +36,10 @@ public class ArcusClientModifier extends AbstractModifier {
aClass.addTraceVariable("__asyncTraceId", "__setAsyncTraceId", "__getAsyncTraceId", "int");
aClass.addConstructorInterceptor(null, new ConstructInterceptor());
-// /**
-// * inject both current and next traceId.
-// */
-// aClass.addTraceVariable("__traceId", "__setTraceId", "__getTraceId", "com.profiler.context.TraceID");
-// aClass.addTraceVariable("__nextTraceId", "__setNextTraceId", "__getNextTraceId", "com.profiler.context.TraceID");
-// aClass.insertCodeAfterConstructor(null, "{ __setTraceId(com.profiler.context.Trace.getCurrentTraceId()); __setNextTraceId(com.profiler.context.Trace.getNextTraceId()); }");
-//
-// /**
-// * inject nano time for checking send time.
-// */
-// aClass.addTraceVariable("__commandCreatedTime", "__setCommandCreatedTime", "__getCommandCreatedTime", "long");
-// aClass.insertCodeAfterConstructor(null, "{ __setCommandCreatedTime(System.nanoTime()); }");
-//
-// /**
-// * inject cancelled time
-// */
-// aClass.addTraceVariable("__cancelledTime", "__setCancelledTime", "__getCancelledTime", "long");
-//
-// /**
-// * insert trace code.
-// */
-// aClass.insertCodeBeforeMethod("transitionState", new String[]{"net.spy.memcached.ops.OperationState"}, getTransitionStateAfterCode());
-// aClass.insertCodeBeforeMethod("cancel", null, getCancelBeforeCode());
+ Interceptor transitionStateInterceptor = newInterceptor(classLoader, protectedDomain, "com.profiler.modifier.arcus.interceptors.BaseOperationTransitionStateInterceptor");
+ aClass.addInterceptor("transitionState", new String[]{"net.spy.memcached.ops.OperationState"}, transitionStateInterceptor);
+ Interceptor cancelInterceptor = newInterceptor(classLoader, protectedDomain, "com.profiler.modifier.arcus.interceptors.BaseOperationCancelInterceptor");
+ aClass.addInterceptor("cancel", null, cancelInterceptor);
return aClass.toBytecode();
} catch (Exception e) {
diff --git a/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationCancelInterceptor.java b/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationCancelInterceptor.java
new file mode 100644
index 000000000..db79b55ba
--- /dev/null
+++ b/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationCancelInterceptor.java
@@ -0,0 +1,43 @@
+package com.profiler.modifier.arcus.interceptors;
+
+import com.profiler.context.AsyncTrace;
+import com.profiler.context.GlobalCallTrace;
+import com.profiler.context.TraceContext;
+import com.profiler.interceptor.StaticAfterInterceptor;
+import com.profiler.util.MetaObject;
+import com.profiler.util.StringUtils;
+import net.spy.memcached.protocol.BaseOperationImpl;
+
+import java.util.Arrays;
+import java.util.logging.Level;
+import java.util.logging.Logger;
+
+/**
+ *
+ */
+public class BaseOperationCancelInterceptor implements StaticAfterInterceptor {
+ private final Logger logger = Logger.getLogger(BaseOperationCancelInterceptor.class.getName());
+ private MetaObject asyncTraceId = new MetaObject("__getAsyncTraceId", null);
+
+ @Override
+ public void after(Object target, String className, String methodName, String parameterDescription, Object[] args, Object result) {
+ if (logger.isLoggable(Level.INFO)) {
+ logger.info("after " + StringUtils.toString(target) + " " + className + "." + methodName + parameterDescription + " args:" + Arrays.toString(args) + " result:" + result);
+ }
+
+ TraceContext traceContext = TraceContext.getTraceContext();
+ GlobalCallTrace globalCallTrace = traceContext.getGlobalCallTrace();
+ Object asyncId = asyncTraceId.invoke(target);
+ if (asyncId == null) {
+ logger.info("asyncId not found");
+ return;
+ }
+ AsyncTrace asyncTrace = globalCallTrace.removeTraceObject((Integer) asyncId);
+ BaseOperationImpl baseOperation = (BaseOperationImpl) target;
+ if (!baseOperation.isCancelled()) {
+ TimeObject timeObject = (TimeObject) asyncTrace.getAttachObject();
+ timeObject.markCancelTime();
+ }
+ }
+}
+
diff --git a/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationImplInterceptor.java b/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationImplInterceptor.java
deleted file mode 100644
index ce9081b8f..000000000
--- a/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationImplInterceptor.java
+++ /dev/null
@@ -1,59 +0,0 @@
-package com.profiler.modifier.arcus.interceptors;
-
-import com.profiler.context.AsyncTrace;
-import com.profiler.context.GlobalCallTrace;
-import com.profiler.context.Trace;
-import com.profiler.context.TraceContext;
-import com.profiler.interceptor.StaticAfterInterceptor;
-
-/**
- *
- */
-public class BaseOperationImplInterceptor implements StaticAfterInterceptor {
- @Override
- public void after(Object target, String className, String methodName, String parameterDescription, Object[] args, Object result) {
-
- TraceContext traceContext = TraceContext.getTraceContext();
- GlobalCallTrace globalCallTrace = traceContext.getGlobalCallTrace();
- AsyncTrace asyncTrace = globalCallTrace.removeTraceObject(1);
-
-// com.profiler.context.Trace.setTraceId(__nextTraceId);
-//
-//
-// /**
-// * After sending command to the Arcus server. now waiting server
-// * response.
-// */
-// code.append("if (newState == net.spy.memcached.ops.OperationState.READING) {");
-//
-// code.append(" java.net.SocketAddress socketAddress = handlingNode.getSocketAddress();");
-// code.append(" if (socketAddress instanceof java.net.InetSocketAddress) {");
-// code.append(" java.net.InetSocketAddress addr = (java.net.InetSocketAddress) handlingNode.getSocketAddress();");
-// code.append(" com.profiler.context.Trace.recordTerminalEndPoint(\"ARCUS:\" + addr.getHostName() + \":\" + addr.getPort());");
-// code.append(" }");
-// code.append(" com.profiler.context.Trace.recordRpcName(\"ARCUS\", this.getClass().getSimpleName());");
-// code.append(" com.profiler.context.Trace.recordAttribute(\"arcus.command\", ((cmd == null) ? \"UNKNOWN\" : new String(cmd.array())));");
-// code.append(" com.profiler.StopWatch.start(this.hashCode());");
-// code.append(" com.profiler.context.Trace.record(com.profiler.context.Annotation.ClientSend, System.nanoTime() - __commandCreatedTime);");
-
- /**
- * Received all response or timed out.
- */
-// code.append("} else if (newState == net.spy.memcached.ops.OperationState.COMPLETE || newState == net.spy.memcached.ops.OperationState.TIMEDOUT) {");
-// code.append(" if (exception != null) { ");
-// code.append(" com.profiler.context.Trace.recordAttribute(\"exception\", com.profiler.util.InterceptorUtils.exceptionToString(exception));");
-// code.append(" }");
-//
-// code.append(" if (!cancelled) {");
-// code.append(" com.profiler.context.Trace.record(com.profiler.context.Annotation.ClientRecv, com.profiler.StopWatch.stopAndGetElapsed(this.hashCode()));");
-// code.append(" } else {");
-// code.append(" com.profiler.context.Trace.recordAttribute(\"exception\", \"cancelled by user\");");
-// code.append(" com.profiler.context.Trace.record(com.profiler.context.Annotation.ClientRecv, System.nanoTime() - __cancelledTime);");
-// code.append(" }");
-//
-// code.append("}");
-//
-//// code.append("com.profiler.context.Trace.traceBlockEnd();");
-// code.append("}");
- }
-}
diff --git a/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationTransitionStateInterceptor.java b/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationTransitionStateInterceptor.java
new file mode 100644
index 000000000..3a9becdb1
--- /dev/null
+++ b/src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationTransitionStateInterceptor.java
@@ -0,0 +1,94 @@
+package com.profiler.modifier.arcus.interceptors;
+
+import com.profiler.context.Annotation;
+import com.profiler.context.AsyncTrace;
+import com.profiler.context.GlobalCallTrace;
+import com.profiler.context.TraceContext;
+import com.profiler.interceptor.StaticAfterInterceptor;
+import com.profiler.util.InterceptorUtils;
+import com.profiler.util.MetaObject;
+import com.profiler.util.StringUtils;
+import net.spy.memcached.MemcachedNode;
+import net.spy.memcached.ops.OperationState;
+import net.spy.memcached.protocol.BaseOperationImpl;
+
+import java.net.InetSocketAddress;
+import java.net.SocketAddress;
+import java.nio.ByteBuffer;
+import java.util.Arrays;
+import java.util.logging.Level;
+import java.util.logging.Logger;
+
+/**
+ *
+ */
+public class BaseOperationTransitionStateInterceptor implements StaticAfterInterceptor {
+
+ private final Logger logger = Logger.getLogger(BaseOperationTransitionStateInterceptor.class.getName());
+ private MetaObject asyncTraceId = new MetaObject("__getAsyncTraceId");
+
+ @Override
+ public void after(Object target, String className, String methodName, String parameterDescription, Object[] args, Object result) {
+ if (logger.isLoggable(Level.INFO)) {
+ logger.info("after " + StringUtils.toString(target) + " " + className + "." + methodName + parameterDescription + " args:" + Arrays.toString(args) + " result:" + result);
+ }
+ TraceContext traceContext = TraceContext.getTraceContext();
+ GlobalCallTrace globalCallTrace = traceContext.getGlobalCallTrace();
+
+ Object asyncId = asyncTraceId.invoke(target);
+ if (asyncId == null) {
+ logger.info("asyncId not found");
+ return;
+ }
+ // saynctrace를 제거한다.
+ AsyncTrace asyncTrace = globalCallTrace.removeTraceObject((Integer) asyncId);
+ OperationState newState = (OperationState) args[0];
+
+ BaseOperationImpl baseOperation = (BaseOperationImpl) target;
+ if (newState == OperationState.READING) {
+ if (logger.isLoggable(Level.FINE)) {
+ logger.log(Level.FINE, "event:" + newState);
+ }
+ MemcachedNode handlingNode = baseOperation.getHandlingNode();
+ SocketAddress socketAddress = handlingNode.getSocketAddress();
+ if (socketAddress instanceof InetSocketAddress) {
+ InetSocketAddress address = (InetSocketAddress) socketAddress;
+ asyncTrace.recordTerminalEndPoint("ARCUS:" + address.getHostName() + ":" + address.getPort());
+ }
+ asyncTrace.recordRpcName("ARCUS", baseOperation.getClass().getSimpleName());
+ String cmd = getCommand(baseOperation);
+ asyncTrace.recordAttibute("arcus.command", cmd);
+
+ TimeObject timeObject = (TimeObject) asyncTrace.getAttachObject();
+ timeObject.markSendTime();
+
+ long createTime = asyncTrace.getSpan().getCreateTime();
+ asyncTrace.record(Annotation.ClientSend, System.currentTimeMillis() - createTime);
+ } else if (newState == OperationState.COMPLETE || newState == OperationState.TIMEDOUT) {
+ if (logger.isLoggable(Level.FINE)) {
+ logger.log(Level.FINE, "event:" + newState);
+ }
+ Exception exception = baseOperation.getException();
+ if (exception != null) {
+ asyncTrace.recordAttibute("exception", InterceptorUtils.exceptionToString(exception));
+ }
+ if (!baseOperation.isCancelled()) {
+ TimeObject timeObject = (TimeObject) asyncTrace.getAttachObject();
+ asyncTrace.record(Annotation.ClientRecv, timeObject.getSendTime());
+ } else {
+ asyncTrace.recordAttribute("exception", "cancelled by user");
+ TimeObject timeObject = (TimeObject) asyncTrace.getAttachObject();
+ asyncTrace.record(Annotation.ClientRecv, timeObject.getCancelTime());
+ }
+ }
+ }
+
+ private String getCommand(BaseOperationImpl baseOperation) {
+ ByteBuffer buffer = baseOperation.getBuffer();
+ if (buffer == null) {
+ return "UNKNOWN";
+ }
+ // TODO 기본 인코딩은 뭔가? 동시성은 괜찮은건가? buffer 사이즈의 compact는 되어있는것인가.
+ return new String(buffer.array());
+ }
+}
diff --git a/src/main/java/com/profiler/modifier/arcus/interceptors/ConstructInterceptor.java b/src/main/java/com/profiler/modifier/arcus/interceptors/ConstructInterceptor.java
index 6b9a5f9e7..8b78d394c 100644
--- a/src/main/java/com/profiler/modifier/arcus/interceptors/ConstructInterceptor.java
+++ b/src/main/java/com/profiler/modifier/arcus/interceptors/ConstructInterceptor.java
@@ -27,11 +27,12 @@ public class ConstructInterceptor implements StaticAfterInterceptor {
if (trace == null) {
return;
}
- GlobalCallTrace globalCallTrace = traceContext.getGlobalCallTrace();
+ GlobalCallTrace globalCallTrace = traceContext.getGlobalCallTrace();
TraceID nextTraceId = trace.getNextTraceId();
Span span = new Span(nextTraceId, null, null);
AsyncTrace asyncTrace = new AsyncTrace(span);
+ asyncTrace.setAttachObject(new TimeObject());
int asyncId = globalCallTrace.registerTraceObject(asyncTrace);
diff --git a/src/main/java/com/profiler/modifier/arcus/interceptors/TimeObject.java b/src/main/java/com/profiler/modifier/arcus/interceptors/TimeObject.java
new file mode 100644
index 000000000..2c5f06cc1
--- /dev/null
+++ b/src/main/java/com/profiler/modifier/arcus/interceptors/TimeObject.java
@@ -0,0 +1,25 @@
+package com.profiler.modifier.arcus.interceptors;
+
+/**
+ *
+ */
+public class TimeObject {
+ private long cancelTime;
+ private long sendTime;
+
+ public void markCancelTime() {
+ cancelTime = System.currentTimeMillis();
+ }
+
+ public long getCancelTime() {
+ return cancelTime;
+ }
+
+ public void markSendTime() {
+ this.sendTime = System.currentTimeMillis();
+ }
+
+ public long getSendTime() {
+ return System.currentTimeMillis() - this.sendTime;
+ }
+}
diff --git a/src/main/java/com/profiler/modifier/connector/interceptors/ExecuteMethodInterceptor.java b/src/main/java/com/profiler/modifier/connector/interceptors/ExecuteMethodInterceptor.java
index 9f271628c..6362fb08e 100644
--- a/src/main/java/com/profiler/modifier/connector/interceptors/ExecuteMethodInterceptor.java
+++ b/src/main/java/com/profiler/modifier/connector/interceptors/ExecuteMethodInterceptor.java
@@ -53,7 +53,7 @@ public class ExecuteMethodInterceptor implements StaticAroundInterceptor {
trace.recordRpcName(request.getProtocolVersion().toString(), "CLIENT");
trace.recordEndPoint(request.getProtocolVersion().toString() + ":" + host.getHostName() + ":" + host.getPort());
- trace.recordAttibute("http.url", request.getRequestLine().getUri());
+ trace.recordAttribute("http.url", request.getRequestLine().getUri());
trace.record(Annotation.ClientSend);
}
diff --git a/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java
index 3650b6028..6d3765681 100644
--- a/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java
+++ b/src/main/java/com/profiler/modifier/db/interceptor/PreparedStatementExecuteQueryInterceptor.java
@@ -45,11 +45,11 @@ public class PreparedStatementExecuteQueryInterceptor implements StaticAroundInt
trace.recordRpcName("MYSQL", url);
trace.recordTerminalEndPoint(url);
String sql = getSql.invoke(target);
- trace.recordAttibute("PreparedStatement", sql);
+ trace.recordAttribute("PreparedStatement", sql);
Map bindValue = getBindValue.invoke(target);
String bindString = toBindVariable(bindValue);
- trace.recordAttibute("BindValue", bindString);
+ trace.recordAttribute("BindValue", bindString);
clean(target);
@@ -96,10 +96,10 @@ public class PreparedStatementExecuteQueryInterceptor implements StaticAroundInt
try {
// TODO 일단 테스트로 실패일경우 종료 아닐경우 resultset fetch까지 계산. fetch count는 옵션으로 빼는게 좋을듯.
boolean success = InterceptorUtils.isSuccess(result);
- trace.recordAttibute("Success", success);
+ trace.recordAttribute("Success", success);
if (!success) {
Throwable th = (Throwable) result;
- trace.recordAttibute("Exception", th.getMessage());
+ trace.recordAttribute("Exception", th.getMessage());
}
trace.record(Annotation.ClientRecv);
} catch (Exception e) {
diff --git a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java
index 25270b28b..9f405fa24 100644
--- a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java
+++ b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteQueryInterceptor.java
@@ -1,6 +1,5 @@
package com.profiler.modifier.db.interceptor;
-import com.profiler.StopWatch;
import com.profiler.context.Annotation;
import com.profiler.context.Trace;
import com.profiler.context.TraceContext;
@@ -47,7 +46,7 @@ public class StatementExecuteQueryInterceptor implements StaticAroundInterceptor
trace.recordRpcName("MYSQL", url);
trace.recordTerminalEndPoint(url);
if (args.length > 0) {
- trace.recordAttibute("Statement", args[0]);
+ trace.recordAttribute("Statement", args[0]);
}
trace.record(Annotation.ClientSend);
} catch (Exception e) {
@@ -72,7 +71,7 @@ public class StatementExecuteQueryInterceptor implements StaticAroundInterceptor
return;
}
- trace.recordAttibute("Success", InterceptorUtils.isSuccess(result));
+ trace.recordAttribute("Success", InterceptorUtils.isSuccess(result));
trace.record(Annotation.ClientRecv, trace.afterTime());
trace.traceBlockEnd();
}
diff --git a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java
index e1f44f0df..821a1c9df 100644
--- a/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java
+++ b/src/main/java/com/profiler/modifier/db/interceptor/StatementExecuteUpdateInterceptor.java
@@ -1,6 +1,5 @@
package com.profiler.modifier.db.interceptor;
-import com.profiler.StopWatch;
import com.profiler.context.Annotation;
import com.profiler.context.Trace;
import com.profiler.context.TraceContext;
@@ -44,7 +43,7 @@ public class StatementExecuteUpdateInterceptor implements StaticAroundIntercepto
if (args.length > 0) {
String url = (String) this.getUrl.invoke(target);
trace.recordRpcName("MYSQL", url);
- trace.recordAttibute("Query", url);
+ trace.recordAttribute("Query", url);
trace.recordTerminalEndPoint(url);
} else {
trace.recordRpcName("MYSQL", "UNKNOWN");
diff --git a/src/main/java/com/profiler/modifier/db/interceptor/TransactionInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/TransactionInterceptor.java
index 9a894bd26..1e0b9b9e2 100644
--- a/src/main/java/com/profiler/modifier/db/interceptor/TransactionInterceptor.java
+++ b/src/main/java/com/profiler/modifier/db/interceptor/TransactionInterceptor.java
@@ -86,20 +86,20 @@ public class TransactionInterceptor implements StaticAroundInterceptor {
if (!autocommit) {
// transaction start;
if (success) {
- trace.recordAttibute("Transaction", "begin");
+ trace.recordAttribute("Transaction", "begin");
} else {
- trace.recordAttibute("Transaction", "begin fail");
+ trace.recordAttribute("Transaction", "begin fail");
Throwable th = (Throwable) result;
- trace.recordAttibute("Exception", th.getMessage());
+ trace.recordAttribute("Exception", th.getMessage());
}
trace.record(Annotation.ClientRecv);
} else {
if (success) {
- trace.recordAttibute("Transaction", "autoCommit:false");
+ trace.recordAttribute("Transaction", "autoCommit:false");
} else {
- trace.recordAttibute("Transaction", "autoCommit:false fail");
+ trace.recordAttribute("Transaction", "autoCommit:false fail");
Throwable th = (Throwable) result;
- trace.recordAttibute("Exception", th.getMessage());
+ trace.recordAttribute("Exception", th.getMessage());
}
trace.record(Annotation.ClientRecv);
}
@@ -129,11 +129,11 @@ public class TransactionInterceptor implements StaticAroundInterceptor {
boolean success = InterceptorUtils.isSuccess(result);
if (success) {
- trace.recordAttibute("Transaction", "commit");
+ trace.recordAttribute("Transaction", "commit");
} else {
- trace.recordAttibute("Transaction", "commit fail");
+ trace.recordAttribute("Transaction", "commit fail");
Throwable th = (Throwable) result;
- trace.recordAttibute("Exception", th.getMessage());
+ trace.recordAttribute("Exception", th.getMessage());
}
trace.record(Annotation.ClientRecv);
} catch (Exception e) {
@@ -163,11 +163,11 @@ public class TransactionInterceptor implements StaticAroundInterceptor {
boolean success = InterceptorUtils.isSuccess(result);
if (success) {
- trace.recordAttibute("Transaction", "rollback");
+ trace.recordAttribute("Transaction", "rollback");
} else {
- trace.recordAttibute("Transaction", "rollback fail");
+ trace.recordAttribute("Transaction", "rollback fail");
Throwable th = (Throwable) result;
- trace.recordAttibute("Exception", th.getMessage());
+ trace.recordAttribute("Exception", th.getMessage());
}
trace.record(Annotation.ClientRecv);
} catch (Exception e) {
diff --git a/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardHostValveInvokeInterceptor.java b/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardHostValveInvokeInterceptor.java
index 9904c87b3..8f26fab93 100644
--- a/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardHostValveInvokeInterceptor.java
+++ b/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardHostValveInvokeInterceptor.java
@@ -54,9 +54,9 @@ public class StandardHostValveInvokeInterceptor implements StaticAroundIntercept
trace.markBeforeTime();
trace.recordRpcName(Agent.getInstance().getApplicationName(), requestURL);
trace.recordEndPoint(request.getProtocol() + ":" + request.getLocalName() + ":" + request.getLocalPort());
- trace.recordAttibute("http.url", request.getRequestURI());
+ trace.recordAttribute("http.url", request.getRequestURI());
if (parameters != null && parameters.length() > 0) {
- trace.recordAttibute("http.params", parameters);
+ trace.recordAttribute("http.params", parameters);
}
trace.record(Annotation.ServerRecv);
diff --git a/src/test/java/com/profiler/context/TraceTest.java b/src/test/java/com/profiler/context/TraceTest.java
index cc1c8c854..e23ec1696 100644
--- a/src/test/java/com/profiler/context/TraceTest.java
+++ b/src/test/java/com/profiler/context/TraceTest.java
@@ -13,7 +13,7 @@ public class TraceTest {
// http server receive
trace.recordRpcName("service_name", "http://");
trace.recordEndPoint("http:localhost:8080");
- trace.recordAttibute("KEY", "VALUE");
+ trace.recordAttribute("KEY", "VALUE");
trace.record(Annotation.ServerRecv);
// get data form db
@@ -31,7 +31,7 @@ public class TraceTest {
// db server request
trace.recordRpcName("mysql", "rpc");
- trace.recordAttibute("mysql.query", "SELECT * FROM TABLE");
+ trace.recordAttribute("mysql.query", "SELECT * FROM TABLE");
// get a db response