From ca804359512e0a3deff87c6ae2b4792680d6d950 Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Thu, 23 Apr 2020 17:48:52 -0400 Subject: [PATCH] GH-3256: Fix interaction with delayed listener reg - Use the test listener before waiting for delayed listener registration. - Don't wait for delayed listener registration if the test fails. --- .../AbstractClientConnectionFactory.java | 1 + .../ip/tcp/connection/TcpConnectionSupport.java | 16 +++++++++++----- 2 files changed, 12 insertions(+), 5 deletions(-) diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractClientConnectionFactory.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractClientConnectionFactory.java index 00fa032e7a..1a220a9bab 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractClientConnectionFactory.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/AbstractClientConnectionFactory.java @@ -183,6 +183,7 @@ public abstract class AbstractClientConnectionFactory extends AbstractConnection connection = buildNewConnection(); if (this.connectionTest != null && !this.connectionTest.test(connection)) { + connection.setTestFailed(true); connection.close(); throw new UncheckedIOException(new IOException("Connection test failed for " + connection)); } diff --git a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java index aecb068fd4..df4c134006 100644 --- a/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java +++ b/spring-integration-ip/src/main/java/org/springframework/integration/ip/tcp/connection/TcpConnectionSupport.java @@ -101,6 +101,8 @@ public abstract class TcpConnectionSupport implements TcpConnection { */ private boolean needsTest; + private volatile boolean testFailed; + public TcpConnectionSupport() { this(null); } @@ -151,6 +153,10 @@ public abstract class TcpConnectionSupport implements TcpConnection { } } + void setTestFailed(boolean testFailed) { + this.testFailed = testFailed; + } + /** * Closes this connection. */ @@ -305,16 +311,16 @@ public abstract class TcpConnectionSupport implements TcpConnection { */ @Override public TcpListener getListener() { - if (this.manualListenerRegistration) { + if (this.needsTest && this.testListener != null) { + this.needsTest = false; + return this.testListener; + } + if (this.manualListenerRegistration && !this.testFailed) { if (this.logger.isDebugEnabled()) { this.logger.debug(getConnectionId() + " Waiting for listener registration"); } waitForListenerRegistration(); } - if (this.needsTest && this.testListener != null) { - this.needsTest = false; - return this.testListener; - } return this.listener; }