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

Reply via email to