This is an automated email from the ASF dual-hosted git repository.
gnodet pushed a commit to branch camel-4.22.x
in repository https://gitbox.apache.org/repos/asf/camel.git
The following commit(s) were added to refs/heads/camel-4.22.x by this push:
new 361ddea9bf48 [backport camel-4.22.x] CAMEL-25164: camel-infinispan -
recover the completed aggregation with the message that completed it (#27126)
361ddea9bf48 is described below
commit 361ddea9bf48733a529fa527665e725659caf6f3
Author: Guillaume Nodet <[email protected]>
AuthorDate: Wed Sep 30 13:36:07 2026 +0200
[backport camel-4.22.x] CAMEL-25164: camel-infinispan - recover the
completed aggregation with the message that completed it (#27126)
---
.../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 b0e0e9588f4f..b1db9fe8fec1 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());
+ }
+}