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 b0bbea701568 CAMEL-25053: camel-core - Stream resequencer stops
delivering after a route restart while a message waits for its timeout (#26933)
b0bbea701568 is described below
commit b0bbea701568e9a07eac6198c98ceb111a45c8bb
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 14:13:16 2026 +0530
CAMEL-25053: camel-core - Stream resequencer stops delivering after a route
restart while a message waits for its timeout (#26933)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../apache/camel/processor/StreamResequencer.java | 26 +++-
.../processor/resequencer/ResequencerEngine.java | 85 ++++++++++--
.../processor/StreamResequencerRestartTest.java | 129 ++++++++++++++++++
.../resequencer/ResequencerEngineTest.java | 147 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 6 +
5 files changed, 381 insertions(+), 12 deletions(-)
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/StreamResequencer.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/StreamResequencer.java
index 834bcd2502b0..6dd3662f7d03 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/StreamResequencer.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/StreamResequencer.java
@@ -18,6 +18,7 @@ package org.apache.camel.processor;
import java.util.ArrayList;
import java.util.List;
+import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
@@ -233,8 +234,12 @@ public class StreamResequencer extends BaseProcessorSupport
@Override
protected void doStop() throws Exception {
// let's stop everything in the reverse order
- // no need to stop the worker thread -- it will stop automatically
when this service is stopped
+ // this also releases the callers waiting for capacity
engine.stop();
+ // end the worker thread, so that a quick restart does not leave it
running next to the new one
+ if (delivery != null) {
+ delivery.cancel();
+ }
ServiceHelper.stopService(processor);
}
@@ -258,6 +263,11 @@ public class StreamResequencer extends BaseProcessorSupport
exchange.setException(e);
callback.done(true);
return true;
+ } catch (RejectedExecutionException e) {
+ // stopped while waiting for capacity
+ exchange.setException(e);
+ callback.done(true);
+ return true;
}
try {
@@ -265,6 +275,10 @@ public class StreamResequencer extends BaseProcessorSupport
Exchange copy = ExchangeHelper.createCorrelatedCopy(exchange,
true);
engine.insert(copy);
delivery.request();
+ } catch (RejectedExecutionException e) {
+ // stopped after the wait for capacity: the exchange is not
queued, so it must fail even when invalid
+ // exchanges are ignored
+ exchange.setException(e);
} catch (Exception e) {
if (isIgnoreInvalidExchanges()) {
LOG.debug("Invalid Exchange. This Exchange will be ignored:
{}", exchange);
@@ -297,6 +311,7 @@ public class StreamResequencer extends BaseProcessorSupport
private final Lock deliveryRequestLock = new ReentrantLock();
private final Condition deliveryRequestCondition =
deliveryRequestLock.newCondition();
+ private volatile boolean cancelled;
Delivery() {
super(camelContext.getExecutorServiceManager().resolveThreadName("Resequencer
Delivery"));
@@ -304,7 +319,7 @@ public class StreamResequencer extends BaseProcessorSupport
@Override
public void run() {
- while (isRunAllowed()) {
+ while (isRunAllowed() && !cancelled) {
try {
deliveryRequestLock.lock();
try {
@@ -316,6 +331,9 @@ public class StreamResequencer extends BaseProcessorSupport
Thread.currentThread().interrupt();
break;
}
+ if (cancelled) {
+ break;
+ }
try {
engine.deliver();
} catch (Exception t) {
@@ -326,7 +344,9 @@ public class StreamResequencer extends BaseProcessorSupport
}
public void cancel() {
- interrupt();
+ // a flag rather than an interrupt, as this thread may be
delivering an exchange right now
+ cancelled = true;
+ request();
}
public void request() {
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java
index cb17b34814b6..5b5f910e7f18 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/resequencer/ResequencerEngine.java
@@ -20,6 +20,7 @@ import java.util.HashMap;
import java.util.Map;
import java.util.Timer;
import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
import java.util.function.Predicate;
@@ -92,6 +93,18 @@ public class ResequencerEngine<E> {
private final Lock lock = new ReentrantLock();
+ /**
+ * Whether this resequencer is stopped. Callers waiting in {@link
#waitUntil(Predicate)} are released on stop, and
+ * no elements are inserted or delivered while stopped.
+ */
+ private volatile boolean stopped;
+
+ /**
+ * Incremented on every stop, under the lock, so that a caller released
from {@link #waitUntil(Predicate)} by a stop
+ * fails even if this resequencer has been started again before it wakes
up.
+ */
+ private volatile long stopCount;
+
/**
* Creates a new resequencer instance with a default timeout of 2000
milliseconds.
*
@@ -103,16 +116,53 @@ public class ResequencerEngine<E> {
this.lastDelivered = null;
}
+ /**
+ * Starts this resequencer. Elements that were still waiting for their
timeout when this resequencer was stopped
+ * wait for a full timeout again from now on.
+ */
public void start() {
- timer = new Timer(
- ThreadHelper.resolveThreadName("Camel Thread ${counter} -
${name}", "Stream Resequencer Timer"), true);
+ lock.lock();
+ try {
+ timer = new Timer(
+ ThreadHelper.resolveThreadName("Camel Thread ${counter} -
${name}", "Stream Resequencer Timer"), true);
+ stopped = false;
+ // the timeouts that were pending on stop were discarded with the
old timer, so schedule them on the new
+ // one, otherwise these elements, and all elements behind them,
would never be delivered
+ for (Element<E> element : sequence) {
+ if (element.scheduled()) {
+ element.schedule(defineTimeout());
+ }
+ }
+ } finally {
+ lock.unlock();
+ }
}
/**
- * Stops this resequencer (i.e. this resequencer's {@link Timer} instance).
+ * Stops this resequencer (i.e. this resequencer's {@link Timer}
instance), and releases the callers waiting in
+ * {@link #waitUntil(Predicate)}, which then fail with a {@link
RejectedExecutionException}. A {@link #deliver()} in
+ * progress ends after the element it is sending, so this method only
waits for that element, not for all the
+ * elements that are ready for delivery. The remaining elements are kept
and delivered after a restart.
*/
public void stop() {
- timer.cancel();
+ // set before taking the lock, which a running deliver() holds while
it sends elements: it then returns after
+ // the element being sent
+ stopped = true;
+ lock.lock();
+ try {
+ // under the lock, so that an insert() that has not seen the flag
can still schedule its timeout, which
+ // start() schedules again
+ if (timer != null) {
+ timer.cancel();
+ }
+ stopCount++;
+ for (CountDownLatch latch : waitConditions.keySet()) {
+ latch.countDown();
+ }
+ waitConditions.clear();
+ } finally {
+ lock.unlock();
+ }
}
/**
@@ -133,22 +183,32 @@ public class ResequencerEngine<E> {
* Wait for the following condition to happen. Do not call this method
while holding a lock on the resequencer
* engine, as it will deadlock. The predicate will be evaluated while
holding a lock on the resequencer engine.
*
- * @param pred the condition to wait for
- * @throws InterruptedException if the thread is interrupted
+ * @param pred the condition to wait for
+ * @throws InterruptedException if the thread is interrupted
+ * @throws RejectedExecutionException if this resequencer is stopped, or
is stopped while waiting
*/
public void waitUntil(Predicate<Sequence<?>> pred) throws
InterruptedException {
CountDownLatch latch;
+ long stopCountAtWait;
lock.lock();
try {
+ if (stopped) {
+ throw new RejectedExecutionException("Resequencer is stopped");
+ }
if (pred.test(sequence)) {
return;
}
latch = new CountDownLatch(1);
waitConditions.put(latch, pred);
+ stopCountAtWait = stopCount;
} finally {
lock.unlock();
}
latch.await();
+ // compare with the stop count rather than read the flag, as a quick
restart may have reset it already
+ if (stopped || stopCount != stopCountAtWait) {
+ throw new RejectedExecutionException("Resequencer is stopped");
+ }
}
private void evaluateConditions() {
@@ -235,12 +295,18 @@ public class ResequencerEngine<E> {
* Inserts the given element into this resequencer. If the element is not
ready for immediate delivery and has no
* immediate presecessor then it is scheduled for timing out. After being
timed out it is ready for delivery.
*
- * @param o an element.
- * @throws IllegalArgumentException if the element cannot be used with
this resequencer engine
+ * @param o an element.
+ * @throws IllegalArgumentException if the element cannot be used with
this resequencer engine
+ * @throws RejectedExecutionException if this resequencer is stopped
*/
public void insert(E o) {
lock.lock();
try {
+ // a stopped resequencer has cancelled its timer, so the element
could not be scheduled for timing out
+ if (stopped) {
+ throw new RejectedExecutionException("Resequencer is stopped");
+ }
+
// wrap object into internal element
Element<E> element = new Element<>(o);
@@ -312,7 +378,8 @@ public class ResequencerEngine<E> {
public boolean deliverNext() throws Exception {
lock.lock();
try {
- if (sequence.isEmpty()) {
+ // once stopped, the elements are kept for a restart, so that
stop() does not wait for all of them
+ if (stopped || sequence.isEmpty()) {
return false;
}
// inspect element with the lowest sequence value
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/StreamResequencerRestartTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/StreamResequencerRestartTest.java
new file mode 100644
index 000000000000..e4b790a4ef8b
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/StreamResequencerRestartTest.java
@@ -0,0 +1,129 @@
+/*
+ * 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.processor;
+
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Future;
+import java.util.concurrent.RejectedExecutionException;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+
+/**
+ * The stream resequencer must keep working when its route is stopped and
started while a message waits for its timeout,
+ * callers waiting for free capacity must be released when the route stops,
and a stop must end the delivery thread.
+ */
+public class StreamResequencerRestartTest extends ContextTestSupport {
+
+ @Test
+ public void testTimeoutExpiresAfterRouteRestart() throws Exception {
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("msg1");
+ template.sendBodyAndHeader("direct:start", "msg1", "seqnum", 1L);
+ // the first message waits for the timeout, as there is no earlier
message to compare it with
+ mock.assertIsSatisfied();
+
+ mock.reset();
+ mock.expectedBodiesReceived("msg3", "msg4", "msg5");
+ // msg2 is missing, so msg3 waits for its timeout, and the route is
restarted meanwhile
+ template.sendBodyAndHeader("direct:start", "msg3", "seqnum", 3L);
+ context.getRouteController().stopRoute("stream");
+ context.getRouteController().startRoute("stream");
+ template.sendBodyAndHeader("direct:start", "msg4", "seqnum", 4L);
+ template.sendBodyAndHeader("direct:start", "msg5", "seqnum", 5L);
+
+ mock.assertIsSatisfied();
+ }
+
+ @Test
+ public void testCallerWaitingForCapacityIsReleasedOnStop() throws
Exception {
+ // the waiting caller is inflight, so the graceful stop of the route
waits for it until the shutdown timeout
+ context.getShutdownStrategy().setTimeout(1);
+
+ // capacity 1: msg2 waits for msg1, and the next caller waits for free
capacity
+ template.sendBodyAndHeader("direct:capacity", "msg2", "seqnum", 2L);
+ ExecutorService executor =
context.getExecutorServiceManager().newSingleThreadExecutor(this, "sender");
+ try {
+ Future<Exchange> sender = executor.submit(() ->
template.send("direct:capacity", e -> {
+ e.getMessage().setBody("msg3");
+ e.getMessage().setHeader("seqnum", 3L);
+ }));
+ await().atMost(5, TimeUnit.SECONDS).until(() ->
context.getInflightRepository().size("capacity") == 1);
+
+ context.getRouteController().stopRoute("capacity");
+
+ Exchange out = sender.get(10, TimeUnit.SECONDS);
+ assertInstanceOf(RejectedExecutionException.class,
out.getException());
+ } finally {
+ context.getExecutorServiceManager().shutdownNow(executor);
+ }
+ }
+
+ @Test
+ public void testDeliveryThreadEndsOnStop() throws Exception {
+ // only the delivery thread of the route under test is left running
+ context.getRouteController().stopRoute("stream");
+ context.getRouteController().stopRoute("capacity");
+ await().atMost(10, TimeUnit.SECONDS).until(() -> deliveryThreads() ==
1);
+
+ // stopped and started again within the delivery attempt interval
+ for (int i = 0; i < 3; i++) {
+ context.getRouteController().stopRoute("delivery");
+ context.getRouteController().startRoute("delivery");
+ }
+
+ // only the delivery thread of the last start is left
+ await().atMost(10, TimeUnit.SECONDS).until(() -> deliveryThreads() ==
1);
+ context.getRouteController().stopRoute("delivery");
+ await().atMost(10, TimeUnit.SECONDS).until(() -> deliveryThreads() ==
0);
+ }
+
+ private long deliveryThreads() {
+ return Thread.getAllStackTraces().keySet().stream()
+ .filter(Thread::isAlive)
+ .filter(t -> t.getName().contains("(" + context.getName() +
")"))
+ .filter(t -> t.getName().endsWith("Resequencer Delivery"))
+ .count();
+ }
+
+ @Override
+ protected RouteBuilder createRouteBuilder() {
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start").routeId("stream")
+
.resequence(header("seqnum")).stream().timeout(1000).deliveryAttemptInterval(10)
+ .to("mock:result");
+
+ from("direct:capacity").routeId("capacity")
+
.resequence(header("seqnum")).stream().capacity(1).timeout(60000).deliveryAttemptInterval(10)
+ .to("mock:capacity");
+
+ from("direct:delivery").routeId("delivery")
+
.resequence(header("seqnum")).stream().deliveryAttemptInterval(1000)
+ .to("mock:delivery");
+ }
+ };
+ }
+}
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/resequencer/ResequencerEngineTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/resequencer/ResequencerEngineTest.java
index 7140748f2c2a..da3c981cb11f 100644
---
a/core/camel-core/src/test/java/org/apache/camel/processor/resequencer/ResequencerEngineTest.java
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/resequencer/ResequencerEngineTest.java
@@ -19,7 +19,13 @@ package org.apache.camel.processor.resequencer;
import java.util.LinkedList;
import java.util.List;
import java.util.Random;
+import java.util.concurrent.CopyOnWriteArrayList;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutionException;
+import java.util.concurrent.FutureTask;
+import java.util.concurrent.RejectedExecutionException;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.TimeoutException;
import org.apache.camel.TestSupport;
import org.apache.camel.util.StopWatch;
@@ -28,8 +34,13 @@ import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.DisabledIf;
import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.junit.jupiter.api.Assertions.fail;
class ResequencerEngineTest extends TestSupport {
@@ -89,6 +100,142 @@ class ResequencerEngineTest extends TestSupport {
assertEquals(4, resequencer.getLastDelivered());
}
+ @Test
+ void testTimeoutAfterRestart() throws Exception {
+ SequenceBuffer<Integer> out = new SequenceBuffer<>();
+ ResequencerEngine<Integer> engine = new ResequencerEngine<>(new
IntegerComparator());
+ engine.setSequenceSender(out);
+ // long enough that the timeout of 4 cannot expire before the stop,
even with a pause of the test
+ engine.setTimeout(2000);
+ engine.start();
+ try {
+ engine.setLastDelivered(2);
+ // 3 is missing, so 4 waits for its timeout
+ engine.insert(4);
+ engine.stop();
+ engine.start();
+ engine.insert(5);
+
+ // the timeout of 4 must still expire after the restart
+ await().atMost(10, TimeUnit.SECONDS).until(engine::deliverNext);
+ engine.deliver();
+ assertEquals(4, out.poll(0));
+ assertEquals(5, out.poll(0));
+ } finally {
+ engine.stop();
+ }
+ }
+
+ @Test
+ void testWaitUntilReleasedOnStop() throws Exception {
+ ResequencerEngine<Integer> engine = new ResequencerEngine<>(new
IntegerComparator());
+ engine.setSequenceSender(new SequenceBuffer<>());
+ engine.start();
+ engine.insert(4);
+
+ // a caller waiting for free capacity
+ FutureTask<Void> waiter = new FutureTask<>(() -> {
+ engine.waitUntil(s -> s.size() < 1);
+ return null;
+ });
+ Thread thread = new Thread(waiter, "waiter");
+ thread.setDaemon(true);
+ thread.start();
+ await().atMost(5, TimeUnit.SECONDS).until(() -> thread.getState() ==
Thread.State.WAITING);
+
+ engine.stop();
+ try {
+ waiter.get(5, TimeUnit.SECONDS);
+ fail("waitUntil should fail when the resequencer is stopped");
+ } catch (ExecutionException e) {
+ assertInstanceOf(RejectedExecutionException.class, e.getCause());
+ } catch (TimeoutException e) {
+ thread.interrupt();
+ fail("waitUntil is still blocked after the resequencer was
stopped");
+ }
+ }
+
+ @Test
+ void testStopDoesNotWaitForTheReadyBacklog() throws Exception {
+ CountDownLatch sendingFirst = new CountDownLatch(1);
+ CountDownLatch releaseFirst = new CountDownLatch(1);
+ CountDownLatch releaseOthers = new CountDownLatch(1);
+ List<Integer> sent = new CopyOnWriteArrayList<>();
+ ResequencerEngine<Integer> engine = new ResequencerEngine<>(new
IntegerComparator());
+ // a slow downstream processor: the first element is held until
released, the others until the end of the test
+ engine.setSequenceSender(o -> {
+ sent.add(o);
+ if (o == 1) {
+ sendingFirst.countDown();
+ releaseFirst.await(10, TimeUnit.SECONDS);
+ } else {
+ releaseOthers.await(10, TimeUnit.SECONDS);
+ }
+ });
+ engine.start();
+ engine.setLastDelivered(0);
+ for (int i = 1; i <= 5; i++) {
+ engine.insert(i);
+ }
+
+ FutureTask<Void> delivery = new FutureTask<>(() -> {
+ engine.deliver();
+ return null;
+ });
+ Thread deliveryThread = new Thread(delivery, "delivery");
+ deliveryThread.setDaemon(true);
+ FutureTask<Void> stop = new FutureTask<>(() -> {
+ engine.stop();
+ return null;
+ });
+ Thread stopThread = new Thread(stop, "stop");
+ stopThread.setDaemon(true);
+ try {
+ deliveryThread.start();
+ assertTrue(sendingFirst.await(5, TimeUnit.SECONDS), "the first
element was not delivered");
+
+ // the route stops while the downstream processor is busy with the
first element
+ stopThread.start();
+ // stop waits for the element being sent
+ await().atMost(5, TimeUnit.SECONDS)
+ .until(() -> stop.isDone() || stopThread.getState() ==
Thread.State.WAITING);
+ releaseFirst.countDown();
+
+ // stop returns after the element being sent, without sending the
other ready elements
+ try {
+ stop.get(5, TimeUnit.SECONDS);
+ } catch (TimeoutException e) {
+ fail("stop waits for the downstream processing of all ready
elements");
+ }
+ delivery.get(5, TimeUnit.SECONDS);
+ assertEquals(List.of(1), sent);
+ assertEquals(4, engine.size());
+ } finally {
+ releaseFirst.countDown();
+ releaseOthers.countDown();
+ }
+ }
+
+ @Test
+ void testInsertRejectedWhenStopped() throws Exception {
+ ResequencerEngine<Integer> engine = new ResequencerEngine<>(new
IntegerComparator());
+ engine.setSequenceSender(new SequenceBuffer<>());
+ engine.start();
+ engine.insert(4);
+ engine.stop();
+
+ assertThrows(RejectedExecutionException.class, () -> engine.insert(5));
+ // not queued, so it is not delivered after a restart
+ assertEquals(1, engine.size());
+ }
+
+ @Test
+ void testStopWithoutStart() {
+ // BaseService.start() calls stop() when the start fails before the
engine was started
+ ResequencerEngine<Integer> engine = new ResequencerEngine<>(new
IntegerComparator());
+ assertDoesNotThrow(engine::stop);
+ }
+
@DisabledIf(value = "isIgnoreLoadTests",
disabledReason = "Enabled only when the System property
'ignore.load.tests' is not set to 'true'")
@Test
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 3361cfdc7f05..08f6cd9b845c 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
@@ -164,6 +164,12 @@ ratios that are all `0`, or ratios whose sum is greater
than `2147483647` now fa
`IllegalArgumentException`. Previously such a route started, but sending to it
could hang the caller, spin a CPU,
or send every message to the same endpoint.
+=== Resequencer EIP - stream resequencing and route stop
+
+When a route with the stream resequencer stops, a caller that is blocked
because the resequencer is full (`capacity`)
+now fails with a `java.util.concurrent.RejectedExecutionException`. Prior to
Camel 4.23 such a caller stayed blocked,
+also after the route was stopped.
+
=== Tokenize language
- A message without a body (or a `source` that has no value) now has no
tokens, the same as an empty body.