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