This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch camel-4.18.x
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/camel-4.18.x by this push:
     new c385ca1c79e2 [backport camel-4.18.x] CAMEL-25164: camel-infinispan - 
recover the completed aggregation with the message that completed it (#27127)
c385ca1c79e2 is described below

commit c385ca1c79e2ffd38beea047fd1add6595b2614b
Author: Guillaume Nodet <[email protected]>
AuthorDate: Wed Sep 30 14:05:30 2026 +0200

    [backport camel-4.18.x] CAMEL-25164: camel-infinispan - recover the 
completed aggregation with the message that completed it (#27127)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../InfinispanAggregationRepository.java           | 13 ++--
 ...mbeddedAggregationRepositoryOperationsTest.java | 20 ++++++
 ...anEmbeddedAggregationRepositoryRecoverTest.java | 84 ++++++++++++++++++++++
 3 files changed, 110 insertions(+), 7 deletions(-)

diff --git 
a/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/InfinispanAggregationRepository.java
 
b/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/InfinispanAggregationRepository.java
index a31a3951d277..367858f2bee4 100644
--- 
a/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/InfinispanAggregationRepository.java
+++ 
b/components/camel-infinispan/camel-infinispan-common/src/main/java/org/apache/camel/component/infinispan/InfinispanAggregationRepository.java
@@ -103,16 +103,15 @@ public abstract class InfinispanAggregationRepository
     @Override
     public void remove(CamelContext camelContext, String key, Exchange 
exchange) {
         LOG.trace("Removing an exchange with ID {} for key {}", 
exchange.getExchangeId(), key);
-        DefaultExchangeHolder holder = getCache().remove(key);
+        getCache().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
-            if (holder == null) {
-                holder = DefaultExchangeHolder.marshal(exchange, true, 
allowSerializedHeaders);
-            }
+            // 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 does not
+            // contain the exchange that completed the aggregation)
             LOG.trace("Putting an exchange with ID {} into the recovery 
store", exchange.getExchangeId());
-            getCache().put(recoveryKey(exchange.getExchangeId()), holder);
+            getCache().put(recoveryKey(exchange.getExchangeId()),
+                    DefaultExchangeHolder.marshal(exchange, true, 
allowSerializedHeaders));
         }
     }
 
diff --git 
a/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/InfinispanEmbeddedAggregationRepositoryOperationsTest.java
 
b/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/InfinispanEmbeddedAggregationRepositoryOperationsTest.java
index 82339c3909e7..a4209acf9624 100644
--- 
a/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/InfinispanEmbeddedAggregationRepositoryOperationsTest.java
+++ 
b/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/InfinispanEmbeddedAggregationRepositoryOperationsTest.java
@@ -236,6 +236,26 @@ public class 
InfinispanEmbeddedAggregationRepositoryOperationsTest extends Infin
         assertNull(unknown);
     }
 
+    @Test
+    public void testRecoverTheExchangeGivenToRemove() {
+        // cleanup
+        aggregationRepository.getCache().clear();
+        // Given - the group in the cache does not contain the message that 
completed it
+        Exchange exchange = new DefaultExchange(context());
+        exchange.setExchangeId("Exchange-RecoverLast");
+        exchange.getIn().setBody("a+b");
+        aggregationRepository.add(context(), "RecoverLast", exchange);
+        exchange.getIn().setBody("a+b+c");
+
+        // When
+        aggregationRepository.remove(context(), "RecoverLast", exchange);
+
+        // Then
+        Exchange recovered = aggregationRepository.recover(context(), 
"Exchange-RecoverLast");
+        assertNotNull(recovered);
+        assertEquals("a+b+c", recovered.getIn().getBody(String.class));
+    }
+
     @Test
     public void testGetKeysIgnoresExchangesToRecover() {
         // cleanup
diff --git 
a/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/InfinispanEmbeddedAggregationRepositoryRecoverTest.java
 
b/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/InfinispanEmbeddedAggregationRepositoryRecoverTest.java
new file mode 100644
index 000000000000..0b57f66ae63e
--- /dev/null
+++ 
b/components/camel-infinispan/camel-infinispan-embedded/src/test/java/org/apache/camel/component/infinispan/embedded/InfinispanEmbeddedAggregationRepositoryRecoverTest.java
@@ -0,0 +1,84 @@
+/*
+ * 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.infinispan.embedded;
+
+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.junit.jupiter.api.Test;
+
+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 InfinispanEmbeddedAggregationRepositoryRecoverTest extends 
InfinispanEmbeddedTestSupport {
+
+    private final AtomicInteger failures = new AtomicInteger();
+
+    @Test
+    public void testRecoveredExchangeContainsTheLastMessage() throws Exception 
{
+        InfinispanEmbeddedConfiguration configuration = new 
InfinispanEmbeddedConfiguration();
+        configuration.setCacheContainer(cacheContainer);
+
+        InfinispanEmbeddedAggregationRepository repository = new 
InfinispanEmbeddedAggregationRepository(getCacheName());
+        repository.setConfiguration(configuration);
+        repository.setRecoveryInterval(100);
+
+        context.addRoutes(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");
+                            }
+                        });
+            }
+        });
+
+        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());
+    }
+}

Reply via email to