From 3fff33c99b89f91ac08a52ab3f619676a4b131ae Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Tue, 29 Nov 2011 11:33:02 -0500 Subject: [PATCH] INT-2268 Fix Spurious Failing Test Occasional test failures caused because test HelloWorldInterceptor could prematurely close the connection, depending on timing. Also added some minor cleanup to the tests themselves. --- .../ip/tcp/TcpSendingMessageHandlerTests.java | 68 +++++++++---------- .../tcp/connection/HelloWorldInterceptor.java | 36 +++++----- 2 files changed, 54 insertions(+), 50 deletions(-) diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java index 23528c5125..7a55e32e4e 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/TcpSendingMessageHandlerTests.java @@ -943,25 +943,25 @@ public class TcpSendingMessageHandlerTests { latch.countDown(); Socket socket = server.accept(); int i = 0; - while (true) { - ObjectInputStream ois = new ObjectInputStream(socket.getInputStream()); - Object in = null; - ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream()); - if (i == 0) { - in = ois.readObject(); - logger.debug("read object: " + in); - oos.writeObject("world!"); - ois = new ObjectInputStream(socket.getInputStream()); - oos = new ObjectOutputStream(socket.getOutputStream()); - in = ois.readObject(); - logger.debug("read object: " + in); - oos.writeObject("world!"); - ois = new ObjectInputStream(socket.getInputStream()); - oos = new ObjectOutputStream(socket.getOutputStream()); - } + ObjectInputStream ois = new ObjectInputStream(socket.getInputStream()); + Object in = null; + ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream()); + if (i == 0) { in = ois.readObject(); - oos.writeObject("Reply" + (++i)); + logger.debug("read object: " + in); + oos.writeObject("world!"); + ois = new ObjectInputStream(socket.getInputStream()); + oos = new ObjectOutputStream(socket.getOutputStream()); + in = ois.readObject(); + logger.debug("read object: " + in); + oos.writeObject("world!"); + ois = new ObjectInputStream(socket.getInputStream()); + oos = new ObjectOutputStream(socket.getOutputStream()); } + in = ois.readObject(); + oos.writeObject("Reply" + (++i)); + socket.close(); + server.close(); } catch (Exception e) { if (!done.get()) { @@ -1000,25 +1000,25 @@ public class TcpSendingMessageHandlerTests { ServerSocket server = ServerSocketFactory.getDefault().createServerSocket(port); latch.countDown(); Socket socket = server.accept(); - while (true) { - ObjectInputStream ois = new ObjectInputStream(socket.getInputStream()); - Object in = null; - ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream()); - if (i == 0) { - in = ois.readObject(); - logger.debug("read object: " + in); - oos.writeObject("world!"); - ois = new ObjectInputStream(socket.getInputStream()); - oos = new ObjectOutputStream(socket.getOutputStream()); - in = ois.readObject(); - logger.debug("read object: " + in); - oos.writeObject("world!"); - ois = new ObjectInputStream(socket.getInputStream()); - oos = new ObjectOutputStream(socket.getOutputStream()); - } + ObjectInputStream ois = new ObjectInputStream(socket.getInputStream()); + Object in = null; + ObjectOutputStream oos = new ObjectOutputStream(socket.getOutputStream()); + if (i == 0) { in = ois.readObject(); - oos.writeObject("Reply" + (++i)); + logger.debug("read object: " + in); + oos.writeObject("world!"); + ois = new ObjectInputStream(socket.getInputStream()); + oos = new ObjectOutputStream(socket.getOutputStream()); + in = ois.readObject(); + logger.debug("read object: " + in); + oos.writeObject("world!"); + ois = new ObjectInputStream(socket.getInputStream()); + oos = new ObjectOutputStream(socket.getOutputStream()); } + in = ois.readObject(); + oos.writeObject("Reply" + (++i)); + socket.close(); + server.close(); } catch (Exception e) { if (i == 0) { e.printStackTrace(); diff --git a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/HelloWorldInterceptor.java b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/HelloWorldInterceptor.java index 33592964f3..ea8250526a 100644 --- a/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/HelloWorldInterceptor.java +++ b/spring-integration-ip/src/test/java/org/springframework/integration/ip/tcp/connection/HelloWorldInterceptor.java @@ -32,18 +32,20 @@ import org.springframework.integration.support.MessageBuilder; public class HelloWorldInterceptor extends AbstractTcpConnectionInterceptor { Log logger = LogFactory.getLog(this.getClass()); - - private boolean negotiated; - - private Semaphore negotiationSemaphore = new Semaphore(0); - - private long timeout = 10000; - - private String hello = "Hello"; - private String world = "world!"; - - private boolean closeReceived; + private volatile boolean negotiated; + + private final Semaphore negotiationSemaphore = new Semaphore(0); + + private volatile long timeout = 10000; + + private volatile String hello = "Hello"; + + private volatile String world = "world!"; + + private volatile boolean closeReceived; + + private volatile boolean pendingSend; public HelloWorldInterceptor() { } @@ -65,7 +67,7 @@ public class HelloWorldInterceptor extends AbstractTcpConnectionInterceptor { if (this.isServer()) { if (payload.equals(hello)) { try { - logger.debug("sending " + this.world); + logger.debug(this.toString() + " sending " + this.world); super.send(MessageBuilder.withPayload(world).build()); this.negotiated = true; return true; @@ -77,7 +79,7 @@ public class HelloWorldInterceptor extends AbstractTcpConnectionInterceptor { "' received '" + payload + "'"); } } else { - logger.debug("received " + payload); + logger.debug(this.toString() + " received " + payload); if (payload.equals(world)) { this.negotiated = true; this.negotiationSemaphore.release(); @@ -92,7 +94,7 @@ public class HelloWorldInterceptor extends AbstractTcpConnectionInterceptor { return super.onMessage(message); } finally { // on the server side, we don't want to close if we are expecting a response - if (!(this.isServer() && this.hasRealSender())) { + if (!(this.isServer() && this.hasRealSender()) && !this.pendingSend) { this.checkDeferredClose(); } } @@ -100,10 +102,11 @@ public class HelloWorldInterceptor extends AbstractTcpConnectionInterceptor { @Override public void send(Message message) throws Exception { + this.pendingSend = true; try { if (!this.negotiated) { if (!this.isServer()) { - logger.debug("Sending " + hello); + logger.debug(this.toString() + " Sending " + hello); super.send(MessageBuilder.withPayload(hello).build()); this.negotiationSemaphore.tryAcquire(this.timeout, TimeUnit.MILLISECONDS); if (!this.negotiated) { @@ -113,6 +116,7 @@ public class HelloWorldInterceptor extends AbstractTcpConnectionInterceptor { } super.send(message); } finally { + this.pendingSend = false; this.checkDeferredClose(); } } @@ -122,7 +126,7 @@ public class HelloWorldInterceptor extends AbstractTcpConnectionInterceptor { */ @Override public void close() { - if (this.negotiated) { + if (this.negotiated && !this.pendingSend) { super.close(); return; }