From b2c37b5242bfad19d32173fca755c11826615918 Mon Sep 17 00:00:00 2001 From: Woonduk Kang Date: Tue, 13 Nov 2012 05:00:44 +0000 Subject: [PATCH] =?UTF-8?q?[=EA=B0=95=EC=9A=B4=EB=8D=95]=20[LUCYSUS-1744]?= =?UTF-8?q?=20arcus=20asynctrace=EB=B6=80=EB=B6=84=20=EA=B0=9C=EB=B0=9C?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit git-svn-id: http://svn.bds.nhncorp.com/pe/hippo-tomcat-profiler/trunk@856 84d0f5b1-2673-498c-a247-62c4ff18d310 --- pom.xml | 7 + .../java/com/profiler/context/AsyncTrace.java | 56 ++-- .../com/profiler/context/DeadlineSpanMap.java | 2 +- .../com/profiler/context/GlobalCallTrace.java | 10 +- src/main/java/com/profiler/context/Span.java | 240 +++++++++--------- src/main/java/com/profiler/context/Trace.java | 62 ++--- .../modifier/arcus/ArcusClientModifier.java | 28 +- .../BaseOperationCancelInterceptor.java | 43 ++++ .../BaseOperationImplInterceptor.java | 59 ----- ...seOperationTransitionStateInterceptor.java | 94 +++++++ .../interceptors/ConstructInterceptor.java | 3 +- .../arcus/interceptors/TimeObject.java | 25 ++ .../ExecuteMethodInterceptor.java | 2 +- ...paredStatementExecuteQueryInterceptor.java | 8 +- .../StatementExecuteQueryInterceptor.java | 5 +- .../StatementExecuteUpdateInterceptor.java | 3 +- .../interceptor/TransactionInterceptor.java | 24 +- .../StandardHostValveInvokeInterceptor.java | 4 +- .../java/com/profiler/context/TraceTest.java | 4 +- 19 files changed, 371 insertions(+), 308 deletions(-) create mode 100644 src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationCancelInterceptor.java delete mode 100644 src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationImplInterceptor.java create mode 100644 src/main/java/com/profiler/modifier/arcus/interceptors/BaseOperationTransitionStateInterceptor.java create mode 100644 src/main/java/com/profiler/modifier/arcus/interceptors/TimeObject.java 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