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