diff --git a/src/main/java/com/profiler/Agent.java b/src/main/java/com/profiler/Agent.java index d9296279c..d976a3ca7 100644 --- a/src/main/java/com/profiler/Agent.java +++ b/src/main/java/com/profiler/Agent.java @@ -9,6 +9,7 @@ import com.profiler.common.dto.thrift.AgentInfo; import com.profiler.common.util.SpanUtils; import com.profiler.context.TraceContext; import com.profiler.sender.DataSender; +import com.profiler.sender.UdpDataSender; import com.profiler.util.NetworkUtils; public class Agent { @@ -21,6 +22,7 @@ public class Agent { private final ServerInfo serverInfo; private final SystemMonitor systemMonitor; + private DataSender dataSender; private final String agentId; private final String nodeName; @@ -111,13 +113,15 @@ public class Agent { agentInfo.setIsAlive(true); agentInfo.setTimestamp(System.currentTimeMillis()); - DataSender.getInstance().addDataToSend(agentInfo); + this.dataSender.send(agentInfo); } public void start() { logger.info("Starting HIPPO Agent."); // trace context 새롭게 생성. TraceContext.initialize(); + this.dataSender = UdpDataSender.getInstance(); + systemMonitor.setDataSender(dataSender); systemMonitor.start(); } @@ -139,14 +143,16 @@ public class Agent { agentInfo.setTimestamp(System.currentTimeMillis()); agentInfo.setAgentId(getAgentId()); - DataSender.getInstance().addDataToSend(agentInfo); + this.dataSender.send(agentInfo); + // 종료 처리 필요. + this.dataSender.stop(); } public static void startAgent() { Agent.getInstance().start(); } - public static void stopAgent() throws Exception { + public static void stopAgent() { Agent.getInstance().stop(); } } diff --git a/src/main/java/com/profiler/LifeCycleEventHandler.java b/src/main/java/com/profiler/LifeCycleEventHandler.java deleted file mode 100644 index 9fbc0b476..000000000 --- a/src/main/java/com/profiler/LifeCycleEventHandler.java +++ /dev/null @@ -1,6 +0,0 @@ -package com.profiler; - -public interface LifeCycleEventHandler { - void start(); - void stop(); -} diff --git a/src/main/java/com/profiler/LifeCycleEventListener.java b/src/main/java/com/profiler/LifeCycleEventListener.java new file mode 100644 index 000000000..709de20f2 --- /dev/null +++ b/src/main/java/com/profiler/LifeCycleEventListener.java @@ -0,0 +1,29 @@ +package com.profiler; + +import java.util.logging.Logger; + +public class LifeCycleEventListener { + + private final static Logger logger = Logger.getLogger(LifeCycleEventListener.class.getName()); + + private static boolean started = false; + + public synchronized static void start() { + if (started) { + logger.info("already started"); + return; + } + + Agent.startAgent(); + started = true; + } + + public synchronized static void stop() { + if (!started) { + logger.info("already stopped"); + return; + } + started = false; + Agent.stopAgent(); + } +} diff --git a/src/main/java/com/profiler/SystemMonitor.java b/src/main/java/com/profiler/SystemMonitor.java index 414f94511..dcb356d77 100644 --- a/src/main/java/com/profiler/SystemMonitor.java +++ b/src/main/java/com/profiler/SystemMonitor.java @@ -22,56 +22,70 @@ import static com.profiler.config.ProfilerConfig.JVM_STAT_GAP; /** * System monitor - * + * * @author netspider - * */ @SuppressWarnings("restriction") public class SystemMonitor { - private static final Logger logger = Logger.getLogger(SystemMonitor.class.getName()); + private static final Logger logger = Logger.getLogger(SystemMonitor.class.getName()); - private final ScheduledExecutorService executor = new ScheduledThreadPoolExecutor(1, new ThreadFactory() { - @Override - public Thread newThread(Runnable runnable) { - Thread t = new Thread(runnable); - t.setName("HIPPO-SystemMonitor"); - t.setDaemon(true); - return t; - } - }); + private final ScheduledExecutorService executor = new ScheduledThreadPoolExecutor(1, new ThreadFactory() { + @Override + public Thread newThread(Runnable runnable) { + Thread t = new Thread(runnable); + t.setName("HIPPO-SystemMonitor"); + t.setDaemon(true); + return t; + } + }); - public void start() { - logger.info("Starting system monitor."); - executor.scheduleAtFixedRate(new Worker(), 5, 5, TimeUnit.SECONDS); - } + private DataSender dataSender; - public void stop() { - logger.info("Stopping system monitor"); - executor.shutdown(); - } + public SystemMonitor() { + } - static class Worker implements Runnable { + public void setDataSender(DataSender udpDataSender) { + this.dataSender = udpDataSender; + } - public void run() { + public void start() { + logger.info("Starting system monitor."); + executor.scheduleAtFixedRate(new Worker(dataSender), 5, 5, TimeUnit.SECONDS); + } + + public void stop() { + logger.info("Stopping system monitor"); + executor.shutdown(); + } + + static class Worker implements Runnable { + + private DataSender dataSender; + + public Worker(DataSender dataSender) { + this.dataSender = dataSender; + } + + public void run() { JVMInfoThriftDTO jvmInfo = new JVMInfoThriftDTO(); - try { - jvmInfo = new JVMInfoThriftDTO(); - jvmInfo.setAgentId(Agent.getInstance().getAgentId()); - jvmInfo.setDataTime(System.currentTimeMillis()); + try { + jvmInfo = new JVMInfoThriftDTO(); + jvmInfo.setAgentId(Agent.getInstance().getAgentId()); + jvmInfo.setDataTime(System.currentTimeMillis()); TraceContext traceContext = TraceContext.getTraceContext(); activeThread(traceContext, jvmInfo); setGCState(jvmInfo); - setMemoryState(jvmInfo); - setProcessCPUUsage(jvmInfo); + setMemoryState(jvmInfo); + setProcessCPUUsage(jvmInfo); - DataSender.getInstance().addDataToSend(jvmInfo); - } catch (Exception e) { + dataSender.send(jvmInfo); + } catch (Exception e) { logger.log(Level.INFO, "JvmInfo collect error Cause:" + e.getMessage(), e); - } - } + } + } private void activeThread(TraceContext traceContext, JVMInfoThriftDTO jvmInfo) { int activeThread = traceContext.getActiveThreadCounter().getActiveThread(); @@ -79,69 +93,68 @@ public class SystemMonitor { } - private void setGCState(JVMInfoThriftDTO jvmInfo) throws Exception { - List list = ManagementFactory.getGarbageCollectorMXBeans(); - if (list.size() == 2) { + private void setGCState(JVMInfoThriftDTO jvmInfo) throws Exception { + List list = ManagementFactory.getGarbageCollectorMXBeans(); + if (list.size() == 2) { // 제네레이션 기반일 경우 young, old 2개. // young.getName() // young gc type - GarbageCollectorMXBean young = list.get(0); + GarbageCollectorMXBean young = list.get(0); jvmInfo.setGc1Count(young.getCollectionCount()); jvmInfo.setGc1Time(young.getCollectionTime()); - GarbageCollectorMXBean old = list.get(1); + GarbageCollectorMXBean old = list.get(1); // old.getName() // old gc type jvmInfo.setGc2Count(old.getCollectionCount()); jvmInfo.setGc2Time(old.getCollectionTime()); - } else { + } else { // g1 ? - if(logger.isLoggable(Level.FINE)) { - logger.fine("unknown gc type. gc collector size:" + list.size()); + if (logger.isLoggable(Level.FINE)) { + logger.fine("unknown gc type. gc collector size:" + list.size()); } } - } + } - public void setMemoryState(JVMInfoThriftDTO jvmInfo) throws Exception { - MemoryMXBean bean = ManagementFactory.getMemoryMXBean(); - MemoryUsage heap = bean.getHeapMemoryUsage(); - MemoryUsage nonHeap = bean.getNonHeapMemoryUsage(); + public void setMemoryState(JVMInfoThriftDTO jvmInfo) throws Exception { + MemoryMXBean bean = ManagementFactory.getMemoryMXBean(); + MemoryUsage heap = bean.getHeapMemoryUsage(); + MemoryUsage nonHeap = bean.getNonHeapMemoryUsage(); jvmInfo.setHeapUsed(heap.getUsed()); jvmInfo.setHeapCommitted(heap.getCommitted()); jvmInfo.setNonHeapUsed(nonHeap.getUsed()); jvmInfo.setNonHeapCommitted(nonHeap.getCommitted()); - } + } - long previousCpuTime = 0; - int processorCount = -1; - boolean processCPUAvailable = true; + long previousCpuTime = 0; + int processorCount = -1; + boolean processCPUAvailable = true; - /** - * I don't know why should I divide by 10 in this result. But it works. - * - * @throws Exception - */ - private void setProcessCPUUsage(JVMInfoThriftDTO jvmInfo) throws Exception { - try { - if (processCPUAvailable) { - OperatingSystemMXBean sunOSMBean = ManagementFactory.newPlatformMXBeanProxy(ManagementFactory.getPlatformMBeanServer(), ManagementFactory.OPERATING_SYSTEM_MXBEAN_NAME, OperatingSystemMXBean.class); - long cpuTime = sunOSMBean.getProcessCpuTime(); + /** + * I don't know why should I divide by 10 in this result. But it works. + * + * @throws Exception + */ + private void setProcessCPUUsage(JVMInfoThriftDTO jvmInfo) throws Exception { + try { + if (processCPUAvailable) { + OperatingSystemMXBean sunOSMBean = ManagementFactory.newPlatformMXBeanProxy(ManagementFactory.getPlatformMBeanServer(), ManagementFactory.OPERATING_SYSTEM_MXBEAN_NAME, OperatingSystemMXBean.class); + long cpuTime = sunOSMBean.getProcessCpuTime(); - if (processorCount == -1) { - processorCount = sunOSMBean.getAvailableProcessors(); - } + if (processorCount == -1) { + processorCount = sunOSMBean.getAvailableProcessors(); + } - if (previousCpuTime != 0) { - long usedCPUTotal = (cpuTime - previousCpuTime) / 1000000; - double usedCPU = (0.1D * usedCPUTotal) / (processorCount * JVM_STAT_GAP / 1000.0); + if (previousCpuTime != 0) { + long usedCPUTotal = (cpuTime - previousCpuTime) / 1000000; + double usedCPU = (0.1D * usedCPUTotal) / (processorCount * JVM_STAT_GAP / 1000.0); jvmInfo.setProcessCPUTime(usedCPU); - } - previousCpuTime = cpuTime; + } + previousCpuTime = cpuTime; - } - } catch (IOException e) { - - processCPUAvailable = false; - } - } - } + } + } catch (IOException e) { + processCPUAvailable = false; + } + } + } } diff --git a/src/main/java/com/profiler/context/Trace.java b/src/main/java/com/profiler/context/Trace.java index baa7d114e..ec611b170 100644 --- a/src/main/java/com/profiler/context/Trace.java +++ b/src/main/java/com/profiler/context/Trace.java @@ -3,6 +3,7 @@ package com.profiler.context; import com.profiler.common.util.AnnotationTranscoder; import com.profiler.common.util.AnnotationTranscoder.Encoded; import com.profiler.sender.DataSender; +import com.profiler.sender.UdpDataSender; import com.profiler.util.NamedThreadLocal; import java.util.logging.Level; @@ -33,7 +34,7 @@ public final class Trace { try { TraceID nextId = getNextTraceId(); - traceIDStack.incr(); + traceIDStack.push(); if (traceIDStack.getTraceId() == null) { if (logger.isLoggable(Level.FINE)) { @@ -46,7 +47,7 @@ public final class Trace { } catch (Exception e) { e.printStackTrace(); } finally { - traceIDStack.decr(); + traceIDStack.pop(); } } @@ -59,7 +60,7 @@ public final class Trace { try { TraceID nextId = getNextTraceId(); - traceIDStack.incr(); + traceIDStack.push(); if (traceIDStack.getTraceId() == null) { traceIDStack.setTraceId(nextId); @@ -71,7 +72,7 @@ public final class Trace { public static void traceBlockEnd() { TraceIDStack traceIDStack = traceIdLocal.get(); - traceIDStack.decr(); + traceIDStack.pop(); } /** @@ -88,10 +89,10 @@ public final class Trace { } if (id == null) { - System.out.println("create new traceid"); - id = TraceID.newTraceId(); - // traceIdLocal.set(id); + if (logger.isLoggable(Level.INFO)) { + logger.info("create new traceid:" + id); + } if (stack == null) { traceIdLocal.set(new TraceIDStack()); @@ -111,7 +112,6 @@ public final class Trace { if (stack != null) { traceId = stack.getTraceId(); } else { - // TODO : remove this log. if (logger.isLoggable(Level.FINE)) { logger.log(Level.FINE, "#############################################################" + @@ -161,7 +161,6 @@ public final class Trace { } public static void setTraceId(TraceID traceId) { - // TODO: remove this, just for debugging. if (getCurrentTraceId() != null) { if (logger.isLoggable(Level.FINE)) { logger.log(Level.FINE, @@ -195,8 +194,9 @@ public final class Trace { static void logSpan(Span span) { try { - // TODO: send span to the server. - System.out.println("\n\n[WRITE SPAN] hashCode=" + span.hashCode() + ",\n\t " + span + ",\n\t SpanMap.size=" + spanMap.size() + ",\n\t CurrentThreadID=" + Thread.currentThread().getId() + ",\n\t CurrentThreadName=" + Thread.currentThread().getName() + "\n\n"); + if (logger.isLoggable(Level.FINE)) { + logger.info("[WRITE SPAN]" + span + " size=" + spanMap.size() + " CurrentThreadID=" + Thread.currentThread().getId() + ",\n\t CurrentThreadName=" + Thread.currentThread().getName() + "\n\n"); + } // TODO: remove this, just for debugging // if (spanMap.size() > 0) { @@ -206,7 +206,7 @@ public final class Trace { // System.out.println("current spamMap=" + spanMap); // } - DataSender.getInstance().addDataToSend(span.toThrift()); + UdpDataSender.getInstance().send(span.toThrift()); span.cancelTimer(); } catch (Exception e) { @@ -240,6 +240,7 @@ public final class Trace { mutate(getTraceIdOrCreateNew(), 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; diff --git a/src/main/java/com/profiler/context/TraceIDStack.java b/src/main/java/com/profiler/context/TraceIDStack.java index ae7183bba..20e364e16 100644 --- a/src/main/java/com/profiler/context/TraceIDStack.java +++ b/src/main/java/com/profiler/context/TraceIDStack.java @@ -1,46 +1,49 @@ package com.profiler.context; /** - * * @author netspider - * */ public class TraceIDStack { - private TraceID[] traceIDs = new TraceID[3]; + private TraceID[] stack = new TraceID[4]; - private volatile int index = 0; + // TODO 개별변수에 volatile 을 건다고 해서 전체의 동시성이 해결되지 않으므로 일단 제거. + private int index = 0; - public TraceID getTraceId() { - return traceIDs[index]; - } - public TraceID getParentTraceId() { - if (index > 0) { - return traceIDs[index - 1]; - } - return null; - } + public TraceID getTraceId() { + return stack[index]; + } - public void setTraceId(TraceID traceId) { - traceIDs[index] = traceId; - } + public TraceID getParentTraceId() { + if (index > 0) { + return stack[index - 1]; + } + return null; + } - public void incr() { - index++; - if (index > traceIDs.length - 1) { - TraceID[] old = traceIDs; - traceIDs = new TraceID[index + 1]; - System.arraycopy(old, 0, traceIDs, 0, old.length); - } - } + public void setTraceId(TraceID traceId) { + stack[index] = traceId; + } - public void decr() { - if (index > 0) - index--; - } + public void push() { + index++; + if (index > stack.length - 1) { + TraceID[] old = stack; + stack = new TraceID[index + 4]; + System.arraycopy(old, 0, stack, 0, old.length); + } + } - public void clear() { - traceIDs[index] = null; - } + public void pop() { + if (index > 0) { +// TODO 이전 reference를 제거해야 될거 같음. +// stack[index] = null; + index--; + } + } + + public void clear() { + stack[index] = null; + } } diff --git a/src/main/java/com/profiler/modifier/DefaultModifierRegistry.java b/src/main/java/com/profiler/modifier/DefaultModifierRegistry.java index 43b19d270..e8b834e9f 100644 --- a/src/main/java/com/profiler/modifier/DefaultModifierRegistry.java +++ b/src/main/java/com/profiler/modifier/DefaultModifierRegistry.java @@ -19,7 +19,7 @@ import com.profiler.modifier.db.oracle.OraclePreparedStatementModifier; import com.profiler.modifier.db.oracle.OracleResultSetModifier; import com.profiler.modifier.db.oracle.OracleStatementModifier; import com.profiler.modifier.tomcat.CatalinaModifier; -import com.profiler.modifier.tomcat.StandardHostValveInvokeInterceptor; +import com.profiler.modifier.tomcat.StandardHostValveInvokeModifier; import com.profiler.modifier.tomcat.TomcatConnectorModifier; import com.profiler.modifier.tomcat.TomcatStandardServiceModifier; @@ -60,8 +60,8 @@ public class DefaultModifierRegistry implements ModifierRegistry { } public void addTomcatModifier() { - StandardHostValveInvokeInterceptor standardHostValveInvokeInterceptor = new StandardHostValveInvokeInterceptor(byteCodeInstrumentor); - addModifier(standardHostValveInvokeInterceptor); + StandardHostValveInvokeModifier standardHostValveInvokeModifier = new StandardHostValveInvokeModifier(byteCodeInstrumentor); + addModifier(standardHostValveInvokeModifier); Modifier tomcatStandardServiceModifier = new TomcatStandardServiceModifier(byteCodeInstrumentor); addModifier(tomcatStandardServiceModifier); diff --git a/src/main/java/com/profiler/modifier/db/interceptor/DataSourceGetConnectionInterceptor.java b/src/main/java/com/profiler/modifier/db/interceptor/DataSourceGetConnectionInterceptor.java new file mode 100644 index 000000000..3235745af --- /dev/null +++ b/src/main/java/com/profiler/modifier/db/interceptor/DataSourceGetConnectionInterceptor.java @@ -0,0 +1,42 @@ +package com.profiler.modifier.db.interceptor; + +import com.profiler.interceptor.StaticAroundInterceptor; +import com.profiler.util.InterceptorUtils; +import com.profiler.util.StringUtils; + +import java.sql.Connection; +import java.util.Arrays; +import java.util.logging.Level; +import java.util.logging.Logger; + +/** + * Datasource의 get을 추적해야 될것으로 예상됨. + */ +public class DataSourceGetConnectionInterceptor implements StaticAroundInterceptor { + + private final Logger logger = Logger.getLogger(DataSourceGetConnectionInterceptor.class.getName()); + + @Override + public void before(Object target, String className, String methodName, String parameterDescription, Object[] args) { + if (logger.isLoggable(Level.INFO)) { + logger.info("before " + StringUtils.toString(target) + " " + className + "." + methodName + parameterDescription + " args:" + Arrays.toString(args)); + } + } + + @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); + } + + if (!InterceptorUtils.isSuccess(result)) { + return; + } + // TODO before도 같이 후킹하여 Connection 생성시간도 측정해야 됨. + // datasource의 pool을 고려할것. + if (result instanceof Connection) { + + } + } + +} diff --git a/src/main/java/com/profiler/modifier/tomcat/StandardHostValveInvokeInterceptor.java b/src/main/java/com/profiler/modifier/tomcat/StandardHostValveInvokeModifier.java similarity index 83% rename from src/main/java/com/profiler/modifier/tomcat/StandardHostValveInvokeInterceptor.java rename to src/main/java/com/profiler/modifier/tomcat/StandardHostValveInvokeModifier.java index 7075c4249..55f743371 100644 --- a/src/main/java/com/profiler/modifier/tomcat/StandardHostValveInvokeInterceptor.java +++ b/src/main/java/com/profiler/modifier/tomcat/StandardHostValveInvokeModifier.java @@ -1,13 +1,10 @@ package com.profiler.modifier.tomcat; -import com.profiler.config.ProfilerConstant; import com.profiler.interceptor.Interceptor; import com.profiler.interceptor.bci.ByteCodeInstrumentor; import com.profiler.interceptor.bci.InstrumentClass; import com.profiler.interceptor.bci.InstrumentException; import com.profiler.modifier.AbstractModifier; -import com.profiler.trace.RequestTracer; -import javassist.ByteArrayClassPath; import java.security.ProtectionDomain; import java.util.logging.Level; @@ -18,11 +15,11 @@ import java.util.logging.Logger; * * @author cowboy93, netspider */ -public class StandardHostValveInvokeInterceptor extends AbstractModifier { +public class StandardHostValveInvokeModifier extends AbstractModifier { - private final Logger logger = Logger.getLogger(StandardHostValveInvokeInterceptor.class.getName()); + private final Logger logger = Logger.getLogger(StandardHostValveInvokeModifier.class.getName()); - public StandardHostValveInvokeInterceptor(ByteCodeInstrumentor byteCodeInstrumentor) { + public StandardHostValveInvokeModifier(ByteCodeInstrumentor byteCodeInstrumentor) { super(byteCodeInstrumentor); } diff --git a/src/main/java/com/profiler/modifier/tomcat/TomcatStandardServiceModifier.java b/src/main/java/com/profiler/modifier/tomcat/TomcatStandardServiceModifier.java index 56ddb75f5..e26f23350 100644 --- a/src/main/java/com/profiler/modifier/tomcat/TomcatStandardServiceModifier.java +++ b/src/main/java/com/profiler/modifier/tomcat/TomcatStandardServiceModifier.java @@ -5,56 +5,51 @@ import java.util.logging.Level; import java.util.logging.Logger; import com.profiler.interceptor.bci.ByteCodeInstrumentor; -import javassist.CtClass; -import javassist.CtMethod; +import com.profiler.interceptor.bci.InstrumentClass; +import com.profiler.interceptor.bci.InstrumentException; +import com.profiler.modifier.tomcat.interceptors.StandardServiceStartInterceptor; +import com.profiler.modifier.tomcat.interceptors.StandardServiceStopInterceptor; -import com.profiler.Agent; import com.profiler.modifier.AbstractModifier; /** * When org.apache.catalina.core.StandardService class is loaded in ClassLoader, * this class modifies methods. - * + * * @author cowboy93, netspider - * */ public class TomcatStandardServiceModifier extends AbstractModifier { - private final Logger logger = Logger.getLogger(TomcatStandardServiceModifier.class.getName()); + private final Logger logger = Logger.getLogger(TomcatStandardServiceModifier.class.getName()); - public TomcatStandardServiceModifier(ByteCodeInstrumentor byteCodeInstrumentor) { - super(byteCodeInstrumentor); - } + public TomcatStandardServiceModifier(ByteCodeInstrumentor byteCodeInstrumentor) { + super(byteCodeInstrumentor); + } - public String getTargetClass() { - return "org/apache/catalina/core/StandardService"; - } + public String getTargetClass() { + return "org/apache/catalina/core/StandardService"; + } - public byte[] modify(ClassLoader classLoader, String javassistClassName, ProtectionDomain protectedDomain, byte[] classFileBuffer) { - if (logger.isLoggable(Level.INFO)) { - logger.info("Modifing. " + javassistClassName); - } - return changeMethod(javassistClassName, classFileBuffer); - } + public byte[] modify(ClassLoader classLoader, String javassistClassName, ProtectionDomain protectedDomain, byte[] classFileBuffer) { + if (logger.isLoggable(Level.INFO)) { + logger.info("Modifing. " + javassistClassName); + } + byteCodeInstrumentor.checkLibrary(classLoader, javassistClassName); - public byte[] changeMethod(String javassistClassName, byte[] classfileBuffer) { - try { - CtClass cc = classPool.get(javassistClassName); -// byteCodeInstrumentor.addInterceptor(, "startAgent", null); - CtMethod startMethod = cc.getDeclaredMethod("start", null); - startMethod.insertBefore("{" + Agent.FQCN + ".startAgent();" + "}"); + try { - CtMethod stopMethod = cc.getDeclaredMethod("stop", null); - stopMethod.insertBefore("{" + Agent.FQCN + ".stopAgent();" + "}"); + InstrumentClass standardService = byteCodeInstrumentor.getClass(javassistClassName); + StandardServiceStartInterceptor start = new StandardServiceStartInterceptor(); + standardService.addInterceptor("start", null, start); - printClassConvertComplete(javassistClassName); - return cc.toBytecode(); - } catch (Exception e) { - if (logger.isLoggable(Level.WARNING)) { - logger.log(Level.WARNING, e.getMessage(), e); - } - } - return null; - } + StandardServiceStopInterceptor stop = new StandardServiceStopInterceptor(); + standardService.addInterceptor("stop", null, stop); + + return standardService.toBytecode(); + } catch (InstrumentException e) { + logger.log(Level.WARNING, "modify fail. Cause:" + e.getMessage(), e); + return null; + } + } } 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 568d59085..e7a916139 100644 --- a/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardHostValveInvokeInterceptor.java +++ b/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardHostValveInvokeInterceptor.java @@ -22,51 +22,52 @@ import com.profiler.util.StringUtils; public class StandardHostValveInvokeInterceptor implements StaticAroundInterceptor { private final Logger logger = Logger.getLogger(StandardHostValveInvokeInterceptor.class.getName()); - @Override - public void before(Object target, String className, String methodName, String parameterDescription, Object[] args) { + @Override + public void before(Object target, String className, String methodName, String parameterDescription, Object[] args) { if (logger.isLoggable(Level.INFO)) { logger.info("before " + StringUtils.toString(target) + " " + className + "." + methodName + parameterDescription + " args:" + Arrays.toString(args)); } - try { + + try { TraceContext traceContext = TraceContext.getTraceContext(); traceContext.getActiveThreadCounter().start(); HttpServletRequest request = (HttpServletRequest) args[0]; - String requestURL = request.getRequestURI(); - String clientIP = request.getRemoteAddr(); - String parameters = getRequestParameter(request); + String requestURL = request.getRequestURI(); + String clientIP = request.getRemoteAddr(); + String parameters = getRequestParameter(request); - TraceID traceId = populateTraceIdFromRequest(request); - if (traceId != null) { - Trace.setTraceId(traceId); - } else { + TraceID traceId = populateTraceIdFromRequest(request); + if (traceId != null) { + Trace.setTraceId(traceId); + } else { TraceID newTraceID = TraceID.newTraceId(); if (logger.isLoggable(Level.INFO)) { logger.info("TraceID not exist. start new trace. " + newTraceID); // 좀더 자세한 정보는 debug레벨로 logger.log(Level.FINE, "requestUrl:" + requestURL + " clientIp" + clientIP + " parameter:" + parameters); } - Trace.setTraceId(newTraceID); - } - - Trace.recordRpcName("TOMCAT", requestURL); - Trace.recordEndPoint(request.getProtocol() + ":" + request.getLocalName() + ":" + request.getLocalPort()); - Trace.recordAttibute("http.url", request.getRequestURI()); - if (parameters != null && parameters.length() > 0) { - Trace.recordAttibute("http.params", parameters); - } - Trace.record(Annotation.ServerRecv); - - StopWatch.start("StandardHostValveInvokeInterceptor-starttime"); - } catch (Exception e) { - if (logger.isLoggable(Level.WARNING)) { - logger.log(Level.WARNING, "Tomcat StandardHostValve trace start fail", e); + Trace.setTraceId(newTraceID); } - } - } - @Override - public void after(Object target, String className, String methodName, String parameterDescription, Object[] args, Object result) { + Trace.recordRpcName("TOMCAT", requestURL); + Trace.recordEndPoint(request.getProtocol() + ":" + request.getLocalName() + ":" + request.getLocalPort()); + Trace.recordAttibute("http.url", request.getRequestURI()); + if (parameters != null && parameters.length() > 0) { + Trace.recordAttibute("http.params", parameters); + } + Trace.record(Annotation.ServerRecv); + + StopWatch.start("StandardHostValveInvokeModifier-starttime"); + } catch (Exception e) { + if (logger.isLoggable(Level.WARNING)) { + logger.log(Level.WARNING, "Tomcat StandardHostValve trace start fail", e); + } + } + } + + @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); } @@ -74,57 +75,57 @@ public class StandardHostValveInvokeInterceptor implements StaticAroundIntercept TraceContext traceContext = TraceContext.getTraceContext(); traceContext.getActiveThreadCounter().end(); - // TODO result 가 Exception 타입일경우 호출 실패임. - Trace.record(Annotation.ServerSend, StopWatch.stopAndGetElapsed("StandardHostValveInvokeInterceptor-starttime")); + // TODO result 가 Exception 타입일경우 호출 실패임. + Trace.record(Annotation.ServerSend, StopWatch.stopAndGetElapsed("StandardHostValveInvokeModifier-starttime")); // RequestTracer.endTransaction(); - - // TODO: I'v changed point of removing. Trace.mutate() - // Trace.removeTraceId(); - } - /** - * Pupulate source trace from HTTP Header. - * - * @param request - * @return - */ - private TraceID populateTraceIdFromRequest(HttpServletRequest request) { - String strUUID = request.getHeader(Header.HTTP_TRACE_ID.toString()); - if (strUUID != null) { - UUID uuid = UUID.fromString(strUUID); - long parentSpanID = NumberUtils.parseLong(request.getHeader(Header.HTTP_PARENT_SPAN_ID.toString()), SpanID.NULL); - long spanID = NumberUtils.parseLong(request.getHeader(Header.HTTP_SPAN_ID.toString()), SpanID.NULL); - boolean sampled = Boolean.parseBoolean(request.getHeader(Header.HTTP_SAMPLED.toString())); - int flags = NumberUtils.parseInteger(request.getHeader(Header.HTTP_FLAGS.toString()), 0); + // TODO: I'v changed point of removing. Trace.mutate() + // Trace.removeTraceId(); + } - TraceID id = new TraceID(uuid, parentSpanID, spanID, sampled, flags); - if (logger.isLoggable(Level.INFO)) { - logger.info("TraceID exist. continue trace. " + id); + /** + * Pupulate source trace from HTTP Header. + * + * @param request + * @return + */ + private TraceID populateTraceIdFromRequest(HttpServletRequest request) { + String strUUID = request.getHeader(Header.HTTP_TRACE_ID.toString()); + if (strUUID != null) { + UUID uuid = UUID.fromString(strUUID); + long parentSpanID = NumberUtils.parseLong(request.getHeader(Header.HTTP_PARENT_SPAN_ID.toString()), SpanID.NULL); + long spanID = NumberUtils.parseLong(request.getHeader(Header.HTTP_SPAN_ID.toString()), SpanID.NULL); + boolean sampled = Boolean.parseBoolean(request.getHeader(Header.HTTP_SAMPLED.toString())); + int flags = NumberUtils.parseInteger(request.getHeader(Header.HTTP_FLAGS.toString()), 0); + + TraceID id = new TraceID(uuid, parentSpanID, spanID, sampled, flags); + if (logger.isLoggable(Level.INFO)) { + logger.info("TraceID exist. continue trace. " + id); } - return id; - } else { - return null; - } - } + return id; + } else { + return null; + } + } - private String getRequestParameter(HttpServletRequest request) { - Enumeration attrs = request.getParameterNames(); + private String getRequestParameter(HttpServletRequest request) { + Enumeration attrs = request.getParameterNames(); - StringBuilder params = new StringBuilder(); + StringBuilder params = new StringBuilder(); - while (attrs.hasMoreElements()) { - String keyString = attrs.nextElement().toString(); - Object value = request.getParameter(keyString); + while (attrs.hasMoreElements()) { + String keyString = attrs.nextElement().toString(); + Object value = request.getParameter(keyString); - if (value != null) { - String valueString = value.toString(); - int valueStringLength = valueString.length(); + if (value != null) { + String valueString = value.toString(); + int valueStringLength = valueString.length(); - if (valueStringLength > 0 && valueStringLength < 100) - params.append(keyString).append("=").append(valueString); - } - } + if (valueStringLength > 0 && valueStringLength < 100) + params.append(keyString).append("=").append(valueString); + } + } - return params.toString(); - } + return params.toString(); + } } diff --git a/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardServiceStartInterceptor.java b/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardServiceStartInterceptor.java new file mode 100644 index 000000000..5943ad26f --- /dev/null +++ b/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardServiceStartInterceptor.java @@ -0,0 +1,28 @@ +package com.profiler.modifier.tomcat.interceptors; + +import com.profiler.LifeCycleEventListener; +import com.profiler.interceptor.StaticAfterInterceptor; +import com.profiler.util.InterceptorUtils; +import com.profiler.util.StringUtils; + +import java.util.Arrays; +import java.util.logging.Level; +import java.util.logging.Logger; + +/** + * + */ +public class StandardServiceStartInterceptor implements StaticAfterInterceptor { + private final Logger logger = Logger.getLogger(StandardServiceStartInterceptor.class.getName()); + + @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); + } +// if (!InterceptorUtils.isSuccess(result)) { +// return; +// } + LifeCycleEventListener.start(); + } +} diff --git a/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardServiceStopInterceptor.java b/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardServiceStopInterceptor.java new file mode 100644 index 000000000..97a37a04c --- /dev/null +++ b/src/main/java/com/profiler/modifier/tomcat/interceptors/StandardServiceStopInterceptor.java @@ -0,0 +1,29 @@ +package com.profiler.modifier.tomcat.interceptors; + +import com.profiler.LifeCycleEventListener; +import com.profiler.interceptor.StaticAfterInterceptor; +import com.profiler.util.InterceptorUtils; +import com.profiler.util.StringUtils; + +import java.util.Arrays; +import java.util.logging.Level; +import java.util.logging.Logger; + +/** + * + */ +public class StandardServiceStopInterceptor implements StaticAfterInterceptor { + private final Logger logger = Logger.getLogger(StandardServiceStopInterceptor.class.getName()); + + @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); + } + // TODO 시작이 실패했을때 stop이 불러 지는가? +// if (!InterceptorUtils.isSuccess(result)) { +// return; +// } + LifeCycleEventListener.stop(); + } +} diff --git a/src/main/java/com/profiler/sampler/RandomSampler.java b/src/main/java/com/profiler/sampler/RandomSampler.java new file mode 100644 index 000000000..860877459 --- /dev/null +++ b/src/main/java/com/profiler/sampler/RandomSampler.java @@ -0,0 +1,38 @@ +package com.profiler.sampler; + +import java.util.Random; +import java.util.logging.Logger; + +/** + * + */ +public class RandomSampler { + + private final Logger logger = Logger.getLogger(this.getClass().getName()); + + private final Random random = new Random(); + private int samplingRate = 10; + + public RandomSampler(int samplingRate) { + if (samplingRate <= 0 || samplingRate >= 100) { + logger.warning("Invalid sampling rate. Expected range(0~100) " + samplingRate); + samplingRate = 0; + } + this.samplingRate = samplingRate; + } + + public boolean sample() { + int i = Math.abs(random.nextInt()) % 101; + return sample(i); + } + + public boolean sample(int seed) { + if (seed == 0) { + return false; + } + if (seed <= samplingRate) { + return true; + } + return false; + } +} diff --git a/src/main/java/com/profiler/sender/DataSender.java b/src/main/java/com/profiler/sender/DataSender.java index 7971a4f0e..f499c3fe8 100644 --- a/src/main/java/com/profiler/sender/DataSender.java +++ b/src/main/java/com/profiler/sender/DataSender.java @@ -1,132 +1,14 @@ -package com.profiler.sender; - -import com.profiler.common.dto.Header; -import com.profiler.common.util.DefaultTBaseLocator; -import com.profiler.common.util.HeaderTBaseSerializer; -import com.profiler.common.util.TBaseLocator; -import com.profiler.config.ProfilerConfig; -import org.apache.thrift.TBase; -import org.apache.thrift.TException; - -import java.io.IOException; -import java.net.DatagramPacket; -import java.net.DatagramSocket; -import java.net.InetSocketAddress; -import java.net.SocketException; -import java.util.concurrent.LinkedBlockingQueue; -import java.util.concurrent.TimeUnit; -import java.util.logging.Level; -import java.util.logging.Logger; - -/** - * @author netspider - */ -public class DataSender extends Thread { - - private final Logger logger = Logger.getLogger(DataSender.class.getName()); - - private final LinkedBlockingQueue> addedQueue = new LinkedBlockingQueue>(4096); - - private final InetSocketAddress serverAddress = new InetSocketAddress(ProfilerConfig.SERVER_IP, ProfilerConfig.SERVER_UDP_PORT); - - private DatagramSocket udpSocket = null; - private TBaseLocator locator = new DefaultTBaseLocator(); - // 주의 single thread용임 - private HeaderTBaseSerializer serializer = new HeaderTBaseSerializer(); - - private static class SingletonHolder { - public static final DataSender INSTANCE = new DataSender(); - } - - public static DataSender getInstance() { - return SingletonHolder.INSTANCE; - } - - private DataSender() { - udpSocket = createSocket(); - setName("HIPPO-DataSender"); - setDaemon(true); - start(); - } - - private DatagramSocket createSocket() { - try { - DatagramSocket datagramSocket = new DatagramSocket(); - datagramSocket.setSoTimeout(1000 * 5); - datagramSocket.connect(serverAddress); - return datagramSocket; - } catch (SocketException e) { - return null; - } - } - - public boolean addDataToSend(TBase data) { - // TODO: addedQueue가 full일 때 IllegalStateException처리. - return addedQueue.add(data); - } - - // TODO: sender thread가 한 개로 충분한가. - public void run() { - while (true) { - try { - TBase dto = take(); - if (dto == null) { - continue; - } - send(dto); - } catch (Exception e) { - logger.log(Level.WARNING, "Unexpected Error", e); - } - } - } - - private void send(TBase dto) { - byte[] sendData = serialize(dto); - if (sendData == null) { - return; - } - DatagramPacket packet = new DatagramPacket(sendData, sendData.length); - if (udpSocket == null) { - // socket생성에 문제가 있으면 재생성? - udpSocket = createSocket(); - } - if (udpSocket != null) { - try { - udpSocket.send(packet); - if (logger.isLoggable(Level.FINE)) { - logger.fine("Data sent. " + dto); - } - } catch (IOException e) { - logger.log(Level.WARNING, "packet send error " + dto, e); - } - } - } - - // TODO: addedqueue에서 bulk로 drain - private TBase take() { - try { - return addedQueue.poll(5, TimeUnit.SECONDS); - } catch (InterruptedException e) { - return null; - } - } - - private byte[] serialize(TBase dto) { - Header header = createHeader(dto); - try { - return serializer.serialize(header, dto); - } catch (TException e) { - if (logger.isLoggable(Level.INFO)) { - logger.log(Level.INFO, "Serialize fail:" + dto, e); - } - return null; - } - } - - private Header createHeader(TBase dto) { - short type = locator.typeLookup(dto); - Header header = new Header(); - header.setType(type); - return header; - } -} +package com.profiler.sender; + +import org.apache.thrift.TBase; + +/** + * + */ +public interface DataSender { + + boolean send(TBase data); + + void stop(); + +} diff --git a/src/main/java/com/profiler/sender/UdpDataSender.java b/src/main/java/com/profiler/sender/UdpDataSender.java new file mode 100644 index 000000000..8022c132b --- /dev/null +++ b/src/main/java/com/profiler/sender/UdpDataSender.java @@ -0,0 +1,152 @@ +package com.profiler.sender; + +import com.profiler.common.dto.Header; +import com.profiler.common.util.DefaultTBaseLocator; +import com.profiler.common.util.HeaderTBaseSerializer; +import com.profiler.common.util.TBaseLocator; +import com.profiler.config.ProfilerConfig; +import org.apache.thrift.TBase; +import org.apache.thrift.TException; + +import java.io.IOException; +import java.net.DatagramPacket; +import java.net.DatagramSocket; +import java.net.InetSocketAddress; +import java.net.SocketException; +import java.util.concurrent.LinkedBlockingQueue; +import java.util.concurrent.TimeUnit; +import java.util.logging.Level; +import java.util.logging.Logger; + +/** + * @author netspider + */ +public class UdpDataSender implements DataSender, Runnable { + + private final Logger logger = Logger.getLogger(UdpDataSender.class.getName()); + + private final LinkedBlockingQueue> queue = new LinkedBlockingQueue>(1024); + + private final InetSocketAddress serverAddress = new InetSocketAddress(ProfilerConfig.SERVER_IP, ProfilerConfig.SERVER_UDP_PORT); + + private DatagramSocket udpSocket = null; + private TBaseLocator locator = new DefaultTBaseLocator(); + // 주의 single thread용임 + private HeaderTBaseSerializer serializer = new HeaderTBaseSerializer(); + + private boolean started = false; + private Object stopLock = new Object(); + + private Thread ioThread; + + private static class SingletonHolder { + public static final UdpDataSender INSTANCE = new UdpDataSender(); + } + + + public static UdpDataSender getInstance() { + return SingletonHolder.INSTANCE; + } + + private UdpDataSender() { + udpSocket = createSocket(); + ioThread = new Thread(this); + ioThread.setName("HIPPO-DataSender"); + ioThread.setDaemon(true); + ioThread.start(); + started = true; + } + + private DatagramSocket createSocket() { + try { + DatagramSocket datagramSocket = new DatagramSocket(); + datagramSocket.setSoTimeout(1000 * 5); + datagramSocket.connect(serverAddress); + return datagramSocket; + } catch (SocketException e) { + return null; + } + } + + public boolean send(TBase data) { + if (!started) { + return false; + } + // TODO: addedQueue가 full일 때 IllegalStateException처리. + return queue.offer(data); + } + + @Override + public void stop() { + if (!started) { + return; + } + started = false; + // io thread 안전 종료. queue 비우기. + } + + // TODO: sender thread가 한 개로 충분한가. + public void run() { + while (true) { + try { + TBase dto = take(); + if (dto == null) { + continue; + } + send0(dto); + } catch (Exception e) { + logger.log(Level.WARNING, "Unexpected Error", e); + } + } + } + + private void send0(TBase dto) { + byte[] sendData = serialize(dto); + if (sendData == null) { + return; + } + DatagramPacket packet = new DatagramPacket(sendData, sendData.length); + if (udpSocket == null) { + // socket생성에 문제가 있으면 재생성? + udpSocket = createSocket(); + } + if (udpSocket != null) { + try { + udpSocket.send(packet); + if (logger.isLoggable(Level.FINE)) { + logger.fine("Data sent. " + dto); + } + } catch (IOException e) { + logger.log(Level.WARNING, "packet send error " + dto, e); + } + } + } + + // TODO: addedqueue에서 bulk로 drain + private TBase take() { + try { + return queue.poll(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + return null; + } + } + + private byte[] serialize(TBase dto) { + Header header = createHeader(dto); + try { + return serializer.serialize(header, dto); + } catch (TException e) { + if (logger.isLoggable(Level.INFO)) { + logger.log(Level.INFO, "Serialize fail:" + dto, e); + } + return null; + } + } + + private Header createHeader(TBase dto) { + short type = locator.typeLookup(dto); + Header header = new Header(); + header.setType(type); + return header; + } +} diff --git a/src/main/java/com/profiler/trace/RequestTracer.java b/src/main/java/com/profiler/trace/RequestTracer.java index 44452662e..042d9dd71 100644 --- a/src/main/java/com/profiler/trace/RequestTracer.java +++ b/src/main/java/com/profiler/trace/RequestTracer.java @@ -4,34 +4,35 @@ import com.profiler.Agent; import com.profiler.common.dto.thrift.RequestDataListThriftDTO; import com.profiler.common.dto.thrift.RequestThriftDTO; import com.profiler.config.ProfilerConstant; -import com.profiler.sender.DataSender; +import com.profiler.sender.UdpDataSender; import com.profiler.util.SystemUtils; import java.util.Collections; import java.util.HashSet; import java.util.Set; +@Deprecated public class RequestTracer { - public static final String FQCN = RequestTracer.class.getName(); + public static final String FQCN = RequestTracer.class.getName(); - private static final ThreadLocal currentRequestID = new ThreadLocal(); - private static final ThreadLocal currentRequestHash = new ThreadLocal(); - private static final Set requestSet = Collections.synchronizedSet(new HashSet()); + private static final ThreadLocal currentRequestID = new ThreadLocal(); + private static final ThreadLocal currentRequestHash = new ThreadLocal(); + private static final Set requestSet = Collections.synchronizedSet(new HashSet()); - public static void startTransaction(String requestURL, String clientIP, long requestTime, String parameters) { - long cpuUserTime[] = SystemUtils.getThreadTime(); + public static void startTransaction(String requestURL, String clientIP, long requestTime, String parameters) { + long cpuUserTime[] = SystemUtils.getThreadTime(); - String tempRequestID = Thread.currentThread().getName() + "_" + System.nanoTime(); - int tempRequestHashCode = tempRequestID.hashCode(); + String tempRequestID = Thread.currentThread().getName() + "_" + System.nanoTime(); + int tempRequestHashCode = tempRequestID.hashCode(); - currentRequestID.set(tempRequestID); - currentRequestHash.set(tempRequestHashCode); - requestSet.add(tempRequestID); + currentRequestID.set(tempRequestID); + currentRequestHash.set(tempRequestHashCode); + requestSet.add(tempRequestID); - RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), tempRequestHashCode, ProfilerConstant.DATA_TYPE_REQUEST, requestTime, cpuUserTime[0], cpuUserTime[1]); - dto.setClientIP(clientIP); - dto.setRequestURL(requestURL); + RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), tempRequestHashCode, ProfilerConstant.DATA_TYPE_REQUEST, requestTime, cpuUserTime[0], cpuUserTime[1]); + dto.setClientIP(clientIP); + dto.setRequestURL(requestURL); // int paramsLength = params.length(); // if (paramsLength > 0) { @@ -39,61 +40,61 @@ public class RequestTracer { // dto.setExtraData1(params.toString()); // } - DataSender.getInstance().addDataToSend(dto); - } + UdpDataSender.getInstance().send(dto); + } - /** - * Transaction is successfully ended. - */ - public static void endTransaction() { - long cpuUserTime[] = SystemUtils.getThreadTime(); - RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), currentRequestHash.get(), ProfilerConstant.DATA_TYPE_RESPONSE, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); + /** + * Transaction is successfully ended. + */ + public static void endTransaction() { + long cpuUserTime[] = SystemUtils.getThreadTime(); + RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), currentRequestHash.get(), ProfilerConstant.DATA_TYPE_RESPONSE, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); - finishTransaction(dto); - } + finishTransaction(dto); + } - /** - * There was an Exception processing transaction. - * - * @param throwable - */ - public static void exceptionTransaction(Throwable throwable) { - long cpuUserTime[] = SystemUtils.getThreadTime(); + /** + * There was an Exception processing transaction. + * + * @param throwable + */ + public static void exceptionTransaction(Throwable throwable) { + long cpuUserTime[] = SystemUtils.getThreadTime(); - RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), currentRequestHash.get(), ProfilerConstant.DATA_TYPE_UNCAUGHT_EXCEPTION, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); + RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), currentRequestHash.get(), ProfilerConstant.DATA_TYPE_UNCAUGHT_EXCEPTION, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); - dto.setExtraData1(throwable.getMessage()); + dto.setExtraData1(throwable.getMessage()); - StackTraceElement[] tempElement = throwable.getStackTrace(); - dto.setExtraData2(tempElement[0].toString()); + StackTraceElement[] tempElement = throwable.getStackTrace(); + dto.setExtraData2(tempElement[0].toString()); - finishTransaction(dto); - } + finishTransaction(dto); + } - /** - * Transaction is ended and send request end data - * - * @param dto - */ - private static void finishTransaction(RequestThriftDTO dto) { - RequestDataListThriftDTO dataListDto = DatabaseRequestTracer.getRequestDataList(); + /** + * Transaction is ended and send request end data + * + * @param dto + */ + private static void finishTransaction(RequestThriftDTO dto) { + RequestDataListThriftDTO dataListDto = DatabaseRequestTracer.getRequestDataList(); - if (dataListDto != null) { - DataSender.getInstance().addDataToSend(dataListDto); - } + if (dataListDto != null) { + UdpDataSender.getInstance().send(dataListDto); + } - DataSender.getInstance().addDataToSend(dto); + UdpDataSender.getInstance().send(dto); - requestSet.remove(currentRequestID.get()); - DatabaseRequestTracer.removeRequestDataList(); - DatabaseRequestTracer.removeFetchCount(); - } + requestSet.remove(currentRequestID.get()); + DatabaseRequestTracer.removeRequestDataList(); + DatabaseRequestTracer.removeFetchCount(); + } - public static int getActiveThreadCount() { - return requestSet.size(); - } + public static int getActiveThreadCount() { + return requestSet.size(); + } - public static Integer getCurrentRequestHash() { - return currentRequestHash.get(); - } + public static Integer getCurrentRequestHash() { + return currentRequestHash.get(); + } } diff --git a/src/test/java/com/profiler/sampler/RandomSamplerTest.java b/src/test/java/com/profiler/sampler/RandomSamplerTest.java new file mode 100644 index 000000000..56130af9a --- /dev/null +++ b/src/test/java/com/profiler/sampler/RandomSamplerTest.java @@ -0,0 +1,41 @@ +package com.profiler.sampler; + + +import org.junit.Assert; +import org.junit.Test; + +/** + * + */ +public class RandomSamplerTest { + @Test + public void test() { + RandomSampler randomSampler = new RandomSampler(10); + assertChoice(randomSampler, 10); + assertChoice(randomSampler, 5); + assertChoice(randomSampler, 1); + + assertDrop(randomSampler, 0); + assertDrop(randomSampler, 11); + assertDrop(randomSampler, 210); + } + + @Test + public void mod() { + int i = 0 % 101; + System.out.println("" + i); + + int j = Math.abs(-102) % 101; + System.out.println("" + j); + } + + private void assertDrop(RandomSampler randomSampler, int seed) { + boolean sample = randomSampler.sample(seed); + Assert.assertFalse(sample); + } + + private void assertChoice(RandomSampler randomSampler, int seed) { + boolean sample = randomSampler.sample(seed); + Assert.assertTrue(sample); + } +} diff --git a/src/test/java/com/profiler/sender/DataSenderTest.java b/src/test/java/com/profiler/sender/UdpDataSenderTest.java similarity index 74% rename from src/test/java/com/profiler/sender/DataSenderTest.java rename to src/test/java/com/profiler/sender/UdpDataSenderTest.java index ac6cf74fb..214029019 100644 --- a/src/test/java/com/profiler/sender/DataSenderTest.java +++ b/src/test/java/com/profiler/sender/UdpDataSenderTest.java @@ -4,12 +4,11 @@ import junit.framework.TestCase; import org.junit.Test; -public class DataSenderTest { +public class UdpDataSenderTest { @Test public void send() { - } }