This is an automated email from the ASF dual-hosted git repository. dominikriemer pushed a commit to branch fix-plc-retry-failed-connection in repository https://gitbox.apache.org/repos/asf/streampipes.git
commit e4d3ae27a02e7dcedff5afbe4c3df860585fe082 Author: Dominik Riemer <[email protected]> AuthorDate: Mon Jun 15 22:03:49 2026 +0200 fix: Improve retry of failed PLC connections --- .../plc/cache/SpConnectionContainer.java | 22 +++++++ .../plc/adapter/ConnectionContainerReproTest.java | 71 ++++++++++++++++++++++ 2 files changed, 93 insertions(+) diff --git a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpConnectionContainer.java b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpConnectionContainer.java index a3a329fc54..67eda0f222 100644 --- a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpConnectionContainer.java +++ b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpConnectionContainer.java @@ -109,6 +109,11 @@ public class SpConnectionContainer { return connectionFuture; } + if (connection != null && !isConnected(connection)) { + closeConnection(connection); + connection = null; + } + // Try to get a new connection, if we haven't got one yet. if (connection == null) { try { @@ -242,6 +247,23 @@ public class SpConnectionContainer { } } + private boolean isConnected(PlcConnection plcConnection) { + try { + return plcConnection.isConnected(); + } catch (Exception e) { + LOGGER.warn("Exception while checking connection state for {}", connectionUrl, e); + return false; + } + } + + private void closeConnection(PlcConnection plcConnection) { + try { + plcConnection.close(); + } catch (Exception e) { + LOGGER.warn("Exception while closing stale connection for {}", connectionUrl, e); + } + } + public void addEventListener(EventListener listener) { if ((connection != null) && (connection instanceof EventPlcConnection)) { diff --git a/streampipes-extensions/streampipes-connectors-plc/src/test/java/org/apache/streampipes/extensions/connectors/plc/adapter/ConnectionContainerReproTest.java b/streampipes-extensions/streampipes-connectors-plc/src/test/java/org/apache/streampipes/extensions/connectors/plc/adapter/ConnectionContainerReproTest.java index c2326fb121..a4573832b6 100644 --- a/streampipes-extensions/streampipes-connectors-plc/src/test/java/org/apache/streampipes/extensions/connectors/plc/adapter/ConnectionContainerReproTest.java +++ b/streampipes-extensions/streampipes-connectors-plc/src/test/java/org/apache/streampipes/extensions/connectors/plc/adapter/ConnectionContainerReproTest.java @@ -47,6 +47,7 @@ import java.util.concurrent.ExecutorService; import java.util.concurrent.Executors; import java.util.concurrent.Future; import java.util.concurrent.TimeUnit; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicInteger; import static org.junit.jupiter.api.Assertions.assertEquals; @@ -161,6 +162,34 @@ class ConnectionContainerReproTest { } } + static class MutableConnection extends DummyConnection { + private final AtomicBoolean connected; + private final AtomicInteger closeCalls; + + MutableConnection(boolean connected) { + this.connected = new AtomicBoolean(connected); + this.closeCalls = new AtomicInteger(); + } + + @Override + public boolean isConnected() { + return connected.get(); + } + + @Override + public void close() { + closeCalls.incrementAndGet(); + } + + void setConnected(boolean connected) { + this.connected.set(connected); + } + + int closeCalls() { + return closeCalls.get(); + } + } + @Test void recoversAfterFailedReconnectAndServesNewLeases() throws Exception { FlakyManager mgr = new FlakyManager(); @@ -194,6 +223,48 @@ class ConnectionContainerReproTest { assertNotNull(lease3); } + @Test + void replacesDisconnectedIdleConnection() throws Exception { + var staleConnection = new MutableConnection(true); + var managerCalls = new AtomicInteger(); + PlcConnectionManager manager = new PlcConnectionManager() { + @Override + public PlcConnection getConnection(String url) { + if (managerCalls.incrementAndGet() == 1) { + return staleConnection; + } + return new DummyConnection(); + } + + @Override + public PlcConnection getConnection(String url, + PlcAuthentication authentication) { + return null; + } + }; + SpConnectionContainer connectionContainer = new SpConnectionContainer( + manager, + "mock://plc", + Duration.ofSeconds(30), + Duration.ofSeconds(30), + url -> null + ); + + SpLeasedPlcConnection firstLease = + (SpLeasedPlcConnection) connectionContainer.lease().get(500, TimeUnit.MILLISECONDS); + connectionContainer.returnConnection(firstLease, false); + + staleConnection.setConnected(false); + + SpLeasedPlcConnection secondLease = + (SpLeasedPlcConnection) connectionContainer.lease().get(500, TimeUnit.MILLISECONDS); + assertNotNull(secondLease); + assertEquals(2, managerCalls.get()); + assertEquals(1, staleConnection.closeCalls()); + connectionContainer.returnConnection(secondLease, false); + connectionContainer.close(); + } + @Test void removingSlowConnectionDoesNotBlockLeasesForOtherUrls() throws Exception { var closeStarted = new CountDownLatch(1);
