diff --git a/src/main/java/com/profiler/Agent.java b/src/main/java/com/profiler/Agent.java index 4c47f88d5..cc1244ba6 100644 --- a/src/main/java/com/profiler/Agent.java +++ b/src/main/java/com/profiler/Agent.java @@ -3,7 +3,8 @@ package com.profiler; import java.util.Map.Entry; import java.util.logging.Logger; -import com.profiler.sender.AgentInfoSender; +import com.profiler.common.dto.thrift.AgentInfo; +import com.profiler.sender.DataSender; public class Agent { @@ -16,8 +17,6 @@ public class Agent { private final ServerInfo serverInfo; private final SystemMonitor systemMonitor; - private int agentHash = 0; - private Agent() { this.serverInfo = new ServerInfo(); this.systemMonitor = new SystemMonitor(); @@ -48,19 +47,11 @@ public class Agent { * * @return */ - public int getAgentHashCode() { - // TODO: string format으로 변경, DTO때문에 일단 int로 구현함. - // TODO: hashcode 생성 방법도 변경해야함. - if (agentHash == 0) { - String portNumbers = ""; - for (Entry entry : serverInfo.getConnectors().entrySet()) { - portNumbers += " " + entry.getKey(); - } - portNumbers = portNumbers.trim(); - agentHash = (serverInfo.getHostip() + portNumbers).hashCode(); - } + public String getAgentId() { + // TODO: agent id 생성 방법 변경이 필요함. + logger.warning("Generating agent id is not implementd. use default 'TEST_AGENT_ID"); - return agentHash; + return "TEST_AGENT_ID"; } /** @@ -71,8 +62,21 @@ public class Agent { public void sendStartupInfo() { logger.info("Send startup information to HIPPO server."); - AgentInfoSender sender = new AgentInfoSender(true); - sender.start(); + String ip = getServerInfo().getHostip(); + String ports = ""; + for (Entry entry : getServerInfo().getConnectors().entrySet()) { + ports += " " + entry.getKey(); + } + + AgentInfo agentInfo = new AgentInfo(); + + agentInfo.setHostname(ip); + agentInfo.setPorts(ports); + agentInfo.setIsAlive(true); + agentInfo.setTimestamp(System.currentTimeMillis()); + agentInfo.setAgentId(getAgentId()); + + DataSender.getInstance().addDataToSend(agentInfo); } public void start() { @@ -84,8 +88,21 @@ public class Agent { logger.info("Stopping HIPPO Agent."); systemMonitor.stop(); - AgentInfoSender sender = new AgentInfoSender(false); - sender.start(); + String ip = getServerInfo().getHostip(); + String ports = ""; + for (Entry entry : getServerInfo().getConnectors().entrySet()) { + ports += " " + entry.getKey(); + } + + AgentInfo agentInfo = new AgentInfo(); + + agentInfo.setHostname(ip); + agentInfo.setPorts(ports); + agentInfo.setIsAlive(false); + agentInfo.setTimestamp(System.currentTimeMillis()); + agentInfo.setAgentId(getAgentId()); + + DataSender.getInstance().addDataToSend(agentInfo); } public static void startAgent() { diff --git a/src/main/java/com/profiler/SystemMonitor.java b/src/main/java/com/profiler/SystemMonitor.java index c566bf450..b38d2a78c 100644 --- a/src/main/java/com/profiler/SystemMonitor.java +++ b/src/main/java/com/profiler/SystemMonitor.java @@ -56,7 +56,7 @@ public class SystemMonitor { public void run() { try { currentDto = new JVMInfoThriftDTO(); - currentDto.setAgentHashCode(Agent.getInstance().getAgentHashCode()); + currentDto.setAgentId(Agent.getInstance().getAgentId()); currentDto.setDataTime(System.currentTimeMillis()); getActiveThreadCount(); diff --git a/src/main/java/com/profiler/context/Span.java b/src/main/java/com/profiler/context/Span.java index 15cc71392..63ef46526 100644 --- a/src/main/java/com/profiler/context/Span.java +++ b/src/main/java/com/profiler/context/Span.java @@ -124,7 +124,7 @@ public class Span { public com.profiler.common.dto.thrift.Span toThrift() { com.profiler.common.dto.thrift.Span span = new com.profiler.common.dto.thrift.Span(); - span.setAgentID(String.valueOf(Agent.getInstance().getAgentHashCode())); + span.setAgentID(Agent.getInstance().getAgentId()); span.setTimestamp(createTime); span.setMostTraceID(traceID.getId().getMostSignificantBits()); span.setLeastTraceID(traceID.getId().getLeastSignificantBits()); diff --git a/src/main/java/com/profiler/sender/AgentInfoSender.java b/src/main/java/com/profiler/sender/AgentInfoSender.java index b52383b69..7c77e35d8 100644 --- a/src/main/java/com/profiler/sender/AgentInfoSender.java +++ b/src/main/java/com/profiler/sender/AgentInfoSender.java @@ -10,6 +10,7 @@ import com.profiler.Agent; import com.profiler.common.dto.AgentInfoDTO; import com.profiler.config.TomcatProfilerConfig; +@Deprecated public class AgentInfoSender extends Thread { private final Logger logger = Logger.getLogger(AgentInfoSender.class.getName()); diff --git a/src/main/java/com/profiler/sender/DataSender.java b/src/main/java/com/profiler/sender/DataSender.java index fb8429fea..0b0c6d53f 100644 --- a/src/main/java/com/profiler/sender/DataSender.java +++ b/src/main/java/com/profiler/sender/DataSender.java @@ -19,119 +19,115 @@ import com.profiler.common.util.HeaderTBaseSerializer; import com.profiler.common.util.TBaseLocator; import com.profiler.config.TomcatProfilerConfig; - /** * @author netspider */ public class DataSender extends Thread { - private final Logger logger = Logger.getLogger(DataSender.class.getName()); + private final Logger logger = Logger.getLogger(DataSender.class.getName()); - private final LinkedBlockingQueue> addedQueue = new LinkedBlockingQueue>(4096); + private final LinkedBlockingQueue> addedQueue = new LinkedBlockingQueue>(4096); - private final InetSocketAddress serverAddress = new InetSocketAddress(TomcatProfilerConfig.SERVER_IP, TomcatProfilerConfig.DEFUALT_PORT); + private final InetSocketAddress serverAddress = new InetSocketAddress(TomcatProfilerConfig.SERVER_IP, TomcatProfilerConfig.DEFUALT_PORT); - private DatagramSocket udpSocket = null; - private TBaseLocator locator = new DefaultTBaseLocator(); - // 주의 single thread용임 - private HeaderTBaseSerializer serializer = new HeaderTBaseSerializer(); + 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(); - } + private static class SingletonHolder { + public static final DataSender INSTANCE = new DataSender(); + } - public static DataSender getInstance() { - return SingletonHolder.INSTANCE; - } + public static DataSender getInstance() { + return SingletonHolder.INSTANCE; + } + private DataSender() { + udpSocket = createSocket(); + setName("HIPPO-DataSender"); + setDaemon(true); + start(); + } - 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; + } + } - 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); + } - 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: 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); - } - } - } + // TODO: addedqueue에서 bulk로 drain + private TBase take() { + try { + return addedQueue.poll(5, TimeUnit.SECONDS); + } catch (InterruptedException e) { + return null; + } + } - 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); - } - } - } + 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; + } + } - - // 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; - } + 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/DatabaseRequestTracer.java b/src/main/java/com/profiler/trace/DatabaseRequestTracer.java index 4abfa4f08..d4f5b1141 100644 --- a/src/main/java/com/profiler/trace/DatabaseRequestTracer.java +++ b/src/main/java/com/profiler/trace/DatabaseRequestTracer.java @@ -3,7 +3,6 @@ package com.profiler.trace; import java.util.ArrayList; import java.util.HashMap; import java.util.HashSet; -import java.util.Hashtable; import java.util.Iterator; import java.util.List; import java.util.Set; @@ -12,7 +11,6 @@ import java.util.concurrent.ConcurrentMap; import java.util.concurrent.CopyOnWriteArraySet; import com.profiler.Agent; -import com.profiler.common.dto.AgentInfoDTO; import com.profiler.common.dto.thrift.RequestDataListThriftDTO; import com.profiler.common.dto.thrift.RequestDataThriftDTO; import com.profiler.config.TomcatProfilerConfig; @@ -261,7 +259,7 @@ public class DatabaseRequestTracer { */ private static RequestDataListThriftDTO checkDTO(RequestDataListThriftDTO dto) { if (dto == null) { - dto = new RequestDataListThriftDTO(Agent.getInstance().getAgentHashCode(), RequestTracer.getCurrentRequestHash(), new ArrayList()); + dto = new RequestDataListThriftDTO(Agent.getInstance().getAgentId(), RequestTracer.getCurrentRequestHash(), new ArrayList()); } return dto; } diff --git a/src/main/java/com/profiler/trace/RequestTracer.java b/src/main/java/com/profiler/trace/RequestTracer.java index 0d6e48fb8..77ab199aa 100644 --- a/src/main/java/com/profiler/trace/RequestTracer.java +++ b/src/main/java/com/profiler/trace/RequestTracer.java @@ -29,7 +29,7 @@ public class RequestTracer { currentRequestHash.set(tempRequestHashCode); requestSet.add(tempRequestID); - RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentHashCode(), tempRequestHashCode, TomcatProfilerConstant.DATA_TYPE_REQUEST, requestTime, cpuUserTime[0], cpuUserTime[1]); + RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), tempRequestHashCode, TomcatProfilerConstant.DATA_TYPE_REQUEST, requestTime, cpuUserTime[0], cpuUserTime[1]); dto.setClientIP(clientIP); dto.setRequestURL(requestURL); @@ -47,7 +47,7 @@ public class RequestTracer { */ public static void endTransaction() { long cpuUserTime[] = SystemUtils.getThreadTime(); - RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentHashCode(), currentRequestHash.get(), TomcatProfilerConstant.DATA_TYPE_RESPONSE, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); + RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), currentRequestHash.get(), TomcatProfilerConstant.DATA_TYPE_RESPONSE, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); finishTransaction(dto); } @@ -60,7 +60,7 @@ public class RequestTracer { public static void exceptionTransaction(Throwable throwable) { long cpuUserTime[] = SystemUtils.getThreadTime(); - RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentHashCode(), currentRequestHash.get(), TomcatProfilerConstant.DATA_TYPE_UNCAUGHT_EXCEPTION, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); + RequestThriftDTO dto = new RequestThriftDTO(Agent.getInstance().getAgentId(), currentRequestHash.get(), TomcatProfilerConstant.DATA_TYPE_UNCAUGHT_EXCEPTION, System.currentTimeMillis(), cpuUserTime[0], cpuUserTime[1]); dto.setExtraData1(throwable.getMessage()); diff --git a/src/test/java/com/profiler/util/HeaderTBaseSerializerTest.java b/src/test/java/com/profiler/util/HeaderTBaseSerializerTest.java index 5ae8ed489..df89a01c0 100644 --- a/src/test/java/com/profiler/util/HeaderTBaseSerializerTest.java +++ b/src/test/java/com/profiler/util/HeaderTBaseSerializerTest.java @@ -1,17 +1,17 @@ package com.profiler.util; +import java.util.Arrays; +import java.util.logging.Logger; + +import org.junit.Assert; +import org.junit.Test; + import com.profiler.common.dto.Header; import com.profiler.common.dto.thrift.JVMInfoThriftDTO; import com.profiler.common.util.DefaultTBaseLocator; import com.profiler.common.util.HeaderTBaseDeserializer; import com.profiler.common.util.HeaderTBaseSerializer; import com.profiler.common.util.TBaseLocator; -import org.apache.thrift.TBase; -import org.junit.Assert; -import org.junit.Test; - -import java.util.Arrays; -import java.util.logging.Logger; public class HeaderTBaseSerializerTest { private final Logger logger = Logger.getLogger(HeaderTBaseSerializerTest.class.getName()); @@ -28,8 +28,8 @@ public class HeaderTBaseSerializerTest { JVMInfoThriftDTO jvmInfoThriftDTO = new JVMInfoThriftDTO(); int activeThreadount = 10; jvmInfoThriftDTO.setActiveThreadCount(activeThreadount); - int agentHashCde = 123; - jvmInfoThriftDTO.setAgentHashCode(agentHashCde); + String agentId = "agentId"; + jvmInfoThriftDTO.setAgentId(agentId); byte[] serialize = serializer.serialize(header, jvmInfoThriftDTO); dump(serialize); @@ -39,7 +39,7 @@ public class HeaderTBaseSerializerTest { logger.info("deserialize:" + deserialize.getClass()); Assert.assertEquals(deserialize.getActiveThreadCount(), activeThreadount); - Assert.assertEquals(deserialize.getAgentHashCode(), agentHashCde); + Assert.assertEquals(deserialize.getAgentId(), agentId); } public void dump(byte[] data) {