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 021b0e1a5883 CAMEL-25153: camel-caffeine, camel-ehcache - keep
completed exchanges for recovery in the aggregation repositories (#27108)
021b0e1a5883 is described below
commit 021b0e1a588343d77bd61bc98badbbd0f8c6b989
Author: allthingssecurity <[email protected]>
AuthorDate: Wed Sep 30 16:28:03 2026 +0530
CAMEL-25153: camel-caffeine, camel-ehcache - keep completed exchanges for
recovery in the aggregation repositories (#27108)
Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
.../aggregate/CaffeineAggregationRepository.java | 52 ++++++++--
...CaffeineAggregationRepositoryOperationTest.java | 107 +++++++++++++++++----
.../CaffeineAggregationRepositoryRecoverTest.java | 103 ++++++++++++++++++++
.../aggregate/EhcacheAggregationRepository.java | 54 ++++++++++-
.../EhcacheAggregationRepositoryOperationTest.java | 107 +++++++++++++++++----
.../EhcacheAggregationRepositoryRecoverTest.java | 105 ++++++++++++++++++++
.../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc | 17 ++++
7 files changed, 493 insertions(+), 52 deletions(-)
diff --git
a/components/camel-caffeine/src/main/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepository.java
b/components/camel-caffeine/src/main/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepository.java
index 3f8d873fb43d..eb7b6f200885 100644
---
a/components/camel-caffeine/src/main/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepository.java
+++
b/components/camel-caffeine/src/main/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepository.java
@@ -19,6 +19,7 @@ package
org.apache.camel.component.caffeine.processor.aggregate;
import java.util.Collections;
import java.util.Set;
import java.util.concurrent.TimeUnit;
+import java.util.stream.Collectors;
import com.github.benmanes.caffeine.cache.Cache;
import com.github.benmanes.caffeine.cache.Caffeine;
@@ -43,6 +44,12 @@ public class CaffeineAggregationRepository extends
ServiceSupport implements Rec
private static final Logger LOG =
LoggerFactory.getLogger(CaffeineAggregationRepository.class);
+ /**
+ * Prefix of the keys under which completed exchanges are kept for
recovery. Recovery entries are keyed by exchange
+ * id, aggregations in progress by correlation key, and the prefix keeps
the two apart in the same cache.
+ */
+ private static final String RECOVERY_KEY_PREFIX = "camel-recovery:";
+
private CamelContext camelContext;
private Cache<String, DefaultExchangeHolder> cache;
@@ -149,33 +156,64 @@ public class CaffeineAggregationRepository extends
ServiceSupport implements Rec
public void remove(CamelContext camelContext, String key, Exchange
exchange) {
LOG.trace("Removing an exchange with ID {} for key {}",
exchange.getExchangeId(), key);
cache.invalidate(key);
+
+ if (useRecovery) {
+ // the aggregation is complete but the exchange has not been
processed yet, so keep a copy that recovery
+ // can pick up if the processing never confirms it (the given
exchange, as the one in the cache may not
+ // contain the exchange that completed the aggregation)
+ LOG.trace("Putting an exchange with ID {} into the recovery
store", exchange.getExchangeId());
+ cache.put(recoveryKey(exchange.getExchangeId()),
+ DefaultExchangeHolder.marshal(exchange, true,
allowSerializedHeaders));
+ }
}
@Override
public void confirm(CamelContext camelContext, String exchangeId) {
LOG.trace("Confirming an exchange with ID {}.", exchangeId);
- cache.invalidate(exchangeId);
+ if (useRecovery) {
+ cache.invalidate(recoveryKey(exchangeId));
+ }
}
@Override
public Set<String> getKeys() {
- Set<String> keys = cache.asMap().keySet();
-
- return Collections.unmodifiableSet(keys);
+ return cache.asMap().keySet().stream()
+ .filter(key -> !isRecoveryKey(key))
+ .collect(Collectors.collectingAndThen(Collectors.toSet(),
Collections::unmodifiableSet));
}
@Override
public Set<String> scan(CamelContext camelContext) {
+ if (!useRecovery) {
+ LOG.debug("Recovery is disabled on the repository of {} context,
nothing to scan", camelContext.getName());
+ return Collections.emptySet();
+ }
+
LOG.trace("Scanning for exchanges to recover in {} context",
camelContext.getName());
- Set<String> scanned = Collections.unmodifiableSet(getKeys());
- LOG.trace("Found {} keys for exchanges to recover in {} context",
scanned.size(), camelContext.getName());
+ Set<String> scanned = cache.asMap().keySet().stream()
+ .filter(CaffeineAggregationRepository::isRecoveryKey)
+ .map(CaffeineAggregationRepository::exchangeIdOf)
+ .collect(Collectors.collectingAndThen(Collectors.toSet(),
Collections::unmodifiableSet));
+ LOG.trace("Found {} exchanges to recover in {} context",
scanned.size(), camelContext.getName());
return scanned;
}
@Override
public Exchange recover(CamelContext camelContext, String exchangeId) {
LOG.trace("Recovering an Exchange with ID {}.", exchangeId);
- return useRecovery ? unmarshallExchange(camelContext,
cache.getIfPresent(exchangeId)) : null;
+ return useRecovery ? unmarshallExchange(camelContext,
cache.getIfPresent(recoveryKey(exchangeId))) : null;
+ }
+
+ private static String recoveryKey(String exchangeId) {
+ return RECOVERY_KEY_PREFIX + exchangeId;
+ }
+
+ private static boolean isRecoveryKey(String key) {
+ return key.startsWith(RECOVERY_KEY_PREFIX);
+ }
+
+ private static String exchangeIdOf(String recoveryKey) {
+ return recoveryKey.substring(RECOVERY_KEY_PREFIX.length());
}
@Override
diff --git
a/components/camel-caffeine/src/test/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepositoryOperationTest.java
b/components/camel-caffeine/src/test/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepositoryOperationTest.java
index e62cf8aec4aa..29ecd1356502 100644
---
a/components/camel-caffeine/src/test/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepositoryOperationTest.java
+++
b/components/camel-caffeine/src/test/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepositoryOperationTest.java
@@ -133,19 +133,20 @@ public class CaffeineAggregationRepositoryOperationTest
extends CamelTestSupport
@Test
void testConfirmExist() {
// Given
- for (int i = 1; i < 4; i++) {
- String key = "Confirm_" + i;
- Exchange exchange = new DefaultExchange(context());
- exchange.setExchangeId("Exchange_" + i);
- aggregationRepository.add(context(), key, exchange);
- assertTrue(exists(key));
- }
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange_Confirm");
+ aggregationRepository.add(context(), "Confirm_1", exchange);
+ // completing the aggregation moves the exchange into the recovery
store
+ aggregationRepository.remove(context(), "Confirm_1", exchange);
+ assertFalse(exists("Confirm_1"));
+ assertNotNull(aggregationRepository.recover(context(),
"Exchange_Confirm"));
+
// When
- aggregationRepository.confirm(context(), "Confirm_2");
+ aggregationRepository.confirm(context(), "Exchange_Confirm");
+
// Then
- assertTrue(exists("Confirm_1"));
- assertFalse(exists("Confirm_2"));
- assertTrue(exists("Confirm_3"));
+ assertNull(aggregationRepository.recover(context(),
"Exchange_Confirm"));
+ assertTrue(aggregationRepository.scan(context()).isEmpty());
}
@Test
@@ -178,26 +179,92 @@ public class CaffeineAggregationRepositoryOperationTest
extends CamelTestSupport
@Test
void testScan() {
// Given
+ String[] keys = { "Scan1", "Scan2", "Scan3" };
+ addExchanges(keys);
+ // the first two aggregations are completed, the third is still in
progress
+ for (int i = 0; i < 2; i++) {
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-" + keys[i]);
+ aggregationRepository.remove(context(), keys[i], exchange);
+ }
+
+ // When
+ Set<String> exchangeIdSet = aggregationRepository.scan(context());
+
+ // Then - the scan reports the exchange ids to recover, not the
correlation keys still aggregating
+ assertEquals(Set.of("Exchange-Scan1", "Exchange-Scan2"),
exchangeIdSet);
+ }
+
+ @Test
+ void testScanWithoutRecovery() {
+ // Given
+ aggregationRepository.setUseRecovery(false);
String[] keys = { "Scan1", "Scan2" };
addExchanges(keys);
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-Scan1");
+ aggregationRepository.remove(context(), "Scan1", exchange);
+
// When
Set<String> exchangeIdSet = aggregationRepository.scan(context());
+
// Then
- for (String key : keys) {
- assertTrue(exchangeIdSet.contains(key));
- }
+ assertTrue(exchangeIdSet.isEmpty());
+ assertNull(aggregationRepository.recover(context(), "Exchange-Scan1"));
+ assertEquals(Set.of("Scan2"), aggregationRepository.getKeys());
}
@Test
void testRecover() {
// Given
- String[] keys = { "Recover1", "Recover2" };
- addExchanges(keys);
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-Recover1");
+ exchange.getIn().setBody("Hello");
+ aggregationRepository.add(context(), "Recover1", exchange);
+ // the exchange that completed the aggregation has been aggregated
after the last add
+ exchange.getIn().setBody("Hello World");
+ aggregationRepository.remove(context(), "Recover1", exchange);
+
// When
- Exchange exchange2 = aggregationRepository.recover(context(),
"Recover2");
- Exchange exchange3 = aggregationRepository.recover(context(),
"Recover3");
+ Exchange recovered = aggregationRepository.recover(context(),
"Exchange-Recover1");
+ Exchange unknown = aggregationRepository.recover(context(),
"Exchange-Recover2");
+ Exchange inProgress = aggregationRepository.recover(context(),
"Recover1");
+
// Then
- assertNotNull(exchange2);
- assertNull(exchange3);
+ assertNotNull(recovered);
+ assertEquals("Exchange-Recover1", recovered.getExchangeId());
+ assertEquals("Hello World", recovered.getIn().getBody());
+ assertNull(unknown);
+ assertNull(inProgress);
+ }
+
+ @Test
+ void testRecoverDoesNotReturnAggregationInProgress() {
+ // Given
+ addExchanges("Recover1");
+
+ // When
+ Set<String> exchangeIdSet = aggregationRepository.scan(context());
+ Exchange recovered = aggregationRepository.recover(context(),
"Recover1");
+
+ // Then
+ assertTrue(exchangeIdSet.isEmpty());
+ assertNull(recovered);
+ }
+
+ @Test
+ void testGetKeysIgnoresExchangesToRecover() {
+ // Given
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-Keys1");
+ aggregationRepository.add(context(), "Keys1", exchange);
+ aggregationRepository.add(context(), "Keys2", exchange);
+ aggregationRepository.remove(context(), "Keys1", exchange);
+
+ // When
+ Set<String> keys = aggregationRepository.getKeys();
+
+ // Then - only the aggregation still in progress is reported
+ assertEquals(Set.of("Keys2"), keys);
}
}
diff --git
a/components/camel-caffeine/src/test/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepositoryRecoverTest.java
b/components/camel-caffeine/src/test/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepositoryRecoverTest.java
new file mode 100644
index 000000000000..e84d9ab06df9
--- /dev/null
+++
b/components/camel-caffeine/src/test/java/org/apache/camel/component/caffeine/processor/aggregate/CaffeineAggregationRepositoryRecoverTest.java
@@ -0,0 +1,103 @@
+/*
+ * 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.caffeine.processor.aggregate;
+
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.test.junit6.CamelTestSupport;
+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.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The recover task of the Aggregate EIP must only recover completed exchanges
that were not confirmed, and never send
+ * an aggregation that is still in progress.
+ */
+public class CaffeineAggregationRepositoryRecoverTest extends CamelTestSupport
{
+
+ private final AtomicInteger scans = new AtomicInteger();
+ private final AtomicInteger failures = new AtomicInteger();
+ private CaffeineAggregationRepository repository;
+
+ @Test
+ void testRecoverCompletedExchangeOnly() throws Exception {
+ MockEndpoint mock = getMockEndpoint("mock:aggregated");
+ // the completed aggregation fails once and is then recovered, the
aggregation in progress is not sent
+ mock.expectedBodiesReceived("a+b", "a+b");
+
+ template.sendBodyAndHeader("direct:start", "c", "id", "inProgress");
+ template.sendBodyAndHeader("direct:start", "a", "id", "completed");
+ template.sendBodyAndHeader("direct:start", "b", "id", "completed");
+
+ mock.assertIsSatisfied();
+
assertNull(mock.getReceivedExchanges().get(0).getIn().getHeader(Exchange.REDELIVERED));
+ assertEquals(Boolean.TRUE,
mock.getReceivedExchanges().get(1).getIn().getHeader(Exchange.REDELIVERED));
+
+ // let the recover task run a few more times: nothing else is sent,
and the recovered exchange was confirmed
+ int scanned = scans.get();
+ await().atMost(10, TimeUnit.SECONDS).until(() -> scans.get() >=
scanned + 3);
+ assertEquals(2, mock.getReceivedCounter());
+ assertTrue(repository.scan(context).isEmpty());
+ assertEquals(Set.of("inProgress"), repository.getKeys());
+ }
+
+ @Override
+ protected RoutesBuilder createRouteBuilder() {
+ repository = new CaffeineAggregationRepository() {
+ @Override
+ public Set<String> scan(CamelContext camelContext) {
+ scans.incrementAndGet();
+ return super.scan(camelContext);
+ }
+ };
+ repository.setRecoveryInterval(100);
+
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"), (oldExchange, newExchange) ->
{
+ if (oldExchange == null) {
+ return newExchange;
+ }
+ String body =
oldExchange.getIn().getBody(String.class) + "+"
+ +
newExchange.getIn().getBody(String.class);
+ oldExchange.getIn().setBody(body);
+ return oldExchange;
+ })
+ .aggregationRepository(repository)
+ .completionSize(2)
+ .to("mock:aggregated")
+ .process(exchange -> {
+ if (failures.getAndIncrement() == 0) {
+ throw new IllegalStateException("Forced
failure after the aggregation");
+ }
+ });
+ }
+ };
+ }
+}
diff --git
a/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepository.java
b/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepository.java
index de9fceb57368..f741121274a7 100644
---
a/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepository.java
+++
b/components/camel-ehcache/src/main/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepository.java
@@ -45,6 +45,12 @@ public class EhcacheAggregationRepository extends
ServiceSupport implements Reco
private static final Logger LOG =
LoggerFactory.getLogger(EhcacheAggregationRepository.class);
+ /**
+ * Prefix of the keys under which completed exchanges are kept for
recovery. Recovery entries are keyed by exchange
+ * id, aggregations in progress by correlation key, and the prefix keeps
the two apart in the same cache.
+ */
+ private static final String RECOVERY_KEY_PREFIX = "camel-recovery:";
+
private CamelContext camelContext;
private CacheManager cacheManager;
@Metadata(description = "Name of cache", required = true)
@@ -173,34 +179,72 @@ public class EhcacheAggregationRepository extends
ServiceSupport implements Reco
public void remove(CamelContext camelContext, String key, Exchange
exchange) {
LOG.trace("Removing an exchange with ID {} for key {}",
exchange.getExchangeId(), key);
cache.remove(key);
+
+ if (useRecovery) {
+ // the aggregation is complete but the exchange has not been
processed yet, so keep a copy that recovery
+ // can pick up if the processing never confirms it (the given
exchange, as the one in the cache may not
+ // contain the exchange that completed the aggregation)
+ LOG.trace("Putting an exchange with ID {} into the recovery
store", exchange.getExchangeId());
+ cache.put(recoveryKey(exchange.getExchangeId()),
+ DefaultExchangeHolder.marshal(exchange, true,
allowSerializedHeaders));
+ }
}
@Override
public void confirm(CamelContext camelContext, String exchangeId) {
LOG.trace("Confirming an exchange with ID {}.", exchangeId);
- cache.remove(exchangeId);
+ if (useRecovery) {
+ cache.remove(recoveryKey(exchangeId));
+ }
}
@Override
public Set<String> getKeys() {
Set<String> keys = new HashSet<>();
- cache.forEach(e -> keys.add(e.getKey()));
+ cache.forEach(e -> {
+ if (!isRecoveryKey(e.getKey())) {
+ keys.add(e.getKey());
+ }
+ });
return Collections.unmodifiableSet(keys);
}
@Override
public Set<String> scan(CamelContext camelContext) {
+ if (!useRecovery) {
+ LOG.debug("Recovery is disabled on the repository of {} context,
nothing to scan", camelContext.getName());
+ return Collections.emptySet();
+ }
+
LOG.trace("Scanning for exchanges to recover in {} context",
camelContext.getName());
- Set<String> scanned = Collections.unmodifiableSet(getKeys());
- LOG.trace("Found {} keys for exchanges to recover in {} context",
scanned.size(), camelContext.getName());
+ Set<String> exchangeIds = new HashSet<>();
+ cache.forEach(e -> {
+ if (isRecoveryKey(e.getKey())) {
+ exchangeIds.add(exchangeIdOf(e.getKey()));
+ }
+ });
+ Set<String> scanned = Collections.unmodifiableSet(exchangeIds);
+ LOG.trace("Found {} exchanges to recover in {} context",
scanned.size(), camelContext.getName());
return scanned;
}
@Override
public Exchange recover(CamelContext camelContext, String exchangeId) {
LOG.trace("Recovering an Exchange with ID {}.", exchangeId);
- return useRecovery ? unmarshallExchange(camelContext,
cache.get(exchangeId)) : null;
+ return useRecovery ? unmarshallExchange(camelContext,
cache.get(recoveryKey(exchangeId))) : null;
+ }
+
+ private static String recoveryKey(String exchangeId) {
+ return RECOVERY_KEY_PREFIX + exchangeId;
+ }
+
+ private static boolean isRecoveryKey(String key) {
+ return key.startsWith(RECOVERY_KEY_PREFIX);
+ }
+
+ private static String exchangeIdOf(String recoveryKey) {
+ return recoveryKey.substring(RECOVERY_KEY_PREFIX.length());
}
@Override
diff --git
a/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepositoryOperationTest.java
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepositoryOperationTest.java
index 1bb7f779cf61..946ff65b3342 100644
---
a/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepositoryOperationTest.java
+++
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepositoryOperationTest.java
@@ -133,19 +133,20 @@ public class EhcacheAggregationRepositoryOperationTest
extends EhcacheTestSuppor
@Test
void testConfirmExist() {
// Given
- for (int i = 1; i < 4; i++) {
- String key = "Confirm_" + i;
- Exchange exchange = new DefaultExchange(context());
- exchange.setExchangeId("Exchange_" + i);
- aggregationRepository.add(context(), key, exchange);
- assertTrue(exists(key));
- }
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange_Confirm");
+ aggregationRepository.add(context(), "Confirm_1", exchange);
+ // completing the aggregation moves the exchange into the recovery
store
+ aggregationRepository.remove(context(), "Confirm_1", exchange);
+ assertFalse(exists("Confirm_1"));
+ assertNotNull(aggregationRepository.recover(context(),
"Exchange_Confirm"));
+
// When
- aggregationRepository.confirm(context(), "Confirm_2");
+ aggregationRepository.confirm(context(), "Exchange_Confirm");
+
// Then
- assertTrue(exists("Confirm_1"));
- assertFalse(exists("Confirm_2"));
- assertTrue(exists("Confirm_3"));
+ assertNull(aggregationRepository.recover(context(),
"Exchange_Confirm"));
+ assertTrue(aggregationRepository.scan(context()).isEmpty());
}
@Test
@@ -178,26 +179,92 @@ public class EhcacheAggregationRepositoryOperationTest
extends EhcacheTestSuppor
@Test
void testScan() {
// Given
+ String[] keys = { "Scan1", "Scan2", "Scan3" };
+ addExchanges(keys);
+ // the first two aggregations are completed, the third is still in
progress
+ for (int i = 0; i < 2; i++) {
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-" + keys[i]);
+ aggregationRepository.remove(context(), keys[i], exchange);
+ }
+
+ // When
+ Set<String> exchangeIdSet = aggregationRepository.scan(context());
+
+ // Then - the scan reports the exchange ids to recover, not the
correlation keys still aggregating
+ assertEquals(Set.of("Exchange-Scan1", "Exchange-Scan2"),
exchangeIdSet);
+ }
+
+ @Test
+ void testScanWithoutRecovery() {
+ // Given
+ aggregationRepository.setUseRecovery(false);
String[] keys = { "Scan1", "Scan2" };
addExchanges(keys);
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-Scan1");
+ aggregationRepository.remove(context(), "Scan1", exchange);
+
// When
Set<String> exchangeIdSet = aggregationRepository.scan(context());
+
// Then
- for (String key : keys) {
- assertTrue(exchangeIdSet.contains(key));
- }
+ assertTrue(exchangeIdSet.isEmpty());
+ assertNull(aggregationRepository.recover(context(), "Exchange-Scan1"));
+ assertEquals(Set.of("Scan2"), aggregationRepository.getKeys());
}
@Test
void testRecover() {
// Given
- String[] keys = { "Recover1", "Recover2" };
- addExchanges(keys);
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-Recover1");
+ exchange.getIn().setBody("Hello");
+ aggregationRepository.add(context(), "Recover1", exchange);
+ // the exchange that completed the aggregation has been aggregated
after the last add
+ exchange.getIn().setBody("Hello World");
+ aggregationRepository.remove(context(), "Recover1", exchange);
+
// When
- Exchange exchange2 = aggregationRepository.recover(context(),
"Recover2");
- Exchange exchange3 = aggregationRepository.recover(context(),
"Recover3");
+ Exchange recovered = aggregationRepository.recover(context(),
"Exchange-Recover1");
+ Exchange unknown = aggregationRepository.recover(context(),
"Exchange-Recover2");
+ Exchange inProgress = aggregationRepository.recover(context(),
"Recover1");
+
// Then
- assertNotNull(exchange2);
- assertNull(exchange3);
+ assertNotNull(recovered);
+ assertEquals("Exchange-Recover1", recovered.getExchangeId());
+ assertEquals("Hello World", recovered.getIn().getBody());
+ assertNull(unknown);
+ assertNull(inProgress);
+ }
+
+ @Test
+ void testRecoverDoesNotReturnAggregationInProgress() {
+ // Given
+ addExchanges("Recover1");
+
+ // When
+ Set<String> exchangeIdSet = aggregationRepository.scan(context());
+ Exchange recovered = aggregationRepository.recover(context(),
"Recover1");
+
+ // Then
+ assertTrue(exchangeIdSet.isEmpty());
+ assertNull(recovered);
+ }
+
+ @Test
+ void testGetKeysIgnoresExchangesToRecover() {
+ // Given
+ Exchange exchange = new DefaultExchange(context());
+ exchange.setExchangeId("Exchange-Keys1");
+ aggregationRepository.add(context(), "Keys1", exchange);
+ aggregationRepository.add(context(), "Keys2", exchange);
+ aggregationRepository.remove(context(), "Keys1", exchange);
+
+ // When
+ Set<String> keys = aggregationRepository.getKeys();
+
+ // Then - only the aggregation still in progress is reported
+ assertEquals(Set.of("Keys2"), keys);
}
}
diff --git
a/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepositoryRecoverTest.java
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepositoryRecoverTest.java
new file mode 100644
index 000000000000..6a66cc5c4e03
--- /dev/null
+++
b/components/camel-ehcache/src/test/java/org/apache/camel/component/ehcache/processor/aggregate/EhcacheAggregationRepositoryRecoverTest.java
@@ -0,0 +1,105 @@
+/*
+ * 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.ehcache.processor.aggregate;
+
+import java.util.Set;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.ehcache.EhcacheTestSupport;
+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.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The recover task of the Aggregate EIP must only recover completed exchanges
that were not confirmed, and never send
+ * an aggregation that is still in progress.
+ */
+public class EhcacheAggregationRepositoryRecoverTest extends
EhcacheTestSupport {
+
+ private final AtomicInteger scans = new AtomicInteger();
+ private final AtomicInteger failures = new AtomicInteger();
+ private EhcacheAggregationRepository repository;
+
+ @Test
+ void testRecoverCompletedExchangeOnly() throws Exception {
+ MockEndpoint mock = getMockEndpoint("mock:aggregated");
+ // the completed aggregation fails once and is then recovered, the
aggregation in progress is not sent
+ mock.expectedBodiesReceived("a+b", "a+b");
+
+ template.sendBodyAndHeader("direct:start", "c", "id", "inProgress");
+ template.sendBodyAndHeader("direct:start", "a", "id", "completed");
+ template.sendBodyAndHeader("direct:start", "b", "id", "completed");
+
+ mock.assertIsSatisfied();
+
assertNull(mock.getReceivedExchanges().get(0).getIn().getHeader(Exchange.REDELIVERED));
+ assertEquals(Boolean.TRUE,
mock.getReceivedExchanges().get(1).getIn().getHeader(Exchange.REDELIVERED));
+
+ // let the recover task run a few more times: nothing else is sent,
and the recovered exchange was confirmed
+ int scanned = scans.get();
+ await().atMost(10, TimeUnit.SECONDS).until(() -> scans.get() >=
scanned + 3);
+ assertEquals(2, mock.getReceivedCounter());
+ assertTrue(repository.scan(context).isEmpty());
+ assertEquals(Set.of("inProgress"), repository.getKeys());
+ }
+
+ @Override
+ protected RoutesBuilder createRouteBuilder() {
+ repository = new EhcacheAggregationRepository() {
+ @Override
+ public Set<String> scan(CamelContext camelContext) {
+ scans.incrementAndGet();
+ return super.scan(camelContext);
+ }
+ };
+ repository.setCache(getAggregateCache());
+ repository.setCacheName(AGGREGATE_TEST_CACHE_NAME);
+ repository.setRecoveryInterval(100);
+
+ return new RouteBuilder() {
+ @Override
+ public void configure() {
+ from("direct:start")
+ .aggregate(header("id"), (oldExchange, newExchange) ->
{
+ if (oldExchange == null) {
+ return newExchange;
+ }
+ String body =
oldExchange.getIn().getBody(String.class) + "+"
+ +
newExchange.getIn().getBody(String.class);
+ oldExchange.getIn().setBody(body);
+ return oldExchange;
+ })
+ .aggregationRepository(repository)
+ .completionSize(2)
+ .to("mock:aggregated")
+ .process(exchange -> {
+ if (failures.getAndIncrement() == 0) {
+ throw new IllegalStateException("Forced
failure after the aggregation");
+ }
+ });
+ }
+ };
+ }
+}
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 1040b2966262..b62a9a431e19 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
@@ -3003,6 +3003,23 @@ default, stop seeing in-progress aggregations
re-delivered, and start seeing gen
now also holds one entry per completed and not yet confirmed exchange; those
entries are removed on
confirmation.
+=== camel-caffeine, camel-ehcache - the aggregation repositories keep
completed exchanges for recovery
+
+`CaffeineAggregationRepository` and `EhcacheAggregationRepository` implement
`RecoverableAggregationRepository`,
+but like the Infinispan repository before 4.23 they had no recovery store:
`remove` deleted the completed exchange,
+`confirm` removed the exchange id from a cache keyed by correlation key, and
`scan` returned the correlation keys of
+the aggregations still in progress. The recovery task therefore sent
aggregations that were still in progress
+(marked `CamelRedelivered`, and to the dead letter channel once
`maximumRedeliveries` was reached), while exchanges
+that failed after completion could never be recovered.
+
+A completed exchange is now kept in the same cache under a
`camel-recovery:<exchange id>` key until it is
+confirmed, which is what `scan` reports and `recover` reads. `getKeys` reports
only the aggregations in progress.
+
+Routes that set `useRecovery=false` are unaffected. Routes that left recovery
enabled, which is the default, no
+longer see aggregations in progress sent by the recovery task, and a completed
exchange whose processing failed is
+now recovered. The cache now also holds one entry per completed and not yet
confirmed exchange; those entries are
+removed on confirmation. A cache with a size limit or an expiry applies it to
these entries as well.
+
=== camel-qdrant - the PayloadSelector header is now honoured
`QdrantHeaders.PAYLOAD_SELECTOR` (`CamelQdrantPointsPayloadSelector`) was
declared and advertised in