This is an automated email from the ASF dual-hosted git repository.
davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/main by this push:
new 302d1e0bba2f CAMEL-25291: camel-plc4x - fail the exchange when a write
or read cannot be done (#27327)
302d1e0bba2f is described below
commit 302d1e0bba2fcef83d3d5e78af419770c12f103d
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:32:32 2026 +0530
CAMEL-25291: camel-plc4x - fail the exchange when a write or read cannot be
done (#27327)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../component/plc4x/Plc4XPollingConsumer.java | 19 ++-
.../camel/component/plc4x/Plc4XProducer.java | 17 +--
.../camel/component/plc4x/Plc4XFailureTest.java | 147 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 14 ++
4 files changed, 174 insertions(+), 23 deletions(-)
diff --git
a/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XPollingConsumer.java
b/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XPollingConsumer.java
index 8c0e10e1bfa1..6fcf7242244f 100644
---
a/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XPollingConsumer.java
+++
b/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XPollingConsumer.java
@@ -106,19 +106,18 @@ public class Plc4XPollingConsumer extends
EventDrivenPollingConsumer {
rsp.put(field, response.getObject(field));
}
exchange.getIn().setBody(rsp);
- } catch (ExecutionException | TimeoutException e) {
- getExceptionHandler().handleException(e);
- exchange.getIn().setBody(new HashMap<>());
+ } catch (TimeoutException e) {
+ // nothing received within the timeout
+ LOGGER.debug("No response from the PLC within {} millis", timeout);
+ return null;
+ } catch (ExecutionException e) {
+ // the read failed: fail the exchange instead of answering an
empty result
+ exchange.setException(e.getCause() != null ? e.getCause() : e);
} catch (InterruptedException e) {
- getExceptionHandler().handleException(e);
Thread.currentThread().interrupt();
+ exchange.setException(e);
} catch (PlcConnectionException e) {
- if (LOGGER.isTraceEnabled()) {
- LOGGER.warn("Unable to reconnect, skipping request", e);
- } else {
- LOGGER.warn("Unable to reconnect, skipping request");
- }
- exchange.getIn().setBody(new HashMap<>());
+ exchange.setException(e);
}
return exchange;
}
diff --git
a/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XProducer.java
b/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XProducer.java
index ac4afc54e406..6b82e831ead0 100644
---
a/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XProducer.java
+++
b/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XProducer.java
@@ -64,18 +64,10 @@ public class Plc4XProducer extends DefaultAsyncProducer {
@Override
public void process(Exchange exchange) throws Exception {
- try {
- plc4XEndpoint.reconnectIfNeeded();
- if (plc4XEndpoint.isConnected() && !plc4XEndpoint.canWrite()) {
- throw new PlcException("This connection (" +
plc4XEndpoint.getUri() + ") doesn't support writing.");
- }
- } catch (PlcConnectionException e) {
- if (log.isTraceEnabled()) {
- log.warn("Unable to reconnect, skipping request", e);
- } else {
- log.warn("Unable to reconnect, skipping request");
- }
- return;
+ // a failed reconnect fails the exchange: the write was not done
+ plc4XEndpoint.reconnectIfNeeded();
+ if (plc4XEndpoint.isConnected() && !plc4XEndpoint.canWrite()) {
+ throw new PlcException("This connection (" +
plc4XEndpoint.getUri() + ") doesn't support writing.");
}
Message in = exchange.getIn();
@@ -112,7 +104,6 @@ public class Plc4XProducer extends DefaultAsyncProducer {
Message out = exchange.getMessage();
out.copyFrom(exchange.getIn());
} catch (Exception e) {
- exchange.setMessage(null);
exchange.setException(e);
}
callback.done(true);
diff --git
a/components/camel-plc4x/src/test/java/org/apache/camel/component/plc4x/Plc4XFailureTest.java
b/components/camel-plc4x/src/test/java/org/apache/camel/component/plc4x/Plc4XFailureTest.java
new file mode 100644
index 000000000000..54ac03fe2253
--- /dev/null
+++
b/components/camel-plc4x/src/test/java/org/apache/camel/component/plc4x/Plc4XFailureTest.java
@@ -0,0 +1,147 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements. See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.component.plc4x;
+
+import java.util.Map;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.ExecutionException;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.plc4x.java.api.PlcConnection;
+import org.apache.plc4x.java.api.exceptions.PlcConnectionException;
+import org.apache.plc4x.java.api.exceptions.PlcRuntimeException;
+import org.apache.plc4x.java.api.messages.PlcReadRequest;
+import org.apache.plc4x.java.api.messages.PlcWriteRequest;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.mockito.Mockito.RETURNS_DEEP_STUBS;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * A write or read that could not be done must fail the exchange, not complete
it as if it had been done.
+ */
+class Plc4XFailureTest {
+
+ private final CamelContext context = new DefaultCamelContext();
+ private final Plc4XEndpoint endpoint = mock(Plc4XEndpoint.class,
RETURNS_DEEP_STUBS);
+
+ @BeforeEach
+ void setUp() {
+ context.start();
+
when(endpoint.getEndpointUri()).thenReturn("plc4x:mock:10.10.10.1/1/1");
+ when(endpoint.createExchange()).thenAnswer(invocation -> new
DefaultExchange(context));
+ }
+
+ @AfterEach
+ void tearDown() {
+ context.stop();
+ }
+
+ @Test
+ void testProducerFailsWhenReconnectFails() throws Exception {
+ PlcConnectionException cause = new PlcConnectionException("PLC
unreachable");
+ doThrow(cause).when(endpoint).reconnectIfNeeded();
+ Map<String, Map<String, Object>> tags = Map.of("test",
Map.of("testAddress", 1));
+ Exchange exchange = new DefaultExchange(context);
+ exchange.getIn().setBody(tags);
+
+ new Plc4XProducer(endpoint).process(exchange, doneSync -> {
+ });
+
+ assertSame(cause, exchange.getException());
+ // the message is kept, so that a redelivery writes it again
+ assertEquals(tags, exchange.getIn().getBody());
+ }
+
+ @Test
+ void testProducerKeepsTheMessageWhenTheWriteFails() throws Exception {
+ PlcWriteRequest request = mock(PlcWriteRequest.class);
+ when(request.execute()).thenAnswer(invocation ->
CompletableFuture.failedFuture(new PlcRuntimeException("denied")));
+ Map<String, Map<String, Object>> tags = Map.of("test",
Map.of("testAddress", 1));
+ when(endpoint.buildPlcWriteRequest(tags)).thenReturn(request);
+ Exchange exchange = new DefaultExchange(context);
+ exchange.getIn().setBody(tags);
+
+ new Plc4XProducer(endpoint).process(exchange, doneSync -> {
+ });
+
+ assertInstanceOf(ExecutionException.class, exchange.getException());
+ assertEquals(tags, exchange.getIn().getBody());
+ }
+
+ @Test
+ void testPollingConsumerFailsWhenReconnectFails() throws Exception {
+ PlcConnectionException cause = new PlcConnectionException("PLC
unreachable");
+ doThrow(cause).when(endpoint).reconnectIfNeeded();
+
+ Exchange exchange = new Plc4XPollingConsumer(endpoint).receive(1000);
+
+ assertNotNull(exchange);
+ assertSame(cause, exchange.getException());
+ }
+
+ @Test
+ void testPollingConsumerFailsWhenReadFails() throws Exception {
+ PlcRuntimeException cause = new PlcRuntimeException("read failed");
+ PlcReadRequest request = mock(PlcReadRequest.class);
+ when(request.execute()).thenAnswer(invocation ->
CompletableFuture.failedFuture(cause));
+ endpoint.connection = mock(PlcConnection.class);
+ when(endpoint.buildPlcReadRequest()).thenReturn(request);
+
+ Exchange exchange = new Plc4XPollingConsumer(endpoint).receive(1000);
+
+ assertNotNull(exchange);
+ assertSame(cause, exchange.getException());
+ }
+
+ @Test
+ void testPollingConsumerReturnsNullWithoutResponse() throws Exception {
+ PlcReadRequest request = mock(PlcReadRequest.class);
+ when(request.execute()).thenAnswer(invocation -> new
CompletableFuture<>());
+ endpoint.connection = mock(PlcConnection.class);
+ when(endpoint.buildPlcReadRequest()).thenReturn(request);
+
+ // receiveNoWait and receive(timeout) return null when nothing is
received in time
+ assertNull(new Plc4XPollingConsumer(endpoint).receiveNoWait());
+ assertNull(new Plc4XPollingConsumer(endpoint).receive(10));
+ }
+
+ @Test
+ void testPollingConsumerReadsTags() throws Exception {
+ endpoint.connection = mock(PlcConnection.class);
+ PlcReadRequest request = mock(PlcReadRequest.class,
RETURNS_DEEP_STUBS);
+ when(endpoint.buildPlcReadRequest()).thenReturn(request);
+
+ Exchange exchange = new Plc4XPollingConsumer(endpoint).receive(1000);
+
+ assertNotNull(exchange);
+ assertNull(exchange.getException());
+ assertInstanceOf(Map.class, exchange.getIn().getBody());
+ }
+}
diff --git
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 858917b437dd..b0d4eeb66d52 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -774,6 +774,20 @@ whose request times out now fails with a
`TimeoutException`, like `GET`, instead
An agent that answers with an OID inside the subtree that does not increase
now also fails the exchange with a
`CamelExchangeException` ("OID not increasing", like net-snmp's `snmpwalk`),
instead of repeating the request.
+=== camel-plc4x - failed writes and reads
+
+A `plc4x` producer with `autoReconnect=true` whose reconnect fails now fails
the exchange with the
+`PlcConnectionException`. Before, it logged "Unable to reconnect, skipping
request" and the exchange completed as if the
+values had been written. A failed write also keeps the message of the exchange
(it was removed before), so that a
+redelivery writes the same values again.
+
+The `plc4x` polling consumer (for example `pollEnrich`) now returns an
exchange that has the exception when the
+connection or the read fails, instead of an exchange with an empty `Map` body,
and it returns `null` when the PLC does
+not answer within the timeout of `receive(timeout)` or `receiveNoWait()`, as
other polling consumers do. With
+`pollEnrich` the exception fails the exchange (unless `aggregateOnException`
is enabled). A route that polls a PLC
+periodically and should go on while the PLC is unreachable can handle the
exception, for example with `onException`
+or `doTry`/`doCatch`.
+
=== camel-vertx - request/reply over the event bus
When a `vertx` consumer receives a message that expects a reply and the route
fails, the sender now gets a failure