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 6836a5367801 CAMEL-25216: camel-redis - recover the completed 
aggregation with the message that completed it (#27165)
6836a5367801 is described below

commit 6836a5367801a761e09057a6a98d75dc31e60d4d
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 12:24:43 2026 +0530

    CAMEL-25216: camel-redis - recover the completed aggregation with the 
message that completed it (#27165)
    
    With recovery enabled (the default) and without optimistic locking,
    RedisAggregationRepository.remove() stored the entry it removed from the
    map as the exchange to recover. When a group is completed by an incoming
    exchange, the Aggregate EIP does not add the aggregated exchange to the
    repository before it removes the group, so that entry is the group before
    the last message. A completed exchange whose processing failed was
    recovered without the message that completed the group.
    
    Store the marshalled exchange passed to remove(), as the optimistic branch
    already does. This is the follow-up named in CAMEL-24946, which made the
    same change in KeyValueAggregationRepository.
    
    AggregateRedisRecoverIT sends a group of three messages whose processing
    fails once: the recovered exchange must hold all three messages. An
    operations test checks that the exchange given to remove() is recovered.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../aggregate/RedisAggregationRepository.java      |  6 +-
 .../integration/AggregateRedisRecoverIT.java       | 93 ++++++++++++++++++++++
 .../RedisAggregationRepositoryOperationsIT.java    | 20 +++++
 3 files changed, 117 insertions(+), 2 deletions(-)

diff --git 
a/components/camel-redis/src/main/java/org/apache/camel/component/redis/processor/aggregate/RedisAggregationRepository.java
 
b/components/camel-redis/src/main/java/org/apache/camel/component/redis/processor/aggregate/RedisAggregationRepository.java
index 214fbf4da58f..91d406fe3af3 100644
--- 
a/components/camel-redis/src/main/java/org/apache/camel/component/redis/processor/aggregate/RedisAggregationRepository.java
+++ 
b/components/camel-redis/src/main/java/org/apache/camel/component/redis/processor/aggregate/RedisAggregationRepository.java
@@ -321,10 +321,12 @@ public class RedisAggregationRepository extends 
ServiceSupport
                     RMap<String, DefaultExchangeHolder> tCache = 
transaction.getMap(mapName);
                     RMap<String, DefaultExchangeHolder> tPersistentCache = 
transaction.getMap(persistenceMapName);
 
-                    DefaultExchangeHolder removedHolder = tCache.remove(key);
+                    tCache.remove(key);
                     LOG.trace("Putting an exchange with ID {} for key {} into 
a recoverable storage in a thread-safe manner.",
                             exchange.getExchangeId(), key);
-                    tPersistentCache.put(exchange.getExchangeId(), 
removedHolder);
+                    // store the given exchange and not the removed entry: 
when a group is completed by an incoming
+                    // exchange, the last aggregated exchange is not added to 
the repository before it is removed
+                    tPersistentCache.put(exchange.getExchangeId(), holder);
 
                     transaction.commit();
                     LOG.trace("Removed an exchange with ID {} for key {} in a 
thread-safe manner.", exchange.getExchangeId(),
diff --git 
a/components/camel-redis/src/test/java/org/apache/camel/component/redis/processor/aggregate/integration/AggregateRedisRecoverIT.java
 
b/components/camel-redis/src/test/java/org/apache/camel/component/redis/processor/aggregate/integration/AggregateRedisRecoverIT.java
new file mode 100644
index 000000000000..f98bf49c363b
--- /dev/null
+++ 
b/components/camel-redis/src/test/java/org/apache/camel/component/redis/processor/aggregate/integration/AggregateRedisRecoverIT.java
@@ -0,0 +1,93 @@
+/*
+ * 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.redis.processor.aggregate.integration;
+
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import 
org.apache.camel.component.redis.processor.aggregate.RedisAggregationRepository;
+import org.apache.camel.test.infra.redis.services.RedisService;
+import org.apache.camel.test.infra.redis.services.RedisServiceFactory;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.extension.RegisterExtension;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * A completed aggregation whose processing fails is recovered with all its 
messages, including the one that completed
+ * the group (which is never added to the repository).
+ */
+public class AggregateRedisRecoverIT extends CamelTestSupport {
+
+    @RegisterExtension
+    static RedisService service = RedisServiceFactory.createSingletonService();
+
+    private final AtomicInteger failures = new AtomicInteger();
+
+    private RedisAggregationRepository repository;
+
+    @Test
+    public void testRecoveredExchangeContainsTheLastMessage() throws Exception 
{
+        MockEndpoint mock = getMockEndpoint("mock:aggregated");
+        // the completed aggregation fails once and is then recovered with all 
three messages
+        mock.expectedBodiesReceived("a+b+c", "a+b+c");
+
+        template.sendBodyAndHeader("direct:start", "a", "id", "group");
+        template.sendBodyAndHeader("direct:start", "b", "id", "group");
+        template.sendBodyAndHeader("direct:start", "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));
+        assertTrue(repository.getKeys().isEmpty());
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        repository = new RedisAggregationRepository("aggregationRecover", 
service.getServiceAddress());
+        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(3)
+                        .to("mock:aggregated")
+                        .process(exchange -> {
+                            if (failures.getAndIncrement() == 0) {
+                                throw new IllegalStateException("Forced 
failure after the aggregation");
+                            }
+                        });
+            }
+        };
+    }
+}
diff --git 
a/components/camel-redis/src/test/java/org/apache/camel/component/redis/processor/aggregate/integration/RedisAggregationRepositoryOperationsIT.java
 
b/components/camel-redis/src/test/java/org/apache/camel/component/redis/processor/aggregate/integration/RedisAggregationRepositoryOperationsIT.java
index 736bf98abe22..2ed20bc09065 100644
--- 
a/components/camel-redis/src/test/java/org/apache/camel/component/redis/processor/aggregate/integration/RedisAggregationRepositoryOperationsIT.java
+++ 
b/components/camel-redis/src/test/java/org/apache/camel/component/redis/processor/aggregate/integration/RedisAggregationRepositoryOperationsIT.java
@@ -177,6 +177,26 @@ public class RedisAggregationRepositoryOperationsIT 
extends CamelTestSupport {
         }
     }
 
+    @Test
+    public void testPessimisticRemoveRecoversTheGivenExchange() {
+        RedisAggregationRepository repo = createRepo("pessimisticRemoveGiven", 
false);
+        try {
+            repo.add(context, "key1", createExchange("a+b"));
+
+            // the Aggregate EIP aggregates the message that completes the 
group into the exchange it read from the
+            // repository, and removes the group without adding that exchange 
first
+            Exchange completed = repo.get(context, "key1");
+            completed.getIn().setBody("a+b+c");
+            repo.remove(context, "key1", completed);
+
+            Exchange recovered = repo.recover(context, 
completed.getExchangeId());
+            assertNotNull(recovered);
+            assertEquals("a+b+c", recovered.getIn().getBody(String.class));
+        } finally {
+            repo.stop();
+        }
+    }
+
     @Test
     public void testRemoveWithoutRecovery() {
         RedisAggregationRepository repo

Reply via email to