Revert "GH-3509 Register TcpSenders on wrapped connection"
This reverts commit 78a0ae8a04.
This commit is contained in:
@@ -1,5 +1,5 @@
|
||||
/*
|
||||
* Copyright 2013-2021 the original author or authors.
|
||||
* Copyright 2013-2019 the original author or authors.
|
||||
*
|
||||
* Licensed under the Apache License, Version 2.0 (the "License");
|
||||
* you may not use this file except in compliance with the License.
|
||||
@@ -16,29 +16,34 @@
|
||||
|
||||
package org.springframework.integration.ip.tcp;
|
||||
|
||||
import org.springframework.context.ApplicationEvent;
|
||||
import org.springframework.context.ApplicationEventPublisher;
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.HelloWorldInterceptorFactory;
|
||||
|
||||
/**
|
||||
* @author Gary Russell
|
||||
* @author Mário Dias
|
||||
* @author Artem Bilan
|
||||
*
|
||||
* @since 3.0
|
||||
*
|
||||
*/
|
||||
public class AbstractTcpChannelAdapterTests {
|
||||
|
||||
private static final ApplicationEventPublisher NOOP_PUBLISHER = event -> { };
|
||||
private static final ApplicationEventPublisher NOOP_PUBLISHER = new ApplicationEventPublisher() {
|
||||
|
||||
@Override
|
||||
public void publishEvent(ApplicationEvent event) {
|
||||
}
|
||||
|
||||
@Override
|
||||
public void publishEvent(Object event) {
|
||||
|
||||
}
|
||||
|
||||
};
|
||||
|
||||
protected HelloWorldInterceptorFactory newInterceptorFactory() {
|
||||
return newInterceptorFactory(NOOP_PUBLISHER);
|
||||
}
|
||||
|
||||
protected HelloWorldInterceptorFactory newInterceptorFactory(ApplicationEventPublisher applicationEventPublisher) {
|
||||
HelloWorldInterceptorFactory factory = new HelloWorldInterceptorFactory();
|
||||
factory.setApplicationEventPublisher(applicationEventPublisher);
|
||||
factory.setApplicationEventPublisher(NOOP_PUBLISHER);
|
||||
return factory;
|
||||
}
|
||||
|
||||
@@ -46,4 +51,5 @@ public class AbstractTcpChannelAdapterTests {
|
||||
connectionFactory.setApplicationEventPublisher(NOOP_PUBLISHER);
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -39,11 +39,10 @@ import java.util.concurrent.atomic.AtomicBoolean;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import javax.net.ServerSocketFactory;
|
||||
import javax.net.SocketFactory;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.Test;
|
||||
import org.mockito.Mockito;
|
||||
|
||||
import org.springframework.beans.factory.BeanFactory;
|
||||
@@ -58,11 +57,9 @@ import org.springframework.integration.config.ConsumerEndpointFactoryBean;
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractClientConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.AbstractServerConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpConnectionCloseEvent;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpConnectionInterceptorFactoryChain;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpNetClientConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpNetServerConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.connection.TcpNioClientConnectionFactory;
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayCrLfSerializer;
|
||||
import org.springframework.integration.ip.tcp.serializer.ByteArrayLengthHeaderSerializer;
|
||||
@@ -80,7 +77,6 @@ import org.springframework.scheduling.concurrent.ThreadPoolTaskScheduler;
|
||||
/**
|
||||
* @author Gary Russell
|
||||
* @author Artem Bilan
|
||||
* @author Mario Dias
|
||||
*
|
||||
* @since 2.0
|
||||
*/
|
||||
@@ -98,7 +94,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetCrLf() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -122,8 +118,8 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
}
|
||||
});
|
||||
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
|
||||
AbstractConnectionFactory ccf =
|
||||
new TcpNetClientConnectionFactory("localhost", serverSocket.get().getLocalPort());
|
||||
AbstractConnectionFactory ccf = new TcpNetClientConnectionFactory("localhost",
|
||||
serverSocket.get().getLocalPort());
|
||||
noopPublisher(ccf);
|
||||
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
|
||||
ccf.setSerializer(serializer);
|
||||
@@ -151,7 +147,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetCrLfClientMode() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -217,7 +213,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioCrLf() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -257,7 +253,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
adapter.setOutputChannel(channel);
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
Set<String> results = new HashSet<>();
|
||||
Set<String> results = new HashSet<String>();
|
||||
Message<?> mOut = channel.receive(10000);
|
||||
assertThat(mOut).isNotNull();
|
||||
results.add(new String((byte[]) mOut.getPayload()));
|
||||
@@ -273,7 +269,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetStxEtx() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -326,7 +322,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioStxEtx() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -366,7 +362,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
adapter.setOutputChannel(channel);
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
Set<String> results = new HashSet<>();
|
||||
Set<String> results = new HashSet<String>();
|
||||
Message<?> mOut = channel.receive(10000);
|
||||
assertThat(mOut).isNotNull();
|
||||
results.add(new String((byte[]) mOut.getPayload()));
|
||||
@@ -438,7 +434,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioLength() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -481,7 +477,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
adapter.setOutputChannel(channel);
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
Set<String> results = new HashSet<>();
|
||||
Set<String> results = new HashSet<String>();
|
||||
Message<?> mOut = channel.receive(10000);
|
||||
assertThat(mOut).isNotNull();
|
||||
results.add(new String((byte[]) mOut.getPayload()));
|
||||
@@ -497,7 +493,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetSerial() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -588,7 +584,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
adapter.setOutputChannel(channel);
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
Set<String> results = new HashSet<>();
|
||||
Set<String> results = new HashSet<String>();
|
||||
Message<?> mOut = channel.receive(10000);
|
||||
assertThat(mOut).isNotNull();
|
||||
results.add((String) mOut.getPayload());
|
||||
@@ -604,7 +600,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetSingleUseNoInbound() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final Semaphore semaphore = new Semaphore(0);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
@@ -651,7 +647,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioSingleUseNoInbound() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final Semaphore semaphore = new Semaphore(0);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
@@ -698,7 +694,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetSingleUseWithInbound() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final Semaphore semaphore = new Semaphore(0);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
@@ -743,7 +739,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
assertThat(semaphore.tryAcquire(2, 10000, TimeUnit.MILLISECONDS)).isTrue();
|
||||
Set<String> replies = new HashSet<>();
|
||||
Set<String> replies = new HashSet<String>();
|
||||
for (int i = 0; i < 2; i++) {
|
||||
Message<?> mOut = channel.receive(10000);
|
||||
assertThat(mOut).isNotNull();
|
||||
@@ -758,7 +754,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioSingleUseWithInbound() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final Semaphore semaphore = new Semaphore(0);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
@@ -803,7 +799,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
assertThat(semaphore.tryAcquire(2, 10000, TimeUnit.MILLISECONDS)).isTrue();
|
||||
Set<String> replies = new HashSet<>();
|
||||
Set<String> replies = new HashSet<String>();
|
||||
for (int i = 0; i < 2; i++) {
|
||||
Message<?> mOut = channel.receive(10000);
|
||||
assertThat(mOut).isNotNull();
|
||||
@@ -818,11 +814,11 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioSingleUseWithInboundMany() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final Semaphore semaphore = new Semaphore(0);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
final List<Socket> serverSockets = new ArrayList<>();
|
||||
final List<Socket> serverSockets = new ArrayList<Socket>();
|
||||
this.executor.execute(() -> {
|
||||
try {
|
||||
ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(0, 100);
|
||||
@@ -904,7 +900,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetNegotiate() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -916,7 +912,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
int i = 0;
|
||||
while (true) {
|
||||
ObjectInputStream ois = new ObjectInputStream(socket.getInputStream());
|
||||
Object in;
|
||||
Object in = null;
|
||||
ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream());
|
||||
if (i == 0) {
|
||||
in = ois.readObject();
|
||||
@@ -948,7 +944,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
ccf.setDeserializer(new DefaultDeserializer());
|
||||
ccf.setSoTimeout(10000);
|
||||
TcpConnectionInterceptorFactoryChain fc = new TcpConnectionInterceptorFactoryChain();
|
||||
fc.setInterceptors(new TcpConnectionInterceptorFactory[]{
|
||||
fc.setInterceptors(new TcpConnectionInterceptorFactory[] {
|
||||
newInterceptorFactory(),
|
||||
newInterceptorFactory()
|
||||
});
|
||||
@@ -975,7 +971,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioNegotiate() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -1015,7 +1011,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
ccf.setDeserializer(new DefaultDeserializer());
|
||||
ccf.setSoTimeout(10000);
|
||||
TcpConnectionInterceptorFactoryChain fc = new TcpConnectionInterceptorFactoryChain();
|
||||
fc.setInterceptors(new TcpConnectionInterceptorFactory[]{ newInterceptorFactory() });
|
||||
fc.setInterceptors(new TcpConnectionInterceptorFactory[] { newInterceptorFactory() });
|
||||
ccf.setInterceptorFactoryChain(fc);
|
||||
ccf.start();
|
||||
TcpSendingMessageHandler handler = new TcpSendingMessageHandler();
|
||||
@@ -1027,7 +1023,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
for (int i = 0; i < 1000; i++) {
|
||||
handler.handleMessage(MessageBuilder.withPayload("Test").build());
|
||||
}
|
||||
Set<String> results = new TreeSet<>();
|
||||
Set<String> results = new TreeSet<String>();
|
||||
for (int i = 0; i < 1000; i++) {
|
||||
Message<?> mOut = channel.receive(10000);
|
||||
assertThat(mOut).isNotNull();
|
||||
@@ -1044,7 +1040,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNetNegotiateSingleNoListen() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -1084,7 +1080,10 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
ccf.setDeserializer(new DefaultDeserializer());
|
||||
ccf.setSoTimeout(10000);
|
||||
TcpConnectionInterceptorFactoryChain fc = new TcpConnectionInterceptorFactoryChain();
|
||||
fc.setInterceptor(newInterceptorFactory(), newInterceptorFactory());
|
||||
fc.setInterceptors(new TcpConnectionInterceptorFactory[] {
|
||||
newInterceptorFactory(),
|
||||
newInterceptorFactory()
|
||||
});
|
||||
ccf.setInterceptorFactoryChain(fc);
|
||||
ccf.setSingleUse(true);
|
||||
ccf.start();
|
||||
@@ -1098,7 +1097,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
|
||||
@Test
|
||||
public void testNioNegotiateSingleNoListen() throws Exception {
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<>();
|
||||
final AtomicReference<ServerSocket> serverSocket = new AtomicReference<ServerSocket>();
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
final AtomicBoolean done = new AtomicBoolean();
|
||||
this.executor.execute(() -> {
|
||||
@@ -1138,7 +1137,10 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
ccf.setDeserializer(new DefaultDeserializer());
|
||||
ccf.setSoTimeout(10000);
|
||||
TcpConnectionInterceptorFactoryChain fc = new TcpConnectionInterceptorFactoryChain();
|
||||
fc.setInterceptor(newInterceptorFactory(), newInterceptorFactory());
|
||||
fc.setInterceptors(new TcpConnectionInterceptorFactory[] {
|
||||
newInterceptorFactory(),
|
||||
newInterceptorFactory()
|
||||
});
|
||||
ccf.setInterceptorFactoryChain(fc);
|
||||
ccf.setSingleUse(true);
|
||||
ccf.start();
|
||||
@@ -1151,18 +1153,18 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testOutboundChannelAdapterWithinChain() {
|
||||
public void testOutboundChannelAdapterWithinChain() throws Exception {
|
||||
AbstractApplicationContext ctx = new ClassPathXmlApplicationContext(
|
||||
"TcpOutboundChannelAdapterWithinChainTests-context.xml", this.getClass());
|
||||
AbstractServerConnectionFactory scf = ctx.getBean(AbstractServerConnectionFactory.class);
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
ctx.getBean(AbstractClientConnectionFactory.class).setPort(scf.getPort());
|
||||
ctx.getBeansOfType(ConsumerEndpointFactoryBean.class).values().forEach(ConsumerEndpointFactoryBean::start);
|
||||
ctx.getBeansOfType(ConsumerEndpointFactoryBean.class).values().forEach(c -> c.start());
|
||||
MessageChannel channelAdapterWithinChain = ctx.getBean("tcpOutboundChannelAdapterWithinChain",
|
||||
MessageChannel.class);
|
||||
PollableChannel inbound = ctx.getBean("inbound", PollableChannel.class);
|
||||
String testPayload = "Hello, world!";
|
||||
channelAdapterWithinChain.send(new GenericMessage<>(testPayload));
|
||||
channelAdapterWithinChain.send(new GenericMessage<String>(testPayload));
|
||||
Message<?> m = inbound.receive(1000);
|
||||
assertThat(m).isNotNull();
|
||||
assertThat(new String((byte[]) m.getPayload())).isEqualTo(testPayload);
|
||||
@@ -1178,7 +1180,7 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
}).when(mockCcf).getConnection();
|
||||
handler.setConnectionFactory(mockCcf);
|
||||
try {
|
||||
handler.handleMessage(new GenericMessage<>("foo"));
|
||||
handler.handleMessage(new GenericMessage<String>("foo"));
|
||||
fail("Expected exception");
|
||||
}
|
||||
catch (Exception e) {
|
||||
@@ -1189,33 +1191,4 @@ public class TcpSendingMessageHandlerTests extends AbstractTcpChannelAdapterTest
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testInterceptedCleanup() throws Exception {
|
||||
final CountDownLatch latch = new CountDownLatch(1);
|
||||
AbstractServerConnectionFactory scf = new TcpNetServerConnectionFactory(0);
|
||||
ByteArrayCrLfSerializer serializer = new ByteArrayCrLfSerializer();
|
||||
scf.setSerializer(serializer);
|
||||
scf.setDeserializer(serializer);
|
||||
TcpReceivingChannelAdapter adapter = new TcpReceivingChannelAdapter();
|
||||
adapter.setConnectionFactory(scf);
|
||||
TcpSendingMessageHandler handler = new TcpSendingMessageHandler();
|
||||
handler.setConnectionFactory(scf);
|
||||
scf.setApplicationEventPublisher(event -> {
|
||||
if (event instanceof TcpConnectionCloseEvent) {
|
||||
latch.countDown();
|
||||
}
|
||||
});
|
||||
TcpConnectionInterceptorFactoryChain fc = new TcpConnectionInterceptorFactoryChain();
|
||||
fc.setInterceptor(newInterceptorFactory(scf.getApplicationEventPublisher()));
|
||||
scf.setInterceptorFactoryChain(fc);
|
||||
scf.start();
|
||||
TestingUtilities.waitListening(scf, null);
|
||||
int port = scf.getPort();
|
||||
Socket socket = SocketFactory.getDefault().createSocket("localhost", port);
|
||||
socket.close();
|
||||
assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue();
|
||||
assertThat(handler.getConnections().isEmpty()).isTrue();
|
||||
scf.stop();
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user