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

dominikriemer pushed a commit to branch avoid-connection-cache-immediate-remove
in repository https://gitbox.apache.org/repos/asf/streampipes.git

commit bdaa8043f531c5392f9dab8c4e1b7a0ebba00840
Author: Dominik Riemer <[email protected]>
AuthorDate: Sun Jun 28 21:48:04 2026 +0200

    fix: PLC4X cache recovery after failed read execution
---
 .../connection/ContinuousPlcRequestReader.java     |  10 +-
 .../plc/cache/SpLeasedPlcConnection.java           |  26 +++--
 .../plc/adapter/ConnectionContainerReproTest.java  | 108 +++++++++++++++++++++
 3 files changed, 129 insertions(+), 15 deletions(-)

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 e5e1a59fea..f32d9dd858 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
@@ -22,7 +22,6 @@ import 
org.apache.streampipes.extensions.api.connect.IEventCollector;
 import org.apache.streampipes.extensions.api.connect.IPollingSettings;
 import org.apache.streampipes.extensions.api.connect.IPullAdapter;
 import 
org.apache.streampipes.extensions.connectors.plc.adapter.generic.model.Plc4xConnectionSettings;
-import 
org.apache.streampipes.extensions.connectors.plc.cache.SpCachedPlcConnectionManager;
 import 
org.apache.streampipes.extensions.management.connect.adapter.util.PollingSettings;
 
 import org.apache.plc4x.java.api.PlcConnection;
@@ -86,16 +85,11 @@ public class ContinuousPlcRequestReader
         handleFailingPlcRead();
       }
     } catch (Exception e) {
-      handleFailingPlcReadAndRemoveFromCache(e.getMessage());
+      handleFailingPlcRead(e.getMessage());
     }
   }
 
-  private void handleFailingPlcReadAndRemoveFromCache(String problem) {
-    // ensure that the cached connection manager removes the broken connection
-    if (connectionManager instanceof SpCachedPlcConnectionManager) {
-      ((SpCachedPlcConnectionManager) 
connectionManager).removeCachedConnection(settings.connectionString());
-    }
-
+  private void handleFailingPlcRead(String problem) {
     LOG.error(
         "Error while reading from PLC with connection string {}, adapter {}. 
Setting adapter to idle for {} attempts. {} ",
         settings.connectionString(), adapterName, idlePullsBeforeNextAttempt, 
problem
diff --git 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpLeasedPlcConnection.java
 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpLeasedPlcConnection.java
index c008518a2b..eabc2df7c0 100644
--- 
a/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpLeasedPlcConnection.java
+++ 
b/streampipes-extensions/streampipes-connectors-plc/src/main/java/org/apache/streampipes/extensions/connectors/plc/cache/SpLeasedPlcConnection.java
@@ -57,6 +57,7 @@ import java.util.Optional;
 import java.util.Timer;
 import java.util.TimerTask;
 import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.CompletionException;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Consumer;
@@ -66,7 +67,7 @@ public class SpLeasedPlcConnection implements 
EventPlcConnection {
   private static final Logger log = 
LoggerFactory.getLogger(SpLeasedPlcConnection.class);
   private final SpConnectionContainer connectionContainer;
   private final AtomicReference<PlcConnection> connection;
-  private boolean invalidateConnection;
+  private volatile boolean invalidateConnection;
   private final Timer usageTimer;
   private final Duration maxUseDuration;
 
@@ -176,17 +177,20 @@ public class SpLeasedPlcConnection implements 
EventPlcConnection {
         return new PlcReadRequest() {
           @Override
           public CompletableFuture<? extends PlcReadResponse> execute() {
-            CompletableFuture<? extends PlcReadResponse> future =
-                innerPlcReadRequest.execute().orTimeout(Math.min(1000, 
maxUseDuration.toMillis()), TimeUnit.MILLISECONDS);
+            CompletableFuture<? extends PlcReadResponse> future;
+            try {
+              future = innerPlcReadRequest.execute()
+                  .orTimeout(Math.min(1000, maxUseDuration.toMillis()), 
TimeUnit.MILLISECONDS);
+            } catch (RuntimeException e) {
+              invalidateConnection(e);
+              return CompletableFuture.failedFuture(e);
+            }
             final CompletableFuture<PlcReadResponse> responseFuture = new 
CompletableFuture<>();
             future.handle((plcReadResponse, throwable) -> {
               if (throwable == null) {
                 responseFuture.complete(plcReadResponse);
               } else {
-                // Mark the connection as invalid.
-                invalidateConnection = true;
-                log.debug("ReadRequest execution completed exceptionally 
invalidateConnection=true",
-                    throwable);
+                invalidateConnection(throwable);
                 responseFuture.completeExceptionally(throwable);
               }
               return null;
@@ -572,4 +576,12 @@ public class SpLeasedPlcConnection implements 
EventPlcConnection {
     connectionContainer.removeEventListener(listener);
   }
 
+  private void invalidateConnection(Throwable throwable) {
+    invalidateConnection = true;
+    Throwable cause = throwable instanceof CompletionException && 
throwable.getCause() != null
+        ? throwable.getCause()
+        : throwable;
+    log.debug("PLC request execution completed exceptionally, invalidating 
leased connection", cause);
+  }
+
 }
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 49decd93cd..cd6cd60276 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
@@ -29,16 +29,20 @@ import 
org.apache.plc4x.java.api.exceptions.PlcConnectionException;
 import org.apache.plc4x.java.api.messages.PlcBrowseRequest;
 import org.apache.plc4x.java.api.messages.PlcPingResponse;
 import org.apache.plc4x.java.api.messages.PlcReadRequest;
+import org.apache.plc4x.java.api.messages.PlcReadResponse;
 import org.apache.plc4x.java.api.messages.PlcSubscriptionRequest;
 import org.apache.plc4x.java.api.messages.PlcUnsubscriptionRequest;
 import org.apache.plc4x.java.api.messages.PlcWriteRequest;
 import org.apache.plc4x.java.api.metadata.PlcConnectionMetadata;
 import org.apache.plc4x.java.api.model.PlcTag;
+import org.apache.plc4x.java.api.types.PlcResponseCode;
 import org.apache.plc4x.java.api.value.PlcValue;
 import 
org.apache.plc4x.java.utils.cache.exceptions.PlcConnectionManagerClosedException;
 import org.junit.jupiter.api.Test;
 
 import java.time.Duration;
+import java.util.LinkedHashSet;
+import java.util.List;
 import java.util.Optional;
 import java.util.concurrent.CompletableFuture;
 import java.util.concurrent.CountDownLatch;
@@ -46,6 +50,7 @@ import java.util.concurrent.ExecutionException;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.Executors;
 import java.util.concurrent.Future;
+import java.util.concurrent.RejectedExecutionException;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicBoolean;
 import java.util.concurrent.atomic.AtomicInteger;
@@ -190,6 +195,65 @@ class ConnectionContainerReproTest {
     }
   }
 
+  static class RejectingReadConnection extends MutableConnection {
+
+    RejectingReadConnection() {
+      super(true);
+    }
+
+    @Override
+    public PlcReadRequest.Builder readRequestBuilder() {
+      return new PlcReadRequest.Builder() {
+        @Override
+        public PlcReadRequest build() {
+          return new PlcReadRequest() {
+            @Override
+            public CompletableFuture<? extends PlcReadResponse> execute() {
+              throw new RejectedExecutionException("transaction executor 
terminated");
+            }
+
+            @Override
+            public int getNumberOfTags() {
+              return 0;
+            }
+
+            @Override
+            public LinkedHashSet<String> getTagNames() {
+              return new LinkedHashSet<>();
+            }
+
+            @Override
+            public PlcResponseCode getTagResponseCode(String tagName) {
+              return null;
+            }
+
+            @Override
+            public PlcTag getTag(String name) {
+              return null;
+            }
+
+            @Override
+            public List<PlcTag> getTags() {
+              return List.of();
+            }
+          };
+        }
+
+        @Override
+        public PlcReadRequest.Builder addTagAddress(String name,
+                                                    String tagAddress) {
+          return this;
+        }
+
+        @Override
+        public PlcReadRequest.Builder addTag(String name,
+                                             PlcTag tag) {
+          return this;
+        }
+      };
+    }
+  }
+
   @Test
   void recoversAfterFailedReconnectAndServesNewLeases() throws Exception {
     FlakyManager mgr = new FlakyManager();
@@ -304,6 +368,50 @@ class ConnectionContainerReproTest {
     connectionContainer.close();
   }
 
+  @Test
+  void invalidatesLeaseWhenReadExecutionIsRejectedSynchronously() throws 
Exception {
+    var firstConnection = new RejectingReadConnection();
+    var managerCalls = new AtomicInteger();
+    PlcConnectionManager manager = new PlcConnectionManager() {
+      @Override
+      public PlcConnection getConnection(String url) {
+        return managerCalls.incrementAndGet() == 1 ? firstConnection : new 
DummyConnection();
+      }
+
+      @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);
+
+    var readRequest = firstLease.readRequestBuilder().build();
+    var exception = assertThrows(
+        ExecutionException.class,
+        () -> readRequest.execute().get(500, TimeUnit.MILLISECONDS)
+    );
+    assertTrue(exception.getCause() instanceof RejectedExecutionException);
+
+    firstLease.close();
+
+    PlcConnection secondLease = connectionContainer.lease().get(500, 
TimeUnit.MILLISECONDS);
+    assertNotNull(secondLease);
+    assertEquals(2, managerCalls.get());
+    assertEquals(1, firstConnection.closeCalls());
+    secondLease.close();
+    connectionContainer.close();
+  }
+
   @Test
   void removingSlowConnectionDoesNotBlockLeasesForOtherUrls() throws Exception 
{
     var closeStarted = new CountDownLatch(1);

Reply via email to