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);

Reply via email to