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 1de5a7c0cb52 CAMEL-25210: camel-cassandraql - keep completed exchanges 
for recovery in the aggregation repository (#27159)
1de5a7c0cb52 is described below

commit 1de5a7c0cb52d87515728f263a7faca2e798b907
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 21:10:55 2026 +0530

    CAMEL-25210: camel-cassandraql - keep completed exchanges for recovery in 
the aggregation repository (#27159)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../cassandra/CassandraAggregationRepository.java  |  88 ++++-----
 .../CassandraAggregationRepositoryIT.java          |  27 ++-
 ...CassandraAggregationRepositoryRecoveryTest.java | 202 +++++++++++++++++++++
 .../NamedCassandraAggregationRepositoryIT.java     |  27 ++-
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  18 ++
 5 files changed, 313 insertions(+), 49 deletions(-)

diff --git 
a/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
 
b/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
index f3c2b9a7d913..92662ae9d252 100644
--- 
a/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
+++ 
b/components/camel-cassandraql/src/main/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepository.java
@@ -18,6 +18,7 @@ package org.apache.camel.processor.aggregate.cassandra;
 
 import java.io.IOException;
 import java.nio.ByteBuffer;
+import java.util.Collections;
 import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
@@ -42,7 +43,6 @@ import 
org.apache.camel.utils.cassandra.CassandraSessionHolder;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import static com.datastax.oss.driver.api.querybuilder.QueryBuilder.bindMarker;
 import static org.apache.camel.utils.cassandra.CassandraUtils.append;
 import static 
org.apache.camel.utils.cassandra.CassandraUtils.applyConsistencyLevel;
 import static org.apache.camel.utils.cassandra.CassandraUtils.concat;
@@ -64,6 +64,11 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
 
     private static final Logger LOGGER = 
LoggerFactory.getLogger(CassandraAggregationRepository.class);
 
+    /**
+     * Prefix of the aggregation key under which a completed exchange is kept, 
by its exchange id, until it is confirmed
+     */
+    private static final String RECOVERY_KEY_PREFIX = "camel-recovery:";
+
     private final CassandraCamelCodec exchangeCodec = new 
CassandraCamelCodec();
 
     @Metadata(description = "Cassandra session", required = true)
@@ -97,10 +102,6 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
      * Prepared statement used to get exchangeIds and exchange ids
      */
     private PreparedStatement selectKeyIdStatement;
-    /**
-     * Prepared statement used to delete with key and exchange id
-     */
-    private PreparedStatement deleteIfIdStatement;
 
     @Metadata(description = "Sets the interval between recovery scans", 
defaultValue = "5000")
     private long recoveryInterval = 5000;
@@ -167,7 +168,6 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
         initSelectStatement();
         initDeleteStatement();
         initSelectKeyIdStatement();
-        initDeleteIfIdStatement();
     }
 
     @Override
@@ -190,13 +190,17 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
      */
     @Override
     public Exchange add(CamelContext camelContext, String key, Exchange 
exchange) {
+        insert(key, exchange);
+        return exchange;
+    }
+
+    private void insert(String key, Exchange exchange) {
         final Object[] idValues = getPKValues(key);
         LOGGER.debug("Inserting key {} exchange {}", idValues, exchange);
         try {
             ByteBuffer marshalledExchange = 
exchangeCodec.marshallExchange(exchange, allowSerializedHeaders);
             Object[] cqlParams = concat(idValues, new Object[] { 
exchange.getExchangeId(), marshalledExchange });
             getSession().execute(insertStatement.bind(cqlParams));
-            return exchange;
         } catch (IOException iOException) {
             throw new CassandraAggregationException("Failed to write 
exchange", exchange, iOException);
         }
@@ -236,29 +240,16 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
 
     // 
-------------------------------------------------------------------------
     // Confirm exchange in repository
-    private void initDeleteIfIdStatement() {
-        Delete delete = generateDelete(table, pkColumns, false);
-        Delete deleteIf = 
delete.ifColumn(exchangeIdColumn).isEqualTo(bindMarker());
-        SimpleStatement statement = applyConsistencyLevel(deleteIf.build(), 
writeConsistencyLevel);
-        LOGGER.debug("Generated Delete If Id {}", statement);
-        deleteIfIdStatement = getSession().prepare(statement);
-    }
 
     /**
-     * Remove exchange by Id from aggregation table.
+     * Remove the completed exchange kept for recovery.
      */
     @Override
     public void confirm(CamelContext camelContext, String exchangeId) {
-        String keyColumn = getKeyColumn();
-        LOGGER.debug("Selecting Ids");
-        List<Row> rows = selectKeyIds();
-        for (Row row : rows) {
-            if (row.getString(exchangeIdColumn).equals(exchangeId)) {
-                String key = row.getString(keyColumn);
-                Object[] cqlParams = append(getPKValues(key), exchangeId);
-                LOGGER.debug("Deleting If Id {}", cqlParams);
-                getSession().execute(deleteIfIdStatement.bind(cqlParams));
-            }
+        if (useRecovery) {
+            Object[] idValues = getPKValues(recoveryKey(exchangeId));
+            LOGGER.debug("Deleting key {}", (Object) idValues);
+            getSession().execute(deleteStatement.bind(idValues));
         }
     }
 
@@ -277,6 +268,13 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
      */
     @Override
     public void remove(CamelContext camelContext, String key, Exchange 
exchange) {
+        if (useRecovery) {
+            // the aggregation is complete but the exchange has not been 
processed yet, so keep it under its exchange
+            // id until it is confirmed, so the recover task can pick it up. 
This is the given exchange, as the stored
+            // one does not contain the exchange that completed the 
aggregation. It is written before the aggregation
+            // is deleted, so a failure in between does not lose it.
+            insert(recoveryKey(exchange.getExchangeId()), exchange);
+        }
         Object[] idValues = getPKValues(key);
         LOGGER.debug("Deleting key {}", (Object) idValues);
         getSession().execute(deleteStatement.bind(idValues));
@@ -304,7 +302,7 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
     }
 
     /**
-     * Get aggregation exchangeIds from aggregation table.
+     * Get the aggregation keys of the aggregations in progress from 
aggregation table.
      */
     @Override
     public Set<String> getKeys() {
@@ -312,7 +310,10 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
         Set<String> keys = new HashSet<>(rows.size());
         String keyColumnName = getKeyColumn();
         for (Row row : rows) {
-            keys.add(row.getString(keyColumnName));
+            String key = row.getString(keyColumnName);
+            if (!isRecoveryKey(key)) {
+                keys.add(key);
+            }
         }
         return keys;
     }
@@ -324,30 +325,35 @@ public class CassandraAggregationRepository extends 
ServiceSupport implements Re
      */
     @Override
     public Set<String> scan(CamelContext camelContext) {
+        if (!useRecovery) {
+            return Collections.emptySet();
+        }
         List<Row> rows = selectKeyIds();
-        Set<String> exchangeIds = new HashSet<>(rows.size());
+        Set<String> exchangeIds = new HashSet<>();
+        String keyColumnName = getKeyColumn();
         for (Row row : rows) {
-            exchangeIds.add(row.getString(exchangeIdColumn));
+            String key = row.getString(keyColumnName);
+            if (isRecoveryKey(key)) {
+                exchangeIds.add(key.substring(RECOVERY_KEY_PREFIX.length()));
+            }
         }
         return exchangeIds;
     }
 
     /**
-     * Get exchange by exchange ID. This is far from optimal.
+     * Get the completed exchange kept for recovery by exchange ID.
      */
     @Override
     public Exchange recover(CamelContext camelContext, String exchangeId) {
-        List<Row> rows = selectKeyIds();
-        String keyColumnName = getKeyColumn();
-        String lKey = null;
-        for (Row row : rows) {
-            String lExchangeId = row.getString(exchangeIdColumn);
-            if (lExchangeId.equals(exchangeId)) {
-                lKey = row.getString(keyColumnName);
-                break;
-            }
-        }
-        return lKey == null ? null : get(camelContext, lKey);
+        return useRecovery ? get(camelContext, recoveryKey(exchangeId)) : null;
+    }
+
+    private static String recoveryKey(String exchangeId) {
+        return RECOVERY_KEY_PREFIX + exchangeId;
+    }
+
+    private static boolean isRecoveryKey(String key) {
+        return key != null && key.startsWith(RECOVERY_KEY_PREFIX);
     }
 
     // 
-------------------------------------------------------------------------
diff --git 
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
 
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
index f08e56bffed0..d2a47bb35dcb 100644
--- 
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
+++ 
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryIT.java
@@ -144,11 +144,16 @@ public class CassandraAggregationRepositoryIT extends 
BaseCassandra {
             aggregationRepository.add(context, key, exchange);
             assertTrue(exists(key));
         }
+        Exchange completed = new DefaultExchange(context);
+        completed.setExchangeId("Exchange_2");
+        aggregationRepository.remove(context, "Confirm_2", completed);
+        assertFalse(exists("Confirm_2"));
+        assertTrue(exists("camel-recovery:Exchange_2"));
         // When
         aggregationRepository.confirm(context, "Exchange_2");
         // Then
+        assertFalse(exists("camel-recovery:Exchange_2"));
         assertTrue(exists("Confirm_1"));
-        assertFalse(exists("Confirm_2"));
         assertTrue(exists("Confirm_3"));
     }
 
@@ -171,6 +176,12 @@ public class CassandraAggregationRepositoryIT extends 
BaseCassandra {
         }
     }
 
+    private void removeExchange(String key) {
+        Exchange exchange = new DefaultExchange(context);
+        exchange.setExchangeId("Exchange-" + key);
+        aggregationRepository.remove(context, key, exchange);
+    }
+
     private void addExchanges(String... keys) {
         for (String key : keys) {
             Exchange exchange = new DefaultExchange(context);
@@ -184,12 +195,16 @@ public class CassandraAggregationRepositoryIT extends 
BaseCassandra {
         // Given
         String[] keys = { "Scan1", "Scan2" };
         addExchanges(keys);
+        // aggregations in progress are not to be recovered
+        
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan1"));
+        
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan2"));
         // When
+        removeExchange("Scan2");
         Set<String> exchangeIdSet = aggregationRepository.scan(context);
         // Then
-        for (String key : keys) {
-            assertTrue(exchangeIdSet.contains("Exchange-" + key));
-        }
+        assertTrue(exchangeIdSet.contains("Exchange-Scan2"));
+        assertFalse(exchangeIdSet.contains("Exchange-Scan1"));
+        
assertFalse(aggregationRepository.getKeys().contains("camel-recovery:Exchange-Scan2"));
     }
 
     @Test
@@ -197,11 +212,15 @@ public class CassandraAggregationRepositoryIT extends 
BaseCassandra {
         // Given
         String[] keys = { "Recover1", "Recover2" };
         addExchanges(keys);
+        removeExchange("Recover2");
         // When
+        Exchange exchange1 = aggregationRepository.recover(context, 
"Exchange-Recover1");
         Exchange exchange2 = aggregationRepository.recover(context, 
"Exchange-Recover2");
         Exchange exchange3 = aggregationRepository.recover(context, 
"Exchange-Recover3");
         // Then
+        assertNull(exchange1);
         assertNotNull(exchange2);
+        assertEquals("Exchange-Recover2", exchange2.getExchangeId());
         assertNull(exchange3);
     }
 
diff --git 
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryRecoveryTest.java
 
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryRecoveryTest.java
new file mode 100644
index 000000000000..edd5a6fdea7f
--- /dev/null
+++ 
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/CassandraAggregationRepositoryRecoveryTest.java
@@ -0,0 +1,202 @@
+/*
+ * 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.aggregate.cassandra;
+
+import java.nio.ByteBuffer;
+import java.util.ArrayList;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import com.datastax.oss.driver.api.core.CqlSession;
+import com.datastax.oss.driver.api.core.cql.BoundStatement;
+import com.datastax.oss.driver.api.core.cql.PreparedStatement;
+import com.datastax.oss.driver.api.core.cql.ResultSet;
+import com.datastax.oss.driver.api.core.cql.Row;
+import com.datastax.oss.driver.api.core.cql.SimpleStatement;
+import org.apache.camel.CamelContext;
+import org.apache.camel.Exchange;
+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;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
+
+/**
+ * The recover task of the Aggregate EIP must only see completed exchanges 
that were not confirmed: an aggregation that
+ * is still open is not recovered, and a completed aggregation whose 
processing fails is recovered with all its
+ * messages. The repository runs against an in-memory table behind a mocked 
{@link CqlSession}.
+ */
+public class CassandraAggregationRepositoryRecoveryTest extends 
CamelTestSupport {
+
+    private final AtomicInteger scans = new AtomicInteger();
+
+    @Test
+    public void testOpenAggregationIsNotRecovered() throws Exception {
+        CassandraAggregationRepository repository = createRepository();
+        addRoute(repository, "open", 3);
+
+        MockEndpoint mock = getMockEndpoint("mock:open");
+        mock.expectedMessageCount(0);
+
+        template.sendBodyAndHeader("direct:open", "a", "id", "group");
+        template.sendBodyAndHeader("direct:open", "b", "id", "group");
+
+        // let the recover task run a few times while the aggregation is open
+        int scanned = scans.get();
+        await().atMost(10, TimeUnit.SECONDS).until(() -> scans.get() >= 
scanned + 3);
+
+        mock.assertIsSatisfied();
+        assertEquals("a+b", repository.get(context, 
"group").getIn().getBody(String.class));
+    }
+
+    @Test
+    public void testCompletedAggregationIsRecoveredWithAllMessages() throws 
Exception {
+        CassandraAggregationRepository repository = createRepository();
+        addRoute(repository, "completed", 3);
+
+        MockEndpoint mock = getMockEndpoint("mock:completed");
+        // the completed aggregation fails once and is then recovered with all 
three messages
+        mock.expectedBodiesReceived("a+b+c", "a+b+c");
+
+        template.sendBodyAndHeader("direct:completed", "a", "id", "group");
+        template.sendBodyAndHeader("direct:completed", "b", "id", "group");
+        template.sendBodyAndHeader("direct:completed", "c", "id", "group");
+
+        mock.assertIsSatisfied();
+        
assertNull(mock.getReceivedExchanges().get(0).getIn().getHeader(Exchange.REDELIVERED));
+        assertEquals(Boolean.TRUE, 
mock.getReceivedExchanges().get(1).getIn().getHeader(Exchange.REDELIVERED));
+        // confirmed after the successful redelivery
+        await().atMost(10, TimeUnit.SECONDS).until(() -> 
repository.scan(context).isEmpty());
+        assertTrue(repository.getKeys().isEmpty());
+    }
+
+    private CassandraAggregationRepository createRepository() {
+        CassandraAggregationRepository repository
+                = new CassandraAggregationRepository(inMemorySession(new 
ConcurrentHashMap<>())) {
+                    @Override
+                    public Set<String> scan(CamelContext camelContext) {
+                        try {
+                            return super.scan(camelContext);
+                        } finally {
+                            scans.incrementAndGet();
+                        }
+                    }
+                };
+        repository.setRecoveryInterval(100);
+        return repository;
+    }
+
+    private void addRoute(CassandraAggregationRepository repository, String 
name, int completionSize) throws Exception {
+        AtomicInteger failures = new AtomicInteger();
+        context.addRoutes(new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:" + name)
+                        .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(completionSize)
+                        .to("mock:" + name)
+                        .process(exchange -> {
+                            if (failures.getAndIncrement() == 0) {
+                                throw new IllegalStateException("Forced 
failure after the aggregation");
+                            }
+                        });
+            }
+        });
+    }
+
+    /**
+     * A session backed by a map of KEY to (KEY, EXCHANGE_ID, EXCHANGE), which 
runs the statements that the repository
+     * prepares for the default table (primary key KEY).
+     */
+    private CqlSession inMemorySession(Map<String, Object[]> table) {
+        CqlSession session = mock(CqlSession.class);
+        Map<BoundStatement, Object[]> boundValues = new ConcurrentHashMap<>();
+        Map<BoundStatement, String> boundQueries = new ConcurrentHashMap<>();
+        
when(session.prepare(any(SimpleStatement.class))).thenAnswer(invocation -> {
+            String query = invocation.getArgument(0, 
SimpleStatement.class).getQuery().toUpperCase();
+            PreparedStatement prepared = mock(PreparedStatement.class);
+            when(prepared.bind(any(Object[].class))).thenAnswer(bind -> {
+                BoundStatement bound = mock(BoundStatement.class);
+                boundValues.put(bound, bind.getArguments());
+                boundQueries.put(bound, query);
+                return bound;
+            });
+            return prepared;
+        });
+        when(session.execute(any(BoundStatement.class))).thenAnswer(invocation 
-> {
+            BoundStatement bound = invocation.getArgument(0);
+            return execute(table, boundQueries.remove(bound), 
boundValues.remove(bound));
+        });
+        return session;
+    }
+
+    private static ResultSet execute(Map<String, Object[]> table, String 
query, Object[] values) {
+        List<Row> rows = new ArrayList<>();
+        if (query.startsWith("INSERT")) {
+            table.put((String) values[0], values.clone());
+        } else if (query.startsWith("DELETE") && query.contains(" IF ")) {
+            table.computeIfPresent((String) values[0], (key, row) -> 
values[1].equals(row[1]) ? null : row);
+        } else if (query.startsWith("DELETE")) {
+            table.remove((String) values[0]);
+        } else if (query.startsWith("SELECT") && query.contains("WHERE")) {
+            Object[] row = table.get((String) values[0]);
+            if (row != null) {
+                rows.add(row(row));
+            }
+        } else if (query.startsWith("SELECT")) {
+            table.values().forEach(row -> rows.add(row(row)));
+        } else {
+            throw new IllegalArgumentException("Unexpected statement " + 
query);
+        }
+        ResultSet resultSet = mock(ResultSet.class);
+        when(resultSet.one()).thenReturn(rows.isEmpty() ? null : rows.get(0));
+        when(resultSet.all()).thenReturn(rows);
+        return resultSet;
+    }
+
+    private static Row row(Object[] values) {
+        Row row = mock(Row.class);
+        when(row.getString(anyString())).thenAnswer(invocation -> switch 
(invocation.getArgument(0, String.class)) {
+            case "KEY" -> values[0];
+            case "EXCHANGE_ID" -> values[1];
+            default -> throw new 
IllegalArgumentException(invocation.getArgument(0, String.class));
+        });
+        when(row.getByteBuffer("EXCHANGE")).thenAnswer(invocation -> 
((ByteBuffer) values[2]).duplicate());
+        return row;
+    }
+}
diff --git 
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
 
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
index 32273139cc7a..a6e4106974d8 100644
--- 
a/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
+++ 
b/components/camel-cassandraql/src/test/java/org/apache/camel/processor/aggregate/cassandra/NamedCassandraAggregationRepositoryIT.java
@@ -145,11 +145,16 @@ public class NamedCassandraAggregationRepositoryIT 
extends BaseCassandra {
             aggregationRepository.add(context, key, exchange);
             assertTrue(exists(key));
         }
+        Exchange completed = new DefaultExchange(context);
+        completed.setExchangeId("Exchange_2");
+        aggregationRepository.remove(context, "Confirm_2", completed);
+        assertFalse(exists("Confirm_2"));
+        assertTrue(exists("camel-recovery:Exchange_2"));
         // When
         aggregationRepository.confirm(context, "Exchange_2");
         // Then
+        assertFalse(exists("camel-recovery:Exchange_2"));
         assertTrue(exists("Confirm_1"));
-        assertFalse(exists("Confirm_2"));
         assertTrue(exists("Confirm_3"));
     }
 
@@ -172,6 +177,12 @@ public class NamedCassandraAggregationRepositoryIT extends 
BaseCassandra {
         }
     }
 
+    private void removeExchange(String key) {
+        Exchange exchange = new DefaultExchange(context);
+        exchange.setExchangeId("Exchange-" + key);
+        aggregationRepository.remove(context, key, exchange);
+    }
+
     private void addExchanges(String... keys) {
         for (String key : keys) {
             Exchange exchange = new DefaultExchange(context);
@@ -185,12 +196,16 @@ public class NamedCassandraAggregationRepositoryIT 
extends BaseCassandra {
         // Given
         String[] keys = { "Scan1", "Scan2" };
         addExchanges(keys);
+        // aggregations in progress are not to be recovered
+        
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan1"));
+        
assertFalse(aggregationRepository.scan(context).contains("Exchange-Scan2"));
         // When
+        removeExchange("Scan2");
         Set<String> exchangeIdSet = aggregationRepository.scan(context);
         // Then
-        for (String key : keys) {
-            assertTrue(exchangeIdSet.contains("Exchange-" + key));
-        }
+        assertTrue(exchangeIdSet.contains("Exchange-Scan2"));
+        assertFalse(exchangeIdSet.contains("Exchange-Scan1"));
+        
assertFalse(aggregationRepository.getKeys().contains("camel-recovery:Exchange-Scan2"));
     }
 
     @Test
@@ -198,11 +213,15 @@ public class NamedCassandraAggregationRepositoryIT 
extends BaseCassandra {
         // Given
         String[] keys = { "Recover1", "Recover2" };
         addExchanges(keys);
+        removeExchange("Recover2");
         // When
+        Exchange exchange1 = aggregationRepository.recover(context, 
"Exchange-Recover1");
         Exchange exchange2 = aggregationRepository.recover(context, 
"Exchange-Recover2");
         Exchange exchange3 = aggregationRepository.recover(context, 
"Exchange-Recover3");
         // Then
+        assertNull(exchange1);
         assertNotNull(exchange2);
+        assertEquals("Exchange-Recover2", exchange2.getExchangeId());
         assertNull(exchange3);
     }
 
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 7fbdd1071cb2..7b807df8f8a2 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
@@ -3048,6 +3048,24 @@ A poll in which the download of every file failed now 
counts as an idle poll, th
 could be acquired. This affects `backoffIdleThreshold`, 
`sendEmptyMessageWhenIdle` (an empty message is sent) and
 `greedy` (the consumer does not poll again immediately).
 
+=== camel-cassandraql - the aggregation repository keeps completed exchanges 
for recovery
+
+`CassandraAggregationRepository` (and `NamedCassandraAggregationRepository`) 
implement
+`RecoverableAggregationRepository` with recovery enabled by default, but had 
no recovery store: `remove` deleted
+the completed exchange, `confirm` deleted the aggregation in progress that had 
the confirmed exchange id, and `scan`
+returned the exchange ids of the aggregations still in progress. The recovery 
task therefore sent aggregations
+that were still in progress (marked `CamelRedelivered`) and then deleted them, 
while exchanges that failed after
+completion could never be recovered.
+
+A completed exchange is now kept in the same table, under the aggregation key 
`camel-recovery:<exchange id>`, until
+it is confirmed, which is what `scan` reports and `recover` reads. `getKeys` 
reports only the aggregations in
+progress. No change of the table is needed.
+
+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 table now also holds one row per completed and not yet 
confirmed exchange; those rows are deleted
+on confirmation. A time to live (`ttl`) applies to these rows as well.
+
 === camel-infinispan - the aggregation repository keeps completed exchanges 
for recovery
 
 `InfinispanAggregationRepository` implements 
`RecoverableAggregationRepository`, but it had no recovery

Reply via email to