TCP: Add Diagnostics

This commit is contained in:
Gary Russell
2015-08-11 15:56:23 -04:00
parent 31b8b37ab4
commit e2a10a8d79

View File

@@ -63,7 +63,11 @@ import javax.net.ServerSocketFactory;
import javax.net.SocketFactory;
import org.apache.commons.logging.Log;
import org.apache.commons.logging.LogFactory;
import org.apache.log4j.Level;
import org.junit.Rule;
import org.junit.Test;
import org.junit.rules.TestName;
import org.mockito.Matchers;
import org.mockito.Mockito;
import org.mockito.invocation.InvocationOnMock;
@@ -78,6 +82,7 @@ import org.springframework.integration.ip.tcp.serializer.MapJsonSerializer;
import org.springframework.integration.ip.util.TestingUtilities;
import org.springframework.integration.support.MessageBuilder;
import org.springframework.integration.support.converter.MapMessageConverter;
import org.springframework.integration.test.rule.Log4jLevelAdjuster;
import org.springframework.integration.test.util.SocketUtils;
import org.springframework.integration.test.util.TestUtils;
import org.springframework.integration.util.CompositeExecutor;
@@ -97,23 +102,22 @@ import org.springframework.util.ReflectionUtils.FieldFilter;
*/
public class TcpNioConnectionTests {
private final ApplicationEventPublisher nullPublisher = new ApplicationEventPublisher() {
private final static Log logger = LogFactory.getLog(TcpNioConnectionTests.class);
@Override
public void publishEvent(ApplicationEvent event) {
}
@Rule
public final Log4jLevelAdjuster adjuster = new Log4jLevelAdjuster(Level.TRACE,
"org.springframework.integration.ip.tcp");
@Override
public void publishEvent(Object event) {
@Rule
public TestName testName = new TestName();
}
};
private final ApplicationEventPublisher nullPublisher = mock(ApplicationEventPublisher.class);
@Test
public void testWriteTimeout() throws Exception {
final int port = SocketUtils.findAvailableServerSocket();
TcpNioClientConnectionFactory factory = new TcpNioClientConnectionFactory("localhost", port);
factory.setApplicationEventPublisher(nullPublisher);
factory.setSoTimeout(1000);
factory.start();
final CountDownLatch latch = new CountDownLatch(1);
@@ -124,6 +128,7 @@ public class TcpNioConnectionTests {
@SuppressWarnings("unused")
public void run() {
try {
logger.debug(testName.getMethodName() + " starting server for " + port);
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(port);
serverSocket.set(server);
latch.countDown();
@@ -153,6 +158,7 @@ public class TcpNioConnectionTests {
public void testReadTimeout() throws Exception {
final int port = SocketUtils.findAvailableServerSocket();
TcpNioClientConnectionFactory factory = new TcpNioClientConnectionFactory("localhost", port);
factory.setApplicationEventPublisher(nullPublisher);
factory.setSoTimeout(1000);
factory.start();
final CountDownLatch latch = new CountDownLatch(1);
@@ -162,6 +168,7 @@ public class TcpNioConnectionTests {
@Override
public void run() {
try {
logger.debug(testName.getMethodName() + " starting server for " + port);
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(port);
serverSocket.set(server);
latch.countDown();
@@ -200,6 +207,7 @@ public class TcpNioConnectionTests {
public void testMemoryLeak() throws Exception {
final int port = SocketUtils.findAvailableServerSocket();
TcpNioClientConnectionFactory factory = new TcpNioClientConnectionFactory("localhost", port);
factory.setApplicationEventPublisher(nullPublisher);
factory.setNioHarvestInterval(100);
factory.start();
final CountDownLatch latch = new CountDownLatch(1);
@@ -208,6 +216,7 @@ public class TcpNioConnectionTests {
@Override
public void run() {
try {
logger.debug(testName.getMethodName() + " starting server for " + port);
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(port);
serverSocket.set(server);
latch.countDown();
@@ -248,6 +257,7 @@ public class TcpNioConnectionTests {
@Test
public void testCleanup() throws Exception {
TcpNioClientConnectionFactory factory = new TcpNioClientConnectionFactory("localhost", 0);
factory.setApplicationEventPublisher(nullPublisher);
factory.setNioHarvestInterval(100);
Map<SocketChannel, TcpNioConnection> connections = new HashMap<SocketChannel, TcpNioConnection>();
SocketChannel chan1 = mock(SocketChannel.class);
@@ -327,7 +337,7 @@ public class TcpNioConnectionTests {
}
}).when(channel).read(Mockito.any(ByteBuffer.class));
when(socket.getReceiveBufferSize()).thenReturn(1024);
final TcpNioConnection connection = new TcpNioConnection(channel, false, false, null, null);
final TcpNioConnection connection = new TcpNioConnection(channel, false, false, nullPublisher, null);
connection.setTaskExecutor(exec);
connection.setPipeTimeout(200);
Method method = TcpNioConnection.class.getDeclaredMethod("doRead");
@@ -578,7 +588,7 @@ public class TcpNioConnectionTests {
public void testAssemblerUsesSecondaryExecutor() throws Exception {
final int port = SocketUtils.findAvailableServerSocket();
TcpNioServerConnectionFactory factory = new TcpNioServerConnectionFactory(port);
factory.setApplicationEventPublisher(mock(ApplicationEventPublisher.class));
factory.setApplicationEventPublisher(nullPublisher);
CompositeExecutor compositeExec = compositeExecutor();
@@ -625,7 +635,7 @@ public class TcpNioConnectionTests {
final int numberOfSockets = 100;
final int port = SocketUtils.findAvailableServerSocket();
TcpNioServerConnectionFactory factory = new TcpNioServerConnectionFactory(port);
factory.setApplicationEventPublisher(mock(ApplicationEventPublisher.class));
factory.setApplicationEventPublisher(nullPublisher);
CompositeExecutor compositeExec = compositeExecutor();