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 358a4d36fa23 CAMEL-25252: camel-timer - only release the shared timer 
when the consumer took it (#27249)
358a4d36fa23 is described below

commit 358a4d36fa23415c45b83724c5098af00f9f88ed
Author: allthingssecurity <[email protected]>
AuthorDate: Fri Oct 2 12:55:46 2026 +0530

    CAMEL-25252: camel-timer - only release the shared timer when the consumer 
took it (#27249)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../camel/component/timer/TimerConsumer.java       | 11 ++-
 .../timer/TimerNegativeDelayStopRouteTest.java     | 80 ++++++++++++++++++++++
 2 files changed, 89 insertions(+), 2 deletions(-)

diff --git 
a/components/camel-timer/src/main/java/org/apache/camel/component/timer/TimerConsumer.java
 
b/components/camel-timer/src/main/java/org/apache/camel/component/timer/TimerConsumer.java
index 396909eb1ad5..5134771fe48f 100644
--- 
a/components/camel-timer/src/main/java/org/apache/camel/component/timer/TimerConsumer.java
+++ 
b/components/camel-timer/src/main/java/org/apache/camel/component/timer/TimerConsumer.java
@@ -48,6 +48,8 @@ public class TimerConsumer extends DefaultConsumer implements 
StartupListener, S
     private ExecutorService executorService;
     private final AtomicLong counter = new AtomicLong();
     private volatile boolean polling;
+    // whether this consumer holds a reference to the shared timer, which it 
must release when it is stopped
+    private volatile boolean timerTaken;
 
     public TimerConsumer(TimerEndpoint endpoint, Processor processor) {
         super(endpoint, processor);
@@ -183,6 +185,7 @@ public class TimerConsumer extends DefaultConsumer 
implements StartupListener, S
             // the StartupListener is configuring the task later
             if (task != null && !configured && 
endpoint.getCamelContext().getStatus().isStarted()) {
                 Timer timer = endpoint.getTimer(this);
+                timerTaken = true;
                 configureTask(task, timer);
             }
         } else {
@@ -213,8 +216,11 @@ public class TimerConsumer extends DefaultConsumer 
implements StartupListener, S
         task = null;
         configured = false;
 
-        // remove timer
-        endpoint.removeTimer(this);
+        // release the timer, but only if this consumer took it (a consumer 
with a negative delay does not use it)
+        if (timerTaken) {
+            timerTaken = false;
+            endpoint.removeTimer(this);
+        }
 
         // if executorService is instantiated then we shutdown it
         if (executorService != null) {
@@ -229,6 +235,7 @@ public class TimerConsumer extends DefaultConsumer 
implements StartupListener, S
     public void onCamelContextStarted(CamelContext context, boolean 
alreadyStarted) throws Exception {
         if (task != null && !configured) {
             Timer timer = endpoint.getTimer(this);
+            timerTaken = true;
             configureTask(task, timer);
         }
     }
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/component/timer/TimerNegativeDelayStopRouteTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/component/timer/TimerNegativeDelayStopRouteTest.java
new file mode 100644
index 000000000000..56a0d392db76
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/component/timer/TimerNegativeDelayStopRouteTest.java
@@ -0,0 +1,80 @@
+/*
+ * 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.timer;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.junit.jupiter.api.Test;
+
+/**
+ * A timer consumer with a negative delay does not use the shared timer of its 
timer name, so stopping it must not
+ * cancel the timer of the other consumers with the same timer name.
+ */
+public class TimerNegativeDelayStopRouteTest extends ContextTestSupport {
+
+    @Test
+    public void testStopNegativeDelayRoute() throws Exception {
+        getMockEndpoint("mock:foo").expectedMinimumMessageCount(1);
+        getMockEndpoint("mock:bar").expectedMessageCount(1);
+
+        context.getRouteController().startRoute("bar");
+
+        assertMockEndpointsSatisfied();
+
+        context.getRouteController().stopRoute("bar");
+
+        resetMocks();
+
+        // stopping bar route, we should still keep getting messages to foo
+        getMockEndpoint("mock:foo").expectedMinimumMessageCount(2);
+
+        assertMockEndpointsSatisfied();
+    }
+
+    @Test
+    public void testRestartNegativeDelayRoute() throws Exception {
+        getMockEndpoint("mock:foo").expectedMinimumMessageCount(1);
+        getMockEndpoint("mock:bar").expectedMessageCount(1);
+
+        context.getRouteController().startRoute("bar");
+
+        assertMockEndpointsSatisfied();
+
+        context.getRouteController().stopRoute("bar");
+        context.getRouteController().startRoute("bar");
+        context.getRouteController().stopRoute("bar");
+
+        resetMocks();
+
+        // the bar route was stopped twice, we should still keep getting 
messages to foo
+        getMockEndpoint("mock:foo").expectedMinimumMessageCount(2);
+
+        assertMockEndpointsSatisfied();
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("timer:mytimer?period=100").routeId("foo").to("mock:foo");
+
+                
from("timer:mytimer?delay=-1&repeatCount=1").routeId("bar").autoStartup(false).to("mock:bar");
+            }
+        };
+    }
+}

Reply via email to