From 67663f15152a0f8ec424efae66e1ccc037100d49 Mon Sep 17 00:00:00 2001 From: koo-taejin Date: Thu, 7 Aug 2014 11:42:00 +0900 Subject: [PATCH] =?UTF-8?q?#71=20profiler=20BO=EC=97=90=EC=84=9C=20?= =?UTF-8?q?=EB=A9=94=EC=8B=9C=EC=A7=80=20=EC=B2=98=EB=A6=AC=20=EA=B0=80?= =?UTF-8?q?=EB=8A=A5=ED=95=98=EB=8F=84=EB=A1=9D=20=ED=95=A8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 1. 코드변경에 따른 연관 코드 변경 --- .../tools/NetworkAvailabilityChecker.java | 58 +++++++++++++++- .../profiler/HeartBitCheckerStressTest.java | 49 ++++++++++++- .../profiler/HeartBitCheckerTest.java | 69 ++++++++++++++++--- .../sender/TcpDataSenderReconnectTest.java | 46 ++++++++++++- .../profiler/sender/TcpDataSenderTest.java | 47 ++++++++++++- .../nhn/pinpoint/profiler/util/MockAgent.java | 5 -- 6 files changed, 255 insertions(+), 19 deletions(-) diff --git a/src/main/java/com/nhn/pinpoint/profiler/tools/NetworkAvailabilityChecker.java b/src/main/java/com/nhn/pinpoint/profiler/tools/NetworkAvailabilityChecker.java index 344020ab1..b8073d9f4 100644 --- a/src/main/java/com/nhn/pinpoint/profiler/tools/NetworkAvailabilityChecker.java +++ b/src/main/java/com/nhn/pinpoint/profiler/tools/NetworkAvailabilityChecker.java @@ -1,9 +1,19 @@ package com.nhn.pinpoint.profiler.tools; +import java.util.Collections; + +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import com.nhn.pinpoint.bootstrap.config.ProfilerConfig; +import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.profiler.sender.DataSender; import com.nhn.pinpoint.profiler.sender.TcpDataSender; import com.nhn.pinpoint.profiler.sender.UdpDataSender; +import com.nhn.pinpoint.rpc.PinpointSocketException; +import com.nhn.pinpoint.rpc.client.MessageListener; +import com.nhn.pinpoint.rpc.client.PinpointSocket; +import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; /** * @@ -11,6 +21,9 @@ import com.nhn.pinpoint.profiler.sender.UdpDataSender; * */ public class NetworkAvailabilityChecker implements PinpointTools { + + private static final Logger LOGGER = LoggerFactory.getLogger(NetworkAvailabilityChecker.class); + public static void main(String[] args) { if (args.length != 1) { System.out.println("usage : " + NetworkAvailabilityChecker.class.getSimpleName() + " AGENT_CONFIG_FILE"); @@ -22,6 +35,9 @@ public class NetworkAvailabilityChecker implements PinpointTools { DataSender udpSender = null; DataSender udpSpanSender = null; DataSender tcpSender = null; + + PinpointSocketFactory socketFactory = null; + PinpointSocket socket = null; try { ProfilerConfig profilerConfig = new ProfilerConfig(); profilerConfig.readConfigFile(configPath); @@ -33,7 +49,11 @@ public class NetworkAvailabilityChecker implements PinpointTools { udpSender = new UdpDataSender(collector, uPort, "UDP", 10); udpSpanSender = new UdpDataSender(collector, usPort, "UDP-SPAN", 10); - tcpSender = new TcpDataSender(collector, tPort); + + socketFactory = createPinpointSocketFactory(); + socket = createPinpointSocket(collector, tPort, socketFactory); + + tcpSender = new TcpDataSender(socket); boolean udpSenderResult = udpSender.isNetworkAvailable(); boolean udpSpanSenderResult = udpSpanSender.isNetworkAvailable(); @@ -53,6 +73,13 @@ public class NetworkAvailabilityChecker implements PinpointTools { closeDataSender(udpSpanSender); closeDataSender(tcpSender); System.out.println("END."); + + if (socket != null) { + socket.close(); + } + if (socketFactory != null) { + socketFactory.release(); + } } } @@ -61,4 +88,33 @@ public class NetworkAvailabilityChecker implements PinpointTools { dataSender.stop(); } } + + private static PinpointSocketFactory createPinpointSocketFactory() { + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + pinpointSocketFactory.setTimeoutMillis(1000 * 5); + pinpointSocketFactory.setAgentProperties(Collections.EMPTY_MAP); + + return pinpointSocketFactory; + } + + + private static PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { + MessageListener messageListener = new CommandDispatcher(); + + PinpointSocket socket = null; + for (int i = 0; i < 3; i++) { + try { + socket = factory.connect(host, port, messageListener); + LOGGER.info("tcp connect success:{}/{}", host, port); + return socket; + } catch (PinpointSocketException e) { + LOGGER.warn("tcp connect fail:{}/{} try reconnect, retryCount:{}", host, port, i); + } + } + LOGGER.warn("change background tcp connect mode {}/{} ", host, port); + socket = factory.scheduledConnect(host, port, messageListener); + + return socket; + } + } diff --git a/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerStressTest.java b/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerStressTest.java index 359fc15ab..9c4acec50 100644 --- a/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerStressTest.java +++ b/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerStressTest.java @@ -1,14 +1,19 @@ package com.nhn.pinpoint.profiler; +import java.util.Collections; import java.util.Random; import java.util.concurrent.atomic.AtomicInteger; -import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializerFactory; import org.apache.thrift.TException; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.profiler.sender.TcpDataSender; +import com.nhn.pinpoint.rpc.PinpointSocketException; +import com.nhn.pinpoint.rpc.client.MessageListener; +import com.nhn.pinpoint.rpc.client.PinpointSocket; +import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.StreamPacket; @@ -19,6 +24,7 @@ import com.nhn.pinpoint.rpc.server.SocketChannel; import com.nhn.pinpoint.thrift.dto.TAgentInfo; import com.nhn.pinpoint.thrift.dto.TResult; import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializer; +import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializerFactory; public class HeartBitCheckerStressTest { @@ -36,7 +42,10 @@ public class HeartBitCheckerStressTest { ResponseServerMessageListener serverListener = new ResponseServerMessageListener(requestCount, successCount); - TcpDataSender sender = new TcpDataSender(HOST, PORT); + PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory); + + TcpDataSender sender = new TcpDataSender(socket); HeartBitChecker checker = new HeartBitChecker(sender, 1000L, getAgentInfo()); long strarTime = System.currentTimeMillis(); @@ -56,6 +65,14 @@ public class HeartBitCheckerStressTest { if (checker != null) { checker.stop(); } + + if (socket != null) { + socket.close(); + } + + if (socketFactory != null) { + socketFactory.release(); + } } } @@ -149,4 +166,32 @@ public class HeartBitCheckerStressTest { } } + private PinpointSocketFactory createPinpointSocketFactory() { + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + pinpointSocketFactory.setTimeoutMillis(1000 * 5); + pinpointSocketFactory.setAgentProperties(Collections.EMPTY_MAP); + + return pinpointSocketFactory; + } + + + private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { + MessageListener messageListener = new CommandDispatcher(); + + PinpointSocket socket = null; + for (int i = 0; i < 3; i++) { + try { + socket = factory.connect(host, port, messageListener); + logger.info("tcp connect success:{}/{}", host, port); + return socket; + } catch (PinpointSocketException e) { + logger.warn("tcp connect fail:{}/{} try reconnect, retryCount:{}", host, port, i); + } + } + logger.warn("change background tcp connect mode {}/{} ", host, port); + socket = factory.scheduledConnect(host, port, messageListener); + + return socket; + } + } diff --git a/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerTest.java b/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerTest.java index a5450f932..2c9827140 100644 --- a/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerTest.java +++ b/src/test/java/com/nhn/pinpoint/profiler/HeartBitCheckerTest.java @@ -1,9 +1,11 @@ package com.nhn.pinpoint.profiler; +import java.util.Collections; import java.util.concurrent.atomic.AtomicInteger; import com.nhn.pinpoint.thrift.io.HeaderTBaseSerializerFactory; + import junit.framework.Assert; import org.apache.thrift.TException; @@ -11,7 +13,12 @@ import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; import com.nhn.pinpoint.profiler.sender.TcpDataSender; +import com.nhn.pinpoint.rpc.PinpointSocketException; +import com.nhn.pinpoint.rpc.client.MessageListener; +import com.nhn.pinpoint.rpc.client.PinpointSocket; +import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.StreamPacket; @@ -38,7 +45,11 @@ public class HeartBitCheckerTest { ResponseServerMessageListener serverListener = new ResponseServerMessageListener(requestCount, successCount); PinpointServerSocket server = createServer(serverListener); - TcpDataSender sender = new TcpDataSender(HOST, PORT); + + PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory); + + TcpDataSender sender = new TcpDataSender(socket); HeartBitChecker checker = new HeartBitChecker(sender, 1000L, getAgentInfo()); try { @@ -48,7 +59,7 @@ public class HeartBitCheckerTest { Assert.assertEquals(1, successCount.get()); } finally { - closeAll(server, checker); + closeAll(server, checker, socket, socketFactory); } } @@ -61,7 +72,10 @@ public class HeartBitCheckerTest { PinpointServerSocket server = createServer(serverListener); - TcpDataSender sender = new TcpDataSender(HOST, PORT); + PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory); + + TcpDataSender sender = new TcpDataSender(socket); HeartBitChecker checker = new HeartBitChecker(sender, 1000L, getAgentInfo()); try { @@ -72,7 +86,7 @@ public class HeartBitCheckerTest { Assert.assertEquals(2, requestCount.get()); Assert.assertEquals(0, successCount.get()); } finally { - closeAll(server, checker); + closeAll(server, checker, socket, socketFactory); } } @@ -83,7 +97,10 @@ public class HeartBitCheckerTest { ResponseServerMessageListener serverListener = new ResponseServerMessageListener(requestCount, successCount); - TcpDataSender sender = new TcpDataSender(HOST, PORT); + PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory); + + TcpDataSender sender = new TcpDataSender(socket); HeartBitChecker checker = new HeartBitChecker(sender, 1000L, getAgentInfo()); try { @@ -97,9 +114,7 @@ public class HeartBitCheckerTest { Assert.assertEquals(3, successCount.get()); } finally { - if (checker != null) { - checker.stop(); - } + closeAll(null, checker, socket, socketFactory); } } @@ -125,7 +140,7 @@ public class HeartBitCheckerTest { } } - private void closeAll(PinpointServerSocket server, HeartBitChecker checker) { + private void closeAll(PinpointServerSocket server, HeartBitChecker checker, PinpointSocket socket, PinpointSocketFactory factory) { if (server != null) { server.close(); } @@ -133,6 +148,14 @@ public class HeartBitCheckerTest { if (checker != null) { checker.stop(); } + + if (socket != null) { + socket.close(); + } + + if (factory != null) { + factory.release(); + } } private TAgentInfo getAgentInfo() { @@ -192,5 +215,33 @@ public class HeartBitCheckerTest { logger.info("handleStreamPacket:{}", streamPacket); } } + + private PinpointSocketFactory createPinpointSocketFactory() { + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + pinpointSocketFactory.setTimeoutMillis(1000 * 5); + pinpointSocketFactory.setAgentProperties(Collections.EMPTY_MAP); + + return pinpointSocketFactory; + } + + + private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { + MessageListener messageListener = new CommandDispatcher(); + + PinpointSocket socket = null; + for (int i = 0; i < 3; i++) { + try { + socket = factory.connect(host, port, messageListener); + logger.info("tcp connect success:{}/{}", host, port); + return socket; + } catch (PinpointSocketException e) { + logger.warn("tcp connect fail:{}/{} try reconnect, retryCount:{}", host, port, i); + } + } + logger.warn("change background tcp connect mode {}/{} ", host, port); + socket = factory.scheduledConnect(host, port, messageListener); + + return socket; + } } diff --git a/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java b/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java index 56ad51a91..174bdc09c 100644 --- a/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java +++ b/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderReconnectTest.java @@ -1,6 +1,13 @@ package com.nhn.pinpoint.profiler.sender; +import java.util.Collections; + import com.nhn.pinpoint.thrift.dto.TApiMetaData; +import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; +import com.nhn.pinpoint.rpc.PinpointSocketException; +import com.nhn.pinpoint.rpc.client.MessageListener; +import com.nhn.pinpoint.rpc.client.PinpointSocket; +import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.StreamPacket; @@ -8,6 +15,9 @@ import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; + +import net.sf.cglib.proxy.Factory; + import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -52,7 +62,11 @@ public class TcpDataSenderReconnectTest { @Test public void connectAndSend() throws InterruptedException { PinpointServerSocket old = serverStart(); - TcpDataSender sender = new TcpDataSender(HOST, PORT); + + PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory); + + TcpDataSender sender = new TcpDataSender(socket); Thread.sleep(500); old.close(); @@ -69,5 +83,35 @@ public class TcpDataSenderReconnectTest { sender.stop(); pinpointServerSocket.close(); + socket.close(); + socketFactory.release(); + } + + private PinpointSocketFactory createPinpointSocketFactory() { + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + pinpointSocketFactory.setTimeoutMillis(1000 * 5); + pinpointSocketFactory.setAgentProperties(Collections.EMPTY_MAP); + + return pinpointSocketFactory; + } + + + private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { + MessageListener messageListener = new CommandDispatcher(); + + PinpointSocket socket = null; + for (int i = 0; i < 3; i++) { + try { + socket = factory.connect(host, port, messageListener); + logger.info("tcp connect success:{}/{}", host, port); + return socket; + } catch (PinpointSocketException e) { + logger.warn("tcp connect fail:{}/{} try reconnect, retryCount:{}", host, port, i); + } + } + logger.warn("change background tcp connect mode {}/{} ", host, port); + socket = factory.scheduledConnect(host, port, messageListener); + + return socket; } } diff --git a/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderTest.java b/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderTest.java index 6abef2738..cca2eb176 100644 --- a/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderTest.java +++ b/src/test/java/com/nhn/pinpoint/profiler/sender/TcpDataSenderTest.java @@ -1,6 +1,11 @@ package com.nhn.pinpoint.profiler.sender; import com.nhn.pinpoint.thrift.dto.TApiMetaData; +import com.nhn.pinpoint.profiler.receiver.CommandDispatcher; +import com.nhn.pinpoint.rpc.PinpointSocketException; +import com.nhn.pinpoint.rpc.client.MessageListener; +import com.nhn.pinpoint.rpc.client.PinpointSocket; +import com.nhn.pinpoint.rpc.client.PinpointSocketFactory; import com.nhn.pinpoint.rpc.packet.RequestPacket; import com.nhn.pinpoint.rpc.packet.SendPacket; import com.nhn.pinpoint.rpc.packet.StreamPacket; @@ -8,13 +13,16 @@ import com.nhn.pinpoint.rpc.server.PinpointServerSocket; import com.nhn.pinpoint.rpc.server.ServerMessageListener; import com.nhn.pinpoint.rpc.server.ServerStreamChannel; import com.nhn.pinpoint.rpc.server.SocketChannel; + import junit.framework.Assert; + import org.junit.After; import org.junit.Before; import org.junit.Test; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.Collections; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -68,7 +76,10 @@ public class TcpDataSenderTest { public void connectAndSend() throws InterruptedException { this.sendLatch = new CountDownLatch(2); - TcpDataSender sender = new TcpDataSender(HOST, PORT); + PinpointSocketFactory socketFactory = createPinpointSocketFactory(); + PinpointSocket socket = createPinpointSocket(HOST, PORT, socketFactory); + + TcpDataSender sender = new TcpDataSender(socket); try { sender.send(new TApiMetaData("test", System.currentTimeMillis(), 1, "TestApi")); sender.send(new TApiMetaData("test", System.currentTimeMillis(), 1, "TestApi")); @@ -78,8 +89,42 @@ public class TcpDataSenderTest { Assert.assertTrue(received); } finally { sender.stop(); + + if (socket != null) { + socket.close(); + } + + if (socketFactory != null) { + socketFactory.release(); + } } + } + + private PinpointSocketFactory createPinpointSocketFactory() { + PinpointSocketFactory pinpointSocketFactory = new PinpointSocketFactory(); + pinpointSocketFactory.setTimeoutMillis(1000 * 5); + pinpointSocketFactory.setAgentProperties(Collections.EMPTY_MAP); + return pinpointSocketFactory; + } + + private PinpointSocket createPinpointSocket(String host, int port, PinpointSocketFactory factory) { + MessageListener messageListener = new CommandDispatcher(); + + PinpointSocket socket = null; + for (int i = 0; i < 3; i++) { + try { + socket = factory.connect(host, port, messageListener); + logger.info("tcp connect success:{}/{}", host, port); + return socket; + } catch (PinpointSocketException e) { + logger.warn("tcp connect fail:{}/{} try reconnect, retryCount:{}", host, port, i); + } + } + logger.warn("change background tcp connect mode {}/{} ", host, port); + socket = factory.scheduledConnect(host, port, messageListener); + + return socket; } } diff --git a/src/test/java/com/nhn/pinpoint/profiler/util/MockAgent.java b/src/test/java/com/nhn/pinpoint/profiler/util/MockAgent.java index 6a73dd8f4..04ebc9f9c 100644 --- a/src/test/java/com/nhn/pinpoint/profiler/util/MockAgent.java +++ b/src/test/java/com/nhn/pinpoint/profiler/util/MockAgent.java @@ -29,11 +29,6 @@ public class MockAgent extends DefaultAgent { return super.createUdpDataSender(port, threadName, writeQueueSize, timeout, sendBufferSize); } - @Override - protected EnhancedDataSender createTcpDataSender() { - return new LoggingDataSender(); - } - @Override protected StorageFactory createStorageFactory() { return new ReadableSpanStorageFactory();