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 369a821f29 fix: PLC4X cache recovery after failed read execution
(#4648)
369a821f29 is described below
commit 369a821f29f811e5f708ae5bc7740805ea284f50
Author: Dominik Riemer <[email protected]>
AuthorDate: Thu Jul 16 17:22:29 2026 +0200
fix: PLC4X cache recovery after failed read execution (#4648)
Co-authored-by: Philipp Zehnder <[email protected]>
---
.../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);