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 37b9750d81ef CAMEL-25292: camel-plc4x - a stopped triggered consumer
must stop its scraper, and the triggered consumer must start (#27323)
37b9750d81ef is described below
commit 37b9750d81ef6d4ee5cf8880b791ea2fbc924f4c
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:33:31 2026 +0530
CAMEL-25292: camel-plc4x - a stopped triggered consumer must stop its
scraper, and the triggered consumer must start (#27323)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
components/camel-plc4x/pom.xml | 6 --
.../camel/component/plc4x/Plc4XConsumer.java | 25 ++++++--
.../component/plc4x/Plc4XConsumerStopTest.java | 67 ++++++++++++++++++++++
3 files changed, 88 insertions(+), 10 deletions(-)
diff --git a/components/camel-plc4x/pom.xml b/components/camel-plc4x/pom.xml
index 062666411015..bb2424bf957a 100644
--- a/components/camel-plc4x/pom.xml
+++ b/components/camel-plc4x/pom.xml
@@ -48,12 +48,6 @@
<groupId>org.apache.plc4x</groupId>
<artifactId>plc4j-scraper</artifactId>
<version>${plc4x-version}</version>
- <exclusions>
- <exclusion>
- <groupId>com.fasterxml.jackson.core</groupId>
- <artifactId>jackson-databind</artifactId>
- </exclusion>
- </exclusions>
</dependency>
<!-- Include all drivers -->
diff --git
a/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XConsumer.java
b/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XConsumer.java
index f16ade351fd5..581191a7929d 100644
---
a/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XConsumer.java
+++
b/components/camel-plc4x/src/main/java/org/apache/camel/component/plc4x/Plc4XConsumer.java
@@ -20,7 +20,6 @@ import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
-import java.util.concurrent.Executors;
import java.util.concurrent.ScheduledExecutorService;
import java.util.concurrent.ScheduledFuture;
import java.util.concurrent.TimeUnit;
@@ -48,8 +47,10 @@ public class Plc4XConsumer extends DefaultConsumer {
private final String trigger;
private final Plc4XEndpoint plc4XEndpoint;
- private final ScheduledExecutorService executorService =
Executors.newSingleThreadScheduledExecutor();
+ private ScheduledExecutorService executorService;
private ScheduledFuture<?> future;
+ private TriggeredScraperImpl scraper;
+ private TriggerCollector collector;
public Plc4XConsumer(Plc4XEndpoint endpoint, Processor processor) {
super(endpoint, processor);
@@ -105,6 +106,8 @@ public class Plc4XConsumer extends DefaultConsumer {
return;
}
+ executorService =
getEndpoint().getCamelContext().getExecutorServiceManager()
+ .newSingleThreadScheduledExecutor(this, "Plc4XConsumer");
future = executorService.schedule(() ->
request.execute().thenAccept(response -> {
try {
Exchange exchange = plc4XEndpoint.createExchange();
@@ -122,9 +125,9 @@ public class Plc4XConsumer extends DefaultConsumer {
private void startTriggered() throws ScraperException {
ScraperConfiguration configuration = getScraperConfig(tags);
- TriggerCollector collector = new
TriggerCollectorImpl(plc4XEndpoint.getPlcDriverManager());
+ collector = new
TriggerCollectorImpl(plc4XEndpoint.getPlcDriverManager());
- TriggeredScraperImpl scraper = new TriggeredScraperImpl(configuration,
(job, alias, response) -> {
+ scraper = new TriggeredScraperImpl(configuration, (job, alias,
response) -> {
try {
plc4XEndpoint.reconnectIfNeeded();
@@ -158,6 +161,20 @@ public class Plc4XConsumer extends DefaultConsumer {
// First stop the polling process
if (future != null) {
future.cancel(true);
+ future = null;
+ }
+ if (executorService != null) {
+
getEndpoint().getCamelContext().getExecutorServiceManager().shutdownNow(executorService);
+ executorService = null;
+ }
+ // a scraper that is not stopped keeps reading the PLC and sending
exchanges to the route
+ if (scraper != null) {
+ scraper.stop();
+ scraper = null;
+ }
+ if (collector != null) {
+ collector.stop();
+ collector = null;
}
super.doStop();
}
diff --git
a/components/camel-plc4x/src/test/java/org/apache/camel/component/plc4x/Plc4XConsumerStopTest.java
b/components/camel-plc4x/src/test/java/org/apache/camel/component/plc4x/Plc4XConsumerStopTest.java
new file mode 100644
index 000000000000..f2362156a147
--- /dev/null
+++
b/components/camel-plc4x/src/test/java/org/apache/camel/component/plc4x/Plc4XConsumerStopTest.java
@@ -0,0 +1,67 @@
+/*
+ * 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 org.apache.camel.Processor;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.plc4x.java.scraper.triggeredscraper.TriggeredScraperImpl;
+import
org.apache.plc4x.java.scraper.triggeredscraper.triggerhandler.collector.TriggerCollectorImpl;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedConstruction;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * A stopped triggered consumer must stop its scraper, or the scraper keeps
reading the PLC and sending exchanges to the
+ * route, also after the route is started again with a new consumer
(duplicates).
+ */
+class Plc4XConsumerStopTest {
+
+ @Test
+ void testStopStopsTheScraper() throws Exception {
+ Plc4XEndpoint endpoint = mock(Plc4XEndpoint.class);
+ when(endpoint.getCamelContext()).thenReturn(new DefaultCamelContext());
+
when(endpoint.getTrigger()).thenReturn("(PLC4X_TRIGGER_VAR,10,(%DB1:DBX0.0:BOOL)==(true))");
+ when(endpoint.getTags()).thenReturn(Map.of("tag1", "%DB1:DBW2:INT"));
+ when(endpoint.getPeriod()).thenReturn(1000);
+ when(endpoint.getUri()).thenReturn("mock:plc");
+
+ try (MockedConstruction<TriggeredScraperImpl> scrapers =
mockConstruction(TriggeredScraperImpl.class);
+ MockedConstruction<TriggerCollectorImpl> collectors =
mockConstruction(TriggerCollectorImpl.class)) {
+ Plc4XConsumer consumer = new Plc4XConsumer(endpoint,
mock(Processor.class));
+ consumer.start();
+
+ assertEquals(1, scrapers.constructed().size());
+ TriggeredScraperImpl scraper = scrapers.constructed().get(0);
+ TriggerCollectorImpl collector = collectors.constructed().get(0);
+ verify(scraper).start();
+ verify(scraper, never()).stop();
+
+ consumer.stop();
+
+ verify(scraper).stop();
+ verify(collector).stop();
+ }
+ }
+}