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 c926fac9e24a CAMEL-24994: camel-core - Enrich and Poll Enrich EIPs: 
the resource on completions are lost when the enrich fails (#26850)
c926fac9e24a is described below

commit c926fac9e24a829bec8d82d93131756a8d9cce7f
Author: Claus Ibsen <[email protected]>
AuthorDate: Sat Sep 26 21:25:11 2026 +0200

    CAMEL-24994: camel-core - Enrich and Poll Enrich EIPs: the resource on 
completions are lost when the enrich fails (#26850)
    
    The on completions registered on the resource exchange were only handed
    over to the original exchange when the enrichment succeeded. When the
    resource failed, or the aggregation strategy threw or returned null,
    they never ran: a file polled by pollEnrich stayed in progress and was
    not polled again, and a producer could not release its resources (such
    as closing a response stream).
    
    Handing the on completions of a failed or not aggregated resource over
    to the exchange tied the release of the resource to how the exchange
    ends. When the exchange then recovered (for example by redelivery), a
    polled file whose content never reached any exchange was committed.
    Now the on completions are only handed over when the resource was used
    (aggregated, or discarded by the aggregation strategy). When the
    resource or the aggregation failed, they run right away as failed, so a
    polled file is rolled back and can be polled again.
    
    Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
    Signed-off-by: Claus Ibsen <[email protected]>
---
 .../java/org/apache/camel/processor/Enricher.java  |  23 +++-
 .../org/apache/camel/processor/PollEnricher.java   |  30 ++++-
 .../enricher/EnricherResourceCompletionTest.java   | 137 +++++++++++++++++++++
 .../PollEnrichFileAggregationFailureTest.java      | 118 ++++++++++++++++++
 4 files changed, 300 insertions(+), 8 deletions(-)

diff --git 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/Enricher.java
 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/Enricher.java
index 5ab1fa795b6d..dbfadfde4fd7 100644
--- 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/Enricher.java
+++ 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/Enricher.java
@@ -37,6 +37,7 @@ import org.apache.camel.spi.StepIdAware;
 import org.apache.camel.support.DefaultExchange;
 import org.apache.camel.support.ExchangeHelper;
 import org.apache.camel.support.MessageHelper;
+import org.apache.camel.support.UnitOfWorkHelper;
 import org.apache.camel.support.service.ServiceHelper;
 
 import static 
org.apache.camel.support.ExchangeHelper.copyResultsPreservePattern;
@@ -228,6 +229,8 @@ public class Enricher extends BaseProcessorSupport
         return sendDynamicProcessor.process(resourceExchange, new 
AsyncCallback() {
             @Override
             public void done(boolean doneSync) {
+                // whether the resource was used (aggregated, or discarded by 
the aggregation strategy)
+                boolean used = false;
                 if (!isAggregateOnException() && resourceExchange.isFailed()) {
                     // copy resource exchange onto original exchange 
(preserving pattern)
                     copyResultsWithoutCorrelationId(exchange, 
resourceExchange);
@@ -249,14 +252,26 @@ public class Enricher extends BaseProcessorSupport
                             }
                             // copy aggregation result onto original exchange 
(preserving pattern)
                             copyResultsWithoutCorrelationId(exchange, 
aggregatedExchange);
-                            // handover any synchronization (if unit of work 
is not shared)
-                            if (resourceExchange != null && 
!isShareUnitOfWork()) {
-                                
resourceExchange.getExchangeExtension().handoverCompletions(exchange);
-                            }
                         }
+                        used = true;
                     } catch (Exception e) {
                         // if the aggregationStrategy threw an exception, set 
it on the original exchange
                         exchange.setException(new 
CamelExchangeException("Error occurred during aggregation", exchange, e));
+                        if (resourceExchange.getException() == null) {
+                            resourceExchange.setException(e);
+                        }
+                    }
+                }
+                if (!isShareUnitOfWork()) {
+                    if (used) {
+                        // handover any synchronization to complete together 
with the exchange
+                        
resourceExchange.getExchangeExtension().handoverCompletions(exchange);
+                    } else {
+                        // the resource or the aggregation failed, so release 
the resource as failed now (such as
+                        // closing a response stream), and not together with 
the exchange, which may still complete
+                        // successfully (such as by redelivery)
+                        UnitOfWorkHelper.doneSynchronizations(resourceExchange,
+                                
resourceExchange.getExchangeExtension().handoverCompletions());
                     }
                 }
 
diff --git 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/PollEnricher.java
 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/PollEnricher.java
index bf75e9a7451f..b1839986616d 100644
--- 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/PollEnricher.java
+++ 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/PollEnricher.java
@@ -51,6 +51,7 @@ import org.apache.camel.support.EndpointHelper;
 import org.apache.camel.support.EventDrivenPollingConsumer;
 import org.apache.camel.support.ExchangeHelper;
 import org.apache.camel.support.NormalizedUri;
+import org.apache.camel.support.UnitOfWorkHelper;
 import org.apache.camel.support.cache.DefaultConsumerCache;
 import org.apache.camel.support.cache.EmptyConsumerCache;
 import org.apache.camel.support.service.ServiceHelper;
@@ -409,6 +410,7 @@ public class PollEnricher extends BaseProcessorSupport 
implements IdAware, Route
                 originalHeaders = 
headersMapFactory.newMap(exchange.getMessage().getHeaders());
             } catch (Exception throwable) {
                 exchange.setException(throwable);
+                doneResourceAsFailed(resourceExchange, throwable);
                 callback.done(true);
                 return true;
             }
@@ -419,6 +421,8 @@ public class PollEnricher extends BaseProcessorSupport 
implements IdAware, Route
                 // copy resource exchange onto original exchange (preserving 
pattern)
                 // and preserve redelivery headers
                 copyResultsPreservePattern(exchange, resourceExchange);
+                // the resource failed, so release it as failed now (such as 
rolling back a polled file)
+                doneResourceAsFailed(resourceExchange, null);
             } else {
                 prepareResult(exchange);
 
@@ -436,11 +440,10 @@ public class PollEnricher extends BaseProcessorSupport 
implements IdAware, Route
                     }
                     // copy aggregation result onto original exchange 
(preserving pattern)
                     copyResultsPreservePattern(exchange, aggregatedExchange);
-                    // handover any synchronization
-                    if (resourceExchange != null) {
-                        
resourceExchange.getExchangeExtension().handoverCompletions(exchange);
-                    }
                 }
+                // the resource has been aggregated (or discarded by the 
aggregation strategy), so handover any
+                // synchronization to complete together with the exchange 
(such as committing a polled file)
+                handoverCompletions(resourceExchange, exchange);
             }
 
             // if we failed then restore caused exception
@@ -467,6 +470,9 @@ public class PollEnricher extends BaseProcessorSupport 
implements IdAware, Route
 
         } catch (Exception e) {
             exchange.setException(new CamelExchangeException("Error occurred 
during aggregation", exchange, e));
+            // the resource was not aggregated, so release it as failed now 
(such as rolling back a polled file),
+            // and not together with the exchange, which may still complete 
successfully (such as by redelivery)
+            doneResourceAsFailed(resourceExchange, e);
             callback.done(true);
             return true;
         }
@@ -475,6 +481,22 @@ public class PollEnricher extends BaseProcessorSupport 
implements IdAware, Route
         return true;
     }
 
+    private static void handoverCompletions(Exchange resourceExchange, 
Exchange exchange) {
+        if (resourceExchange != null) {
+            
resourceExchange.getExchangeExtension().handoverCompletions(exchange);
+        }
+    }
+
+    private static void doneResourceAsFailed(Exchange resourceExchange, 
Exception cause) {
+        if (resourceExchange != null) {
+            if (cause != null && resourceExchange.getException() == null) {
+                resourceExchange.setException(cause);
+            }
+            UnitOfWorkHelper.doneSynchronizations(resourceExchange,
+                    
resourceExchange.getExchangeExtension().handoverCompletions());
+        }
+    }
+
     private static boolean isBridgeErrorHandler(PollingConsumer consumer) {
         Consumer delegate = consumer;
         if (consumer instanceof EventDrivenPollingConsumer 
eventDrivenPollingConsumer) {
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/enricher/EnricherResourceCompletionTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/enricher/EnricherResourceCompletionTest.java
new file mode 100644
index 000000000000..f8452b416b12
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/enricher/EnricherResourceCompletionTest.java
@@ -0,0 +1,137 @@
+/*
+ * 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.enricher;
+
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.AggregationStrategy;
+import org.apache.camel.Consumer;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Endpoint;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.support.DefaultComponent;
+import org.apache.camel.support.DefaultEndpoint;
+import org.apache.camel.support.DefaultProducer;
+import org.apache.camel.support.SynchronizationAdapter;
+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.assertThrows;
+
+/**
+ * The on completions registered on the resource exchange (for example to 
close a response stream) run when the enrich
+ * fails (the resource fails, or the aggregation fails), when the aggregation 
strategy returns null, and on success.
+ */
+public class EnricherResourceCompletionTest extends ContextTestSupport {
+
+    private final AtomicInteger completions = new AtomicInteger();
+    private final AtomicInteger failures = new AtomicInteger();
+
+    @Test
+    public void testResourceFailed() {
+        assertThrows(Exception.class, () -> 
template.requestBody("direct:resourceFails", "Hello"));
+        await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, completions.get()));
+        assertEquals(1, failures.get(), "the resource should be completed as 
failed");
+    }
+
+    @Test
+    public void testAggregationFailed() {
+        assertThrows(Exception.class, () -> 
template.requestBody("direct:aggregationFails", "Hello"));
+        await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, completions.get()));
+        assertEquals(1, failures.get(), "the resource should be completed as 
failed");
+    }
+
+    @Test
+    public void testAggregationReturnsNull() {
+        // the strategy discards the resource, so the original message 
continues unchanged
+        String out = template.requestBody("direct:aggregationReturnsNull", 
"Hello", String.class);
+        assertEquals("Hello", out);
+        await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, completions.get()));
+    }
+
+    @Test
+    public void testSuccess() {
+        template.requestBody("direct:ok", "Hello");
+        await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> 
assertEquals(1, completions.get()));
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        // a producer that registers an on completion on the exchange it is 
given (such as a producer that must close
+        // a response stream when the exchange is done), and fails when its 
path is "failing"
+        context.addComponent("resource", new DefaultComponent() {
+            @Override
+            protected Endpoint createEndpoint(String uri, String remaining, 
Map<String, Object> parameters) {
+                return new DefaultEndpoint(uri, this) {
+                    @Override
+                    public Producer createProducer() {
+                        return new DefaultProducer(this) {
+                            @Override
+                            public void process(Exchange exchange) {
+                                addCompletion(exchange);
+                                if ("failing".equals(remaining)) {
+                                    throw new 
IllegalArgumentException("Forced");
+                                }
+                                exchange.getMessage().setBody("Resource");
+                            }
+                        };
+                    }
+
+                    @Override
+                    public Consumer createConsumer(Processor processor) {
+                        throw new UnsupportedOperationException();
+                    }
+                };
+            }
+        });
+
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:resourceFails").enrich("resource:failing");
+                from("direct:aggregationFails").enrich("resource:ok", new 
FailingStrategy());
+                from("direct:aggregationReturnsNull").enrich("resource:ok", 
(oldExchange, newExchange) -> null);
+                from("direct:ok").enrich("resource:ok");
+            }
+        };
+    }
+
+    private void addCompletion(Exchange exchange) {
+        exchange.getExchangeExtension().addOnCompletion(new 
SynchronizationAdapter() {
+            @Override
+            public void onDone(Exchange exchange) {
+                completions.incrementAndGet();
+                if (exchange.isFailed()) {
+                    failures.incrementAndGet();
+                }
+            }
+        });
+    }
+
+    private static class FailingStrategy implements AggregationStrategy {
+        @Override
+        public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
+            throw new IllegalArgumentException("Forced");
+        }
+    }
+}
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/enricher/PollEnrichFileAggregationFailureTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/enricher/PollEnrichFileAggregationFailureTest.java
new file mode 100644
index 000000000000..a5daf681d175
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/enricher/PollEnrichFileAggregationFailureTest.java
@@ -0,0 +1,118 @@
+/*
+ * 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.enricher;
+
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.AggregationStrategy;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+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.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * When the aggregation of a polled file fails, the file is released (its on 
completion runs), so it can be polled
+ * again.
+ */
+public class PollEnrichFileAggregationFailureTest extends ContextTestSupport {
+
+    private final AtomicInteger attempts = new AtomicInteger();
+    private final AtomicInteger redeliveryAttempts = new AtomicInteger();
+
+    @Test
+    public void testFileCanBePolledAgainAfterAggregationFailure() {
+        template.sendBodyAndHeader(fileUri("data"), "Big file", 
Exchange.FILE_NAME, "AAA.dat");
+
+        // the first aggregation fails
+        assertThrows(Exception.class, () -> 
template.requestBody("direct:start", "Start"));
+
+        // the file must be released so the second poll gets it again
+        String out = template.requestBody("direct:start", "Start", 
String.class);
+        assertEquals("Big file", out);
+    }
+
+    @Test
+    public void testFileNotAggregatedIsNotCommittedWhenRedeliverySucceeds() 
throws Exception {
+        template.sendBodyAndHeader(fileUri("inbox"), "A", Exchange.FILE_NAME, 
"a.txt");
+        template.sendBodyAndHeader(fileUri("inbox"), "B", Exchange.FILE_NAME, 
"b.txt");
+
+        // the first aggregation fails, and the redelivery polls and 
aggregates the other file
+        String out = template.requestBody("direct:redelivery", "Start", 
String.class);
+        assertTrue("A".equals(out) || "B".equals(out), "one of the files 
should be aggregated, was: " + out);
+
+        // the file that was not aggregated must not be committed (moved to 
.camel), so it can be polled again
+        String notAggregated = "A".equals(out) ? "b.txt" : "a.txt";
+        assertFileExists(testFile("inbox/" + notAggregated));
+        assertFileNotExists(testFile("inbox/" + out.toLowerCase() + ".txt"));
+    }
+
+    @Test
+    public void testFileIsCommittedWhenAggregationReturnsNull() {
+        template.sendBodyAndHeader(fileUri("discard"), "Content", 
Exchange.FILE_NAME, "x.dat");
+
+        // the strategy discards the polled file (returns null), so the 
original message continues unchanged
+        String out = template.requestBody("direct:discard", "Start", 
String.class);
+        assertEquals("Start", out);
+
+        // the polled file was used (discarded by the strategy), so it is 
committed (moved to .camel)
+        await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> 
assertFileNotExists(testFile("discard/x.dat")));
+        assertFileExists(testFile("discard/.camel/x.dat"));
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:start")
+                        .pollEnrich(fileUri("data?initialDelay=0&delay=10"), 
2000, new FailOnceStrategy());
+
+                from("direct:discard")
+                        
.pollEnrich(fileUri("discard?initialDelay=0&delay=10"), 2000, (original, 
resource) -> null);
+
+                from("direct:redelivery")
+                        
.errorHandler(defaultErrorHandler().maximumRedeliveries(1).redeliveryDelay(0))
+                        .pollEnrich(fileUri("inbox?initialDelay=0&delay=10"), 
2000, (original, resource) -> {
+                            if (redeliveryAttempts.incrementAndGet() == 1) {
+                                throw new IllegalStateException("Transient 
failure");
+                            }
+                            
original.getMessage().setBody(resource.getMessage().getBody(String.class));
+                            return original;
+                        });
+            }
+        };
+    }
+
+    private class FailOnceStrategy implements AggregationStrategy {
+        @Override
+        public Exchange aggregate(Exchange original, Exchange resource) {
+            if (attempts.incrementAndGet() == 1) {
+                throw new IllegalArgumentException("Forced");
+            }
+            if (resource != null) {
+                
original.getMessage().setBody(resource.getMessage().getBody(String.class));
+            }
+            return original;
+        }
+    }
+}

Reply via email to