From 1d5ccb7a0e61f1fae963d84aae75b3b65509a9ad Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 23 Aug 2010 18:40:25 +0000 Subject: [PATCH] INT-1367 Fix Deadlock on NIO Reading Thread --- .../ip/tcp/connection/TcpNioConnection.java | 33 ++++++++++++------- 1 file changed, 21 insertions(+), 12 deletions(-) diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpNioConnection.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpNioConnection.java index 0cde646120..8e7a582142 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpNioConnection.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpNioConnection.java @@ -63,6 +63,8 @@ public class TcpNioConnection extends AbstractTcpConnection { private long lastRead; private AtomicInteger executionControl = new AtomicInteger(); + + private boolean writingToPipe; /** * Constructs a TcpNetConnection for the SocketChannel. @@ -156,8 +158,7 @@ public class TcpNioConnection extends AbstractTcpConnection { } while (active) { try { - while (this.socketChannel.isOpen() && - this.pipedInputStream.available() > 0) { + while (dataAvailable()) { convertAndSend(); } } catch (IOException e) { @@ -174,8 +175,13 @@ public class TcpNioConnection extends AbstractTcpConnection { } } + private boolean dataAvailable() throws IOException { + return this.socketChannel.isOpen() && + (this.pipedInputStream.available() > 0 || writingToPipe); + } + private synchronized void convertAndSend() throws IOException { - if (!this.socketChannel.isOpen() || this.pipedInputStream.available() <= 0) { + if (!dataAvailable()) { return; } Message message = null; @@ -231,6 +237,16 @@ public class TcpNioConnection extends AbstractTcpConnection { if (rawBuffer == null) { rawBuffer = allocate(maxMessageSize); } + + writingToPipe = true; + if (this.taskExecutor == null) { + this.taskExecutor = Executors.newSingleThreadExecutor(); + } + if (this.executionControl.incrementAndGet() <= 1) { + // only execute run() if we don't already have one running + this.executionControl.set(1); + this.taskExecutor.execute(this); + } rawBuffer.clear(); int len = socketChannel.read(rawBuffer); if (len < 0) { @@ -243,15 +259,8 @@ public class TcpNioConnection extends AbstractTcpConnection { } pipedOutputStream.write(rawBuffer.array(), 0, rawBuffer.limit()); pipedOutputStream.flush(); - - if (this.taskExecutor == null) { - this.taskExecutor = Executors.newSingleThreadExecutor(); - } - if (this.executionControl.incrementAndGet() <= 1) { - // only execute run() if we don't already have one running - this.executionControl.set(1); - this.taskExecutor.execute(this); - } + writingToPipe = false; + } /**