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 dd1c491e71b9 CAMEL-24945: camel-core - Aggregate EIP: with optimistic
locking a completion should not remove the timeout of a new group (#26790)
dd1c491e71b9 is described below
commit dd1c491e71b9e4f12999dfeb5de0a04e57ead6a7
Author: allthingssecurity <[email protected]>
AuthorDate: Mon Sep 28 19:12:25 2026 +0530
CAMEL-24945: camel-core - Aggregate EIP: with optimistic locking a
completion should not remove the timeout of a new group (#26790)
The completion timeout map is keyed by correlation key. When a group was
completed by an incoming exchange (or by force or interval completion),
onCompletion removed the group from the repository and then called
timeoutMap.remove(key). With optimistic locking there is no lock around
these
calls, so an exchange for the same key could start a new group in between:
it
registers its timeout (trackTimeout runs before the repository add) and adds
the group, and the completing thread then removed that new entry. The new
group
never timed out; it stayed in the repository until another exchange for the
key arrived or the route was stopped with forceCompletionOnStop.
With optimistic locking onCompletion no longer removes the timeout entry. An
entry left behind is harmless: it expires, and the timeout checker
completes a
group only if it can remove the current group from the repository (compare
and
set), while every new group refreshes the entry of its key before it is
added.
In the same mode onEviction no longer skips an entry because its exchange
id is
in progress: that id can belong to a group that was completed concurrently
while the repository holds a newer group for the key that relies on this
entry
(two exchanges starting a new group at the same time), which would then be
stranded as well. The compare-and-set removal already prevents completing a
group twice. Pessimistic locking is unchanged.
Both parts were checked with a TLA+ model of the aggregator, which finds a
stranded group with either part alone.
With a completionTimeoutExpression a new group can have no timeout, so the
left over entry of the previous group could complete it. With optimistic
locking and a completionTimeoutExpression the group now records its
completion timeout in the CamelAggregatedTimeout property of the stored
exchange (updated with the group by compare and set), and the timeout
checker does not complete a group without it. The 4.23 upgrade guide
documents that groups persisted before the upgrade do not have it.
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../processor/aggregate/AggregateProcessor.java | 55 +++-
.../AggregateOptimisticLockingTimeoutTest.java | 355 +++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 12 +
3 files changed, 419 insertions(+), 3 deletions(-)
diff --git
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
index be77dbbfb0b3..572eeb75d559 100644
---
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
+++
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
@@ -547,6 +547,13 @@ public class AggregateProcessor extends
BaseProcessorSupport
// prepare the exchanges for aggregation
ExchangeHelper.prepareAggregation(oldExchange, newExchange);
+ // the group records its completion timeout (see updateGroupTimeout),
so only a timeout that this aggregator
+ // tracks for the exchange counts
+ boolean groupTimeout = isGroupTimeout();
+ if (groupTimeout) {
+ newExchange.removeProperty(ExchangePropertyKey.AGGREGATED_TIMEOUT);
+ }
+
// check if we are pre complete
if (preCompletion) {
try {
@@ -652,6 +659,9 @@ public class AggregateProcessor extends BaseProcessorSupport
}
if (!aggregateFailed && complete == null) {
+ if (groupTimeout) {
+ updateGroupTimeout(newExchange, originalExchange, answer);
+ }
// only need to update aggregation repository if we are not
complete
doAggregationRepositoryAdd(newExchange.getContext(), key,
originalExchange, answer);
} else {
@@ -787,6 +797,35 @@ public class AggregateProcessor extends
BaseProcessorSupport
return null;
}
+ /**
+ * Whether the group records if it has a completion timeout, and the
timeout checker only completes a group that has
+ * one.
+ * <p/>
+ * With optimistic locking the timeout entry of a completed group is not
removed (see
+ * {@link #onCompletion(String, Exchange, Exchange, boolean, boolean)}),
so it can be left over when a new group for
+ * the same correlation key starts. A new group normally replaces the
entry with its own timeout, but with a
+ * completion timeout expression a new group can have no timeout, and the
left over entry must then not complete it.
+ */
+ private boolean isGroupTimeout() {
+ return optimisticLocking && completionTimeoutExpression != null;
+ }
+
+ /**
+ * Records the completion timeout of the group on the aggregated exchange,
which is stored in the repository: the
+ * timeout tracked for the new exchange, otherwise the timeout of the
group so far.
+ */
+ private static void updateGroupTimeout(Exchange newExchange, Exchange
originalExchange, Exchange answer) {
+ Object timeout =
newExchange.getProperty(ExchangePropertyKey.AGGREGATED_TIMEOUT);
+ if (timeout == null && originalExchange != null) {
+ timeout =
originalExchange.getProperty(ExchangePropertyKey.AGGREGATED_TIMEOUT);
+ }
+ if (timeout != null) {
+ answer.setProperty(ExchangePropertyKey.AGGREGATED_TIMEOUT,
timeout);
+ } else {
+ answer.removeProperty(ExchangePropertyKey.AGGREGATED_TIMEOUT);
+ }
+ }
+
protected void trackTimeout(String key, Exchange exchange) {
// timeout can be either evaluated based on an expression or from a
fixed value
// expression takes precedence
@@ -832,8 +871,11 @@ public class AggregateProcessor extends
BaseProcessorSupport
aggregationRepository.remove(aggregated.getContext(), key,
original);
}
- if (!fromTimeout && timeoutMap != null) {
- // cleanup timeout map if it was a incoming exchange which
triggered the timeout (and not the timeout checker)
+ // cleanup timeout map if it was a incoming exchange which triggered
the timeout (and not the timeout checker)
+ // but not with optimistic locking: the timeout map is keyed by
correlation key, and without a lock the entry
+ // may already belong to a new group for the same key, which would
then never time out. An entry left behind
+ // does no harm, as the timeout checker completes a group only if it
can remove it from the repository.
+ if (!fromTimeout && timeoutMap != null && !optimisticLocking) {
LOG.trace("Removing correlation key {} from timeout", key);
timeoutMap.remove(key);
}
@@ -1362,7 +1404,9 @@ public class AggregateProcessor extends
BaseProcessorSupport
}
log.debug("Completion timeout triggered for correlation key: {}",
key);
- boolean inProgress =
inProgressCompleteExchanges.contains(exchangeId);
+ // with optimistic locking the exchange id in the entry can belong
to a group that has been completed
+ // while a newer group for the same key is in the repository, so
the repository decides (see below)
+ boolean inProgress = !optimisticLocking &&
inProgressCompleteExchanges.contains(exchangeId);
if (inProgress) {
log.trace("Aggregated exchange with id: {} is already in
progress.", exchangeId);
return;
@@ -1373,6 +1417,11 @@ public class AggregateProcessor extends
BaseProcessorSupport
Exchange answer = aggregationRepository.get(camelContext, key);
if (answer == null) {
evictionStolen = true;
+ } else if (isGroupTimeout() &&
answer.getProperty(ExchangePropertyKey.AGGREGATED_TIMEOUT) == null) {
+ // the entry is left over from a completed group, and the
current group has no completion timeout
+ log.debug("Completion timeout for correlation key: {} is left
over from a completed group, as the group"
+ + " has no completion timeout",
+ key);
} else {
// indicate it was completed by timeout
answer.setProperty(ExchangePropertyKey.AGGREGATED_COMPLETED_BY,
COMPLETED_BY_TIMEOUT);
diff --git
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateOptimisticLockingTimeoutTest.java
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateOptimisticLockingTimeoutTest.java
new file mode 100644
index 000000000000..1bb49143f862
--- /dev/null
+++
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateOptimisticLockingTimeoutTest.java
@@ -0,0 +1,355 @@
+/*
+ * 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.aggregator;
+
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.ScheduledFuture;
+import java.util.concurrent.ScheduledThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+import java.util.concurrent.atomic.AtomicLong;
+
+import org.apache.camel.AggregationStrategy;
+import org.apache.camel.CamelContext;
+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.apache.camel.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.spi.AggregationRepository;
+import org.apache.camel.spi.OptimisticLockingAggregationRepository;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.support.DefaultExchangeHolder;
+import org.apache.camel.support.service.ServiceSupport;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * With optimistic locking a group must not lose its completion timeout
because a group of the same correlation key is
+ * completed concurrently.
+ */
+public class AggregateOptimisticLockingTimeoutTest extends ContextTestSupport {
+
+ private final CountDownLatch paused = new CountDownLatch(1);
+ private final CountDownLatch release = new CountDownLatch(1);
+ private final CountDownLatch releaseDownstream = new CountDownLatch(1);
+ private final CountDownLatch timeoutGate = new CountDownLatch(1);
+ private final AtomicBoolean pause = new AtomicBoolean(true);
+ private final AtomicInteger timeoutCheckerRuns = new AtomicInteger();
+ private ScheduledExecutorService timeoutChecker;
+
+ @Override
+ public boolean isUseRouteBuilder() {
+ return false;
+ }
+
+ @Override
+ @BeforeEach
+ public void setUp() throws Exception {
+ super.setUp();
+ // the timeout checker only runs when the gate is open, so no group
times out before the test is ready
+ timeoutChecker = new ScheduledThreadPoolExecutor(1) {
+ @Override
+ public ScheduledFuture<?> scheduleWithFixedDelay(
+ Runnable command, long initialDelay, long delay, TimeUnit
unit) {
+ return super.scheduleWithFixedDelay(() -> {
+ if (awaitLatch(timeoutGate)) {
+ command.run();
+ timeoutCheckerRuns.incrementAndGet();
+ }
+ }, initialDelay, delay, unit);
+ }
+ };
+ }
+
+ @Override
+ @AfterEach
+ public void tearDown() throws Exception {
+ release.countDown();
+ releaseDownstream.countDown();
+ timeoutGate.countDown();
+ super.tearDown();
+ timeoutChecker.shutdownNow();
+ }
+
+ @Test
+ public void testCompletionDoesNotRemoveTimeoutOfNewGroup() throws
Exception {
+ // pauses the thread of b right after it removed the completed group
[a, b] from the repository
+ MemoryAggregationRepository repository = new
MemoryAggregationRepository(true) {
+ @Override
+ public void remove(CamelContext camelContext, String key, Exchange
exchange) {
+ super.remove(camelContext, key, exchange);
+ pauseIfThread("producer-b");
+ }
+ };
+ addRoute(repository);
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("a+b", "c");
+
mock.message(0).exchangeProperty(Exchange.AGGREGATED_COMPLETED_BY).isEqualTo("size");
+
mock.message(1).exchangeProperty(Exchange.AGGREGATED_COMPLETED_BY).isEqualTo("timeout");
+
+ send("a");
+ Thread producer = new Thread(() -> send("b"), "producer-b");
+ producer.start();
+ assertTrue(awaitLatch(paused));
+
+ // c starts a new group for the same key, which registers its
completion timeout
+ send("c");
+ // then the thread of b continues completing the group [a, b]
+ release.countDown();
+ producer.join(10000);
+
+ timeoutGate.countDown();
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testTimeoutWithExchangeIdOfCompletedGroup() throws Exception {
+ // pauses the thread of m1 before it adds its new group [m1], after it
registered its completion timeout
+ IdPreservingRepository repository = new IdPreservingRepository();
+ addRoute(repository);
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("m2+m3", "m1");
+
mock.message(1).exchangeProperty(Exchange.AGGREGATED_COMPLETED_BY).isEqualTo("timeout");
+
+ Thread producer = new Thread(() -> send("m1"), "producer-m1");
+ producer.start();
+ assertTrue(awaitLatch(paused));
+
+ // m2 starts a group (the timeout entry of the key now holds the
exchange id of m2), and m3 completes it;
+ // the aggregated exchange (with the exchange id of m2) is in progress
until it is released downstream
+ send("m2");
+ send("m3");
+ // m1 adds its group [m1]; the timeout entry of the key still holds
the exchange id of m2
+ release.countDown();
+ producer.join(10000);
+
+ timeoutGate.countDown();
+ // the group [m1] must be completed by the timeout while the exchange
m2+m3 is in progress
+ await("group [m1] completed by timeout").atMost(10,
TimeUnit.SECONDS).until(() -> repository.getKeys().isEmpty());
+
+ releaseDownstream.countDown();
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testLeftoverTimeoutDoesNotCompleteNewGroupWithoutTimeout()
throws Exception {
+ MemoryAggregationRepository repository = new
MemoryAggregationRepository(true);
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"), new
BodyInAggregatingStrategy()).aggregationRepository(repository)
+ .optimisticLocking().completionSize(2)
+
.completionTimeout(header("timeout")).completionTimeoutCheckerInterval(10)
+ .timeoutCheckerExecutorService(timeoutChecker)
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("a+b", "c+d");
+
mock.message(0).exchangeProperty(Exchange.AGGREGATED_COMPLETED_BY).isEqualTo("size");
+
mock.message(1).exchangeProperty(Exchange.AGGREGATED_COMPLETED_BY).isEqualTo("size");
+
+ // the group [a, b] has a completion timeout (from a), but is
completed by size
+ template.sendBodyAndHeaders("direct:start", "a", Map.of("id", "1",
"timeout", 1));
+ send("b");
+ // c starts a new group for the same key, whose completion timeout
expression returns no timeout
+ send("c");
+
+ // let the timeout checker run twice, so the timeout of a has expired
and been processed
+ timeoutGate.countDown();
+ await().atMost(10, TimeUnit.SECONDS).until(() ->
timeoutCheckerRuns.get() >= 2);
+ assertEquals(Set.of("1"), repository.getKeys(), "The group [c] has no
completion timeout and should not be completed");
+
+ send("d");
+ assertMockEndpointsSatisfied();
+ }
+
+ @Test
+ public void testGroupKeepsTimeoutWhenLaterExchangeHasNoTimeout() throws
Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"), new
BodyInLatestStrategy()).optimisticLocking().completionSize(3)
+
.completionTimeout(header("timeout")).completionTimeoutCheckerInterval(10)
+ .timeoutCheckerExecutorService(timeoutChecker)
+ .to("mock:result");
+ }
+ });
+ context.start();
+
+ MockEndpoint mock = getMockEndpoint("mock:result");
+ mock.expectedBodiesReceived("a+b");
+
mock.message(0).exchangeProperty(Exchange.AGGREGATED_COMPLETED_BY).isEqualTo("timeout");
+
+ // the group has the completion timeout of a, and the stored exchange
is b, which has no timeout
+ template.sendBodyAndHeaders("direct:start", "a", Map.of("id", "1",
"timeout", 1));
+ send("b");
+
+ timeoutGate.countDown();
+ assertMockEndpointsSatisfied();
+ }
+
+ private void addRoute(AggregationRepository repository) throws Exception {
+ context.addRoutes(new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"), new
BodyInAggregatingStrategy()).aggregationRepository(repository)
+ .optimisticLocking().completionSize(2)
+
.completionTimeout(10).completionTimeoutCheckerInterval(10)
+ .timeoutCheckerExecutorService(timeoutChecker)
+ .process(exchange -> {
+ if
("m2+m3".equals(exchange.getIn().getBody(String.class))) {
+ awaitLatch(releaseDownstream);
+ }
+ })
+ .to("mock:result");
+ }
+ });
+ context.start();
+ }
+
+ private void send(String body) {
+ template.sendBodyAndHeader("direct:start", body, "id", "1");
+ }
+
+ private void pauseIfThread(String name) {
+ if (Thread.currentThread().getName().equals(name) &&
pause.compareAndSet(true, false)) {
+ paused.countDown();
+ awaitLatch(release);
+ }
+ }
+
+ private static boolean awaitLatch(CountDownLatch latch) {
+ try {
+ return latch.await(20, TimeUnit.SECONDS);
+ } catch (InterruptedException e) {
+ Thread.currentThread().interrupt();
+ return false;
+ }
+ }
+
+ private static class BodyInAggregatingStrategy implements
AggregationStrategy {
+ @Override
+ public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
+ if (oldExchange == null) {
+ return newExchange;
+ }
+ oldExchange.getIn().setBody(
+ oldExchange.getIn().getBody(String.class) + "+" +
newExchange.getIn().getBody(String.class));
+ return oldExchange;
+ }
+ }
+
+ /**
+ * Aggregates the bodies into the latest exchange.
+ */
+ private static class BodyInLatestStrategy implements AggregationStrategy {
+ @Override
+ public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
+ if (oldExchange != null) {
+ newExchange.getIn().setBody(
+ oldExchange.getIn().getBody(String.class) + "+" +
newExchange.getIn().getBody(String.class));
+ }
+ return newExchange;
+ }
+ }
+
+ /**
+ * An optimistic locking repository which stores a copy of the exchange
(so an exchange keeps its exchange id, like
+ * the JDBC repository), with a version per key.
+ */
+ private final class IdPreservingRepository extends ServiceSupport
implements OptimisticLockingAggregationRepository {
+
+ private static final String VERSION = "IdPreservingRepositoryVersion";
+
+ private final Map<String, DefaultExchangeHolder> groups = new
HashMap<>();
+ private final Map<String, Long> versions = new HashMap<>();
+ private final AtomicLong version = new AtomicLong();
+
+ @Override
+ public Exchange add(CamelContext camelContext, String key, Exchange
oldExchange, Exchange newExchange) {
+ pauseIfThread("producer-m1");
+ synchronized (this) {
+ Long current = versions.get(key);
+ Long expected = oldExchange == null ? null :
oldExchange.getProperty(VERSION, Long.class);
+ if (current == null ? expected != null :
!current.equals(expected)) {
+ throw new OptimisticLockingException();
+ }
+ long next = version.incrementAndGet();
+ newExchange.setProperty(VERSION, next);
+ groups.put(key, DefaultExchangeHolder.marshal(newExchange,
true));
+ versions.put(key, next);
+ return oldExchange;
+ }
+ }
+
+ @Override
+ public Exchange add(CamelContext camelContext, String key, Exchange
exchange) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public synchronized Exchange get(CamelContext camelContext, String
key) {
+ DefaultExchangeHolder holder = groups.get(key);
+ if (holder == null) {
+ return null;
+ }
+ Exchange answer = new DefaultExchange(camelContext);
+ DefaultExchangeHolder.unmarshal(answer, holder);
+ answer.setProperty(VERSION, versions.get(key));
+ return answer;
+ }
+
+ @Override
+ public synchronized void remove(CamelContext camelContext, String key,
Exchange exchange) {
+ Long expected = exchange.getProperty(VERSION, Long.class);
+ if (expected == null || !expected.equals(versions.get(key))) {
+ throw new OptimisticLockingException();
+ }
+ groups.remove(key);
+ versions.remove(key);
+ }
+
+ @Override
+ public void confirm(CamelContext camelContext, String exchangeId) {
+ // noop
+ }
+
+ @Override
+ public synchronized Set<String> getKeys() {
+ return Set.copyOf(groups.keySet());
+ }
+ }
+}
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 202fee19bf36..dfb0e42ac18a 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
@@ -162,6 +162,18 @@ redeliveries and which fails 4 times before the exception
changes leaves the new
`onException(IllegalArgumentException.class).maximumRedeliveries(2)` is
already exhausted and the message goes to
the dead letter channel on the next failure.
+=== Aggregate EIP - completion timeout with optimistic locking
+
+When the aggregator uses optimistic locking together with
`completionTimeoutExpression`, a group now records its
+completion timeout in the `CamelAggregatedTimeout` exchange property of the
aggregated exchange stored in the
+aggregation repository, and a group only completes by timeout when it has this
property. This prevents the timeout of
+an already completed group from completing a new group for the same
correlation key.
+
+A group that was persisted in the aggregation repository before the upgrade
does not have this property, so it does
+not complete by timeout until another exchange that has a completion timeout
arrives for the group. A custom
+`AggregationRepository` that does not keep exchange properties never completes
a group by timeout. This only affects
+optimistic locking combined with `completionTimeoutExpression`.
+
=== Weighted Load Balancer EIP
The distribution ratios of the weighted load balancer are now validated when
the route starts. A negative ratio,