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();
+        }
+    }
+}

Reply via email to