This is an automated email from the ASF dual-hosted git repository.

dominikriemer pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/streampipes.git


The following commit(s) were added to refs/heads/dev by this push:
     new e7dafb922b fix: Improve retry of failed PLC connections (#4572)
e7dafb922b is described below

commit e7dafb922bb279520db9a96c570a3261586062b6
Author: Dominik Riemer <[email protected]>
AuthorDate: Fri Jun 26 13:39:11 2026 +0200

    fix: Improve retry of failed PLC connections (#4572)
---
 .../plc/adapter/generic/GenericPlc4xAdapter.java   |   9 +-
 .../connection/ContinuousPlcRequestReader.java     |  15 ++-
 .../connectors/plc/adapter/s7/Plc4xS7Adapter.java  |   9 +-
 .../plc/cache/SpConnectionContainer.java           |  72 +++++++++-----
 .../plc/adapter/ConnectionContainerReproTest.java  | 110 +++++++++++++++++++++
 5 files changed, 183 insertions(+), 32 deletions(-)

diff --git 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/GenericPlc4xAdapter.java
 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/GenericPlc4xAdapter.java
index ffd48c218a..351a6c86ae 100644
--- 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/GenericPlc4xAdapter.java
+++ 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/GenericPlc4xAdapter.java
@@ -79,7 +79,14 @@ public class GenericPlc4xAdapter implements 
StreamPipesAdapter, SupportsRuntimeC
     var settings = new Plc4xConnectionExtractor(
         extractor.getStaticPropertyExtractor(), driver.getProtocolCode()
     ).makeSettings();
-    var plcRequestReader = new ContinuousPlcRequestReader(connectionManager, 
settings, requestProvider, collector);
+    var adapterName = extractor.getAdapterDescription().getName();
+    var plcRequestReader = new ContinuousPlcRequestReader(
+        connectionManager,
+        settings,
+        requestProvider,
+        collector,
+        adapterName
+    );
     this.pullAdapterScheduler = new PullAdapterScheduler();
     this.pullAdapterScheduler.schedule(plcRequestReader, 
extractor.getAdapterDescription().getElementId());
   }
diff --git 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/connection/ContinuousPlcRequestReader.java
 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/connection/ContinuousPlcRequestReader.java
index 89c302eed7..e5e1a59fea 100644
--- 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/connection/ContinuousPlcRequestReader.java
+++ 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/generic/connection/ContinuousPlcRequestReader.java
@@ -42,6 +42,7 @@ public class ContinuousPlcRequestReader
   private final IEventCollector collector;
   private int idlePullsBeforeNextAttempt = 0;
   private int currentIdlePulls = 0;
+  private final String adapterName;
 
   /**
    *  Failure and recovery strategy:
@@ -52,10 +53,12 @@ public class ContinuousPlcRequestReader
       PlcConnectionManager connectionManager,
       Plc4xConnectionSettings settings,
       PlcRequestProvider requestProvider,
-      IEventCollector collector
+      IEventCollector collector,
+      String adapterName
   ) {
     super(connectionManager, settings, requestProvider);
     this.collector = collector;
+    this.adapterName = adapterName;
   }
 
   @Override
@@ -75,7 +78,11 @@ public class ContinuousPlcRequestReader
             .get(5000, TimeUnit.MILLISECONDS);
         processPlcReadResponse(readResponse);
       } else {
-        LOG.error("Not connected to PLC with connection string {}", 
settings.connectionString());
+        LOG.error(
+            "Not connected to PLC with connection string {}, adapter {}",
+            settings.connectionString(),
+            adapterName
+        );
         handleFailingPlcRead();
       }
     } catch (Exception e) {
@@ -90,8 +97,8 @@ public class ContinuousPlcRequestReader
     }
 
     LOG.error(
-        "Error while reading from PLC with connection string {}. Setting 
adapter to idle for {} attempts. {} ",
-        settings.connectionString(), idlePullsBeforeNextAttempt, problem
+        "Error while reading from PLC with connection string {}, adapter {}. 
Setting adapter to idle for {} attempts. {} ",
+        settings.connectionString(), adapterName, idlePullsBeforeNextAttempt, 
problem
     );
 
     handleFailingPlcRead();
diff --git 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/s7/Plc4xS7Adapter.java
 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/s7/Plc4xS7Adapter.java
index d371e837a2..f5666fab17 100644
--- 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/s7/Plc4xS7Adapter.java
+++ 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/adapter/s7/Plc4xS7Adapter.java
@@ -158,7 +158,14 @@ public class Plc4xS7Adapter implements StreamPipesAdapter {
       IAdapterRuntimeContext adapterRuntimeContext
   ) {
     var settings = getConfigurations(extractor.getStaticPropertyExtractor());
-    var plcRequestReader = new ContinuousPlcRequestReader(connectionManager, 
settings, requestProvider, collector);
+    var adapterName = extractor.getAdapterDescription().getName();
+    var plcRequestReader = new ContinuousPlcRequestReader(
+        connectionManager,
+        settings,
+        requestProvider,
+        collector,
+        adapterName
+    );
     this.pullAdapterScheduler = new PullAdapterScheduler();
     this.pullAdapterScheduler.schedule(
         plcRequestReader,
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..ec1f0371a7 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 {
@@ -164,27 +169,23 @@ public class SpConnectionContainer {
       return;
     }
 
-    // If something happened while using the connection, invalidate this one 
and create a new connection.
+    // If something happened while using the connection, invalidate it.
     if (invalidateConnection) {
       // Close the old connection.
-      try {
-        connection.close();
-      } catch (Exception e) {
-        // We're ignoring this as we have no idea, what state the connection 
is in.
-        // Nevertheless, it is polite to say something in logs about this 
situation.
-        LOGGER.warn("Exception while closing connection", e);
-      }
+      closeConnection(connection);
+      connection = null;
 
-      // Try to get a new connection.
-      try {
-        connection = connectionManager.getConnection(connectionUrl);
-      } catch (PlcConnectionException e) {
-        // If something goes wrong, close all waiting futures exceptionally.
-        LOGGER.warn("Can't get connection for {} complete queue items 
exceptionally", connectionUrl, e);
-        queue.forEach(future -> future.completeExceptionally(e));
-        queue.clear();
-        leasedConnection = null;
-        connection = null;
+      // Only reconnect immediately when another client is waiting for the 
connection.
+      if (!queue.isEmpty()) {
+        try {
+          connection = connectionManager.getConnection(connectionUrl);
+        } catch (PlcConnectionException e) {
+          // If something goes wrong, close all waiting futures exceptionally.
+          LOGGER.warn("Can't get connection for {} complete queue items 
exceptionally", connectionUrl, e);
+          queue.forEach(future -> future.completeExceptionally(e));
+          queue.clear();
+          leasedConnection = null;
+        }
       }
     }
 
@@ -192,14 +193,16 @@ public class SpConnectionContainer {
     if (queue.isEmpty()) {
       leasedConnection = null;
 
-      // Start a timer to invalidate this connection if it's idle for too long.
-      idleTimer = new Timer("CC-Idle-Timer-" + Thread.currentThread().getId());
-      idleTimer.schedule(new TimerTask() {
-        @Override
-        public void run() {
-          closeIdleConnection();
-        }
-      }, maxIdleTime.toMillis());
+      if (connection != null) {
+        // Start a timer to invalidate this connection if it's idle for too 
long.
+        idleTimer = new Timer("CC-Idle-Timer-" + 
Thread.currentThread().getId());
+        idleTimer.schedule(new TimerTask() {
+          @Override
+          public void run() {
+            closeIdleConnection();
+          }
+        }, maxIdleTime.toMillis());
+      }
       return;
     }
 
@@ -242,6 +245,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..49decd93cd 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,87 @@ 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 doesNotEagerlyReplaceInvalidConnectionWithoutWaitingClient() throws 
Exception {
+    var firstConnection = new MutableConnection(true);
+    var secondConnection = new MutableConnection(true);
+    var managerCalls = new AtomicInteger();
+    PlcConnectionManager manager = new PlcConnectionManager() {
+      @Override
+      public PlcConnection getConnection(String url) {
+        return managerCalls.incrementAndGet() == 1 ? firstConnection : 
secondConnection;
+      }
+
+      @Override
+      public PlcConnection getConnection(String url,
+                                         PlcAuthentication authentication) {
+        return null;
+      }
+    };
+    var 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, true);
+
+    assertEquals(1, managerCalls.get());
+    assertEquals(1, firstConnection.closeCalls());
+
+    SpLeasedPlcConnection secondLease =
+        (SpLeasedPlcConnection) connectionContainer.lease().get(500, 
TimeUnit.MILLISECONDS);
+    assertEquals(2, managerCalls.get());
+    connectionContainer.returnConnection(secondLease, false);
+    connectionContainer.close();
+  }
+
   @Test
   void removingSlowConnectionDoesNotBlockLeasesForOtherUrls() throws Exception 
{
     var closeStarted = new CountDownLatch(1);

Reply via email to