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 f78a93999a24 CAMEL-25124: camel-core - Aggregate EIP: keep spooled 
stream cache bodies until the aggregated exchange is done (#27033)
f78a93999a24 is described below

commit f78a93999a24d1fb0f1e85474b2fdc8aa9b71a9d
Author: allthingssecurity <[email protected]>
AuthorDate: Tue Sep 29 17:25:35 2026 +0530

    CAMEL-25124: camel-core - Aggregate EIP: keep spooled stream cache bodies 
until the aggregated exchange is done (#27033)
    
    * CAMEL-25124: camel-core - Aggregate EIP: keep spooled stream cache bodies 
until the aggregated exchange is done
    * CAMEL-25124: camel-core - Aggregate EIP: release only the spool 
references the aggregator took
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../TarAggregationStrategyRepositoryTest.java      | 147 +++++++
 .../ZipAggregationStrategyRepositoryTest.java      | 147 +++++++
 .../processor/aggregate/AggregateProcessor.java    | 159 ++++++-
 .../AggregateStreamCachingSpoolTest.java           | 458 +++++++++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  16 +
 5 files changed, 925 insertions(+), 2 deletions(-)

diff --git 
a/components/camel-tarfile/src/test/java/org/apache/camel/processor/aggregate/tarfile/TarAggregationStrategyRepositoryTest.java
 
b/components/camel-tarfile/src/test/java/org/apache/camel/processor/aggregate/tarfile/TarAggregationStrategyRepositoryTest.java
new file mode 100644
index 000000000000..f8b8c85d8954
--- /dev/null
+++ 
b/components/camel-tarfile/src/test/java/org/apache/camel/processor/aggregate/tarfile/TarAggregationStrategyRepositoryTest.java
@@ -0,0 +1,147 @@
+/*
+ * 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.tarfile;
+
+import java.io.File;
+import java.io.FileInputStream;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+
+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.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.spi.AggregationRepository;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.apache.commons.compress.archivers.tar.TarArchiveEntry;
+import org.apache.commons.compress.archivers.tar.TarArchiveInputStream;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.apache.camel.test.junit6.TestSupport.deleteDirectory;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+/**
+ * The tar file that the strategy appends to must not be deleted before the 
group completes, when the aggregation
+ * repository does not keep the exchange instances (optimistic locking, or a 
persistent repository).
+ */
+public class TarAggregationStrategyRepositoryTest extends CamelTestSupport {
+
+    private static final String TEST_DIR = 
"target/out_TarAggregationStrategyRepositoryTest";
+
+    @BeforeEach
+    public void deleteTestDirs() {
+        deleteDirectory(TEST_DIR);
+    }
+
+    @Test
+    public void testOptimisticLocking() throws Exception {
+        sendAndAssertZip("optimistic");
+    }
+
+    @Test
+    public void testRepositoryStoringCopies() throws Exception {
+        sendAndAssertZip("copies");
+    }
+
+    private void sendAndAssertZip(String name) throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:" + name);
+        mock.expectedMessageCount(1);
+
+        template.sendBody("direct:" + name, "Hello");
+        template.sendBody("direct:" + name, "Hello again");
+        template.sendBody("direct:" + name, "Bye");
+
+        MockEndpoint.assertIsSatisfied(context);
+
+        File[] files = new File(TEST_DIR, name).listFiles();
+        assertNotNull(files);
+        assertEquals(1, files.length);
+        int count = 0;
+        try (TarArchiveInputStream tin = new TarArchiveInputStream(new 
FileInputStream(files[0]))) {
+            for (TarArchiveEntry te = tin.getNextEntry(); te != null; te = 
tin.getNextEntry()) {
+                count++;
+            }
+        }
+        assertEquals(3, count, "Tar file should contain 3 files");
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:optimistic")
+                        .aggregate(constant(true), strategy())
+                        .aggregationRepository(new 
MemoryAggregationRepository(true)).optimisticLocking()
+                        .completionSize(3)
+                        .to("file:" + TEST_DIR + "/optimistic")
+                        .to("mock:optimistic");
+
+                from("direct:copies")
+                        .aggregate(constant(true), strategy())
+                        .aggregationRepository(new CopyingRepository())
+                        .completionSize(3)
+                        .to("file:" + TEST_DIR + "/copies")
+                        .to("mock:copies");
+            }
+        };
+    }
+
+    private static TarAggregationStrategy strategy() {
+        TarAggregationStrategy strategy = new TarAggregationStrategy();
+        strategy.setParentDir(TEST_DIR + "/temp");
+        return strategy;
+    }
+
+    /**
+     * A repository that stores a copy of the exchange, as the persistent 
repositories do.
+     */
+    private static final class CopyingRepository extends ServiceSupport 
implements AggregationRepository {
+        private final Map<String, Exchange> exchanges = new 
ConcurrentHashMap<>();
+
+        @Override
+        public Exchange add(CamelContext camelContext, String key, Exchange 
exchange) {
+            return exchanges.put(key, exchange.copy());
+        }
+
+        @Override
+        public Exchange get(CamelContext camelContext, String key) {
+            Exchange exchange = exchanges.get(key);
+            return exchange != null ? exchange.copy() : null;
+        }
+
+        @Override
+        public void remove(CamelContext camelContext, String key, Exchange 
exchange) {
+            exchanges.remove(key);
+        }
+
+        @Override
+        public void confirm(CamelContext camelContext, String exchangeId) {
+            // noop
+        }
+
+        @Override
+        public Set<String> getKeys() {
+            return exchanges.keySet();
+        }
+    }
+}
diff --git 
a/components/camel-zipfile/src/test/java/org/apache/camel/processor/aggregate/zipfile/ZipAggregationStrategyRepositoryTest.java
 
b/components/camel-zipfile/src/test/java/org/apache/camel/processor/aggregate/zipfile/ZipAggregationStrategyRepositoryTest.java
new file mode 100644
index 000000000000..bf28b2944cab
--- /dev/null
+++ 
b/components/camel-zipfile/src/test/java/org/apache/camel/processor/aggregate/zipfile/ZipAggregationStrategyRepositoryTest.java
@@ -0,0 +1,147 @@
+/*
+ * 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.zipfile;
+
+import java.io.File;
+import java.io.FileInputStream;
+import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
+import java.util.zip.ZipEntry;
+import java.util.zip.ZipInputStream;
+
+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.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.spi.AggregationRepository;
+import org.apache.camel.support.service.ServiceSupport;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.apache.camel.test.junit6.TestSupport.deleteDirectory;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+/**
+ * The zip file that the strategy appends to must not be deleted before the 
group completes, when the aggregation
+ * repository does not keep the exchange instances (optimistic locking, or a 
persistent repository).
+ */
+public class ZipAggregationStrategyRepositoryTest extends CamelTestSupport {
+
+    private static final String TEST_DIR = 
"target/out_ZipAggregationStrategyRepositoryTest";
+
+    @BeforeEach
+    public void deleteTestDirs() {
+        deleteDirectory(TEST_DIR);
+    }
+
+    @Test
+    public void testOptimisticLocking() throws Exception {
+        sendAndAssertZip("optimistic");
+    }
+
+    @Test
+    public void testRepositoryStoringCopies() throws Exception {
+        sendAndAssertZip("copies");
+    }
+
+    private void sendAndAssertZip(String name) throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:" + name);
+        mock.expectedMessageCount(1);
+
+        template.sendBody("direct:" + name, "Hello");
+        template.sendBody("direct:" + name, "Hello again");
+        template.sendBody("direct:" + name, "Bye");
+
+        MockEndpoint.assertIsSatisfied(context);
+
+        File[] files = new File(TEST_DIR, name).listFiles();
+        assertNotNull(files);
+        assertEquals(1, files.length);
+        int count = 0;
+        try (ZipInputStream zin = new ZipInputStream(new 
FileInputStream(files[0]))) {
+            for (ZipEntry ze = zin.getNextEntry(); ze != null; ze = 
zin.getNextEntry()) {
+                count++;
+            }
+        }
+        assertEquals(3, count, "Zip file should contain 3 files");
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:optimistic")
+                        .aggregate(constant(true), strategy())
+                        .aggregationRepository(new 
MemoryAggregationRepository(true)).optimisticLocking()
+                        .completionSize(3)
+                        .to("file:" + TEST_DIR + "/optimistic")
+                        .to("mock:optimistic");
+
+                from("direct:copies")
+                        .aggregate(constant(true), strategy())
+                        .aggregationRepository(new CopyingRepository())
+                        .completionSize(3)
+                        .to("file:" + TEST_DIR + "/copies")
+                        .to("mock:copies");
+            }
+        };
+    }
+
+    private static ZipAggregationStrategy strategy() {
+        ZipAggregationStrategy strategy = new ZipAggregationStrategy();
+        strategy.setParentDir(TEST_DIR + "/temp");
+        return strategy;
+    }
+
+    /**
+     * A repository that stores a copy of the exchange, as the persistent 
repositories do.
+     */
+    private static final class CopyingRepository extends ServiceSupport 
implements AggregationRepository {
+        private final Map<String, Exchange> exchanges = new 
ConcurrentHashMap<>();
+
+        @Override
+        public Exchange add(CamelContext camelContext, String key, Exchange 
exchange) {
+            return exchanges.put(key, exchange.copy());
+        }
+
+        @Override
+        public Exchange get(CamelContext camelContext, String key) {
+            Exchange exchange = exchanges.get(key);
+            return exchange != null ? exchange.copy() : null;
+        }
+
+        @Override
+        public void remove(CamelContext camelContext, String key, Exchange 
exchange) {
+            exchanges.remove(key);
+        }
+
+        @Override
+        public void confirm(CamelContext camelContext, String exchangeId) {
+            // noop
+        }
+
+        @Override
+        public Set<String> getKeys() {
+            return exchanges.keySet();
+        }
+    }
+}
diff --git 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
index 574c03376e9b..93392d6c66a3 100644
--- 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
+++ 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java
@@ -16,7 +16,9 @@
  */
 package org.apache.camel.processor.aggregate;
 
+import java.io.IOException;
 import java.util.ArrayList;
+import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
@@ -48,6 +50,7 @@ import org.apache.camel.Predicate;
 import org.apache.camel.Processor;
 import org.apache.camel.ProducerTemplate;
 import org.apache.camel.ShutdownRunningTask;
+import org.apache.camel.StreamCache;
 import org.apache.camel.TimeoutMap;
 import org.apache.camel.Traceable;
 import org.apache.camel.processor.BaseProcessorSupport;
@@ -62,12 +65,15 @@ import org.apache.camel.spi.RouteIdAware;
 import org.apache.camel.spi.ShutdownAware;
 import org.apache.camel.spi.StepIdAware;
 import org.apache.camel.spi.Synchronization;
+import org.apache.camel.support.DefaultExchange;
 import org.apache.camel.support.DefaultTimeoutMap;
 import org.apache.camel.support.ExchangeHelper;
 import org.apache.camel.support.KeyValueAggregationRepository;
 import org.apache.camel.support.LRUCacheFactory;
 import org.apache.camel.support.LoggingExceptionHandler;
 import org.apache.camel.support.NoLock;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.apache.camel.support.UnitOfWorkHelper;
 import org.apache.camel.support.service.ServiceHelper;
 import org.apache.camel.util.ObjectHelper;
 import org.apache.camel.util.StopWatch;
@@ -135,6 +141,9 @@ public class AggregateProcessor extends BaseProcessorSupport
     // aggregated exchanges being completed (exchange id -> count): moved to 
the completed store of a recoverable
     // repository, but not yet registered in inProgressCompleteExchanges by 
onSubmitCompletion
     private final Map<String, Integer> completingExchanges = new 
ConcurrentHashMap<>();
+    // the references to spooled stream caches that the groups in the 
repository hold, by correlation key (only when
+    // isKeepingReferences, and guarded by the lock)
+    private final Map<String, SpooledStreamCaches> groupSpooledStreamCaches = 
new HashMap<>();
     private final Map<String, RedeliveryData> redeliveryState = new 
ConcurrentHashMap<>();
 
     private final AggregateProcessorStatistics statistics = new Statistics();
@@ -451,6 +460,31 @@ public class AggregateProcessor extends 
BaseProcessorSupport
         removeFlagCompleteAllGroups(copy);
         removeFlagCompleteAllGroupsInclusive(copy);
 
+        // a stream cache spooled to disk is deleted when the incoming 
exchange is done, but the group keeps the body
+        // for later, so the copy takes its own reference, which the 
aggregator releases when it is done with the body
+        SpooledStreamCaches spooled = null;
+        if (copy.getIn().getBody() instanceof StreamCache sc && 
!sc.inMemory()) {
+            // the copy is independent of the unit of work that the incoming 
exchange may release its stream caches with
+            // (such as the parent of a split), as the wire tap does
+            copy.removeProperty(ExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORK);
+            try {
+                StreamCache copied = sc.copy(copy);
+                if (copied != null) {
+                    copy.getIn().setBody(copied);
+                }
+            } catch (IOException e) {
+                UnitOfWorkHelper.doneSynchronizations(copy, 
copy.getExchangeExtension().handoverCompletions());
+                exchange.setException(e);
+                callback.done(sync);
+                return sync;
+            }
+            // the copy has no on completions of its own, so these are only 
the releases of the reference, which are
+            // kept apart from the on completions that the aggregation 
strategy may add to the exchanges (such as the
+            // ZipAggregationStrategy deleting its zip file when the 
aggregated exchange is done)
+            spooled = new SpooledStreamCaches();
+            spooled.add(copy.getExchangeExtension().handoverCompletions());
+        }
+
         List<Exchange> aggregated = new ArrayList<>();
         lock.lock();
         try {
@@ -459,11 +493,17 @@ public class AggregateProcessor extends 
BaseProcessorSupport
             if (closedCorrelationKeys != null && 
closedCorrelationKeys.containsKey(key)) {
                 throw new ClosedCorrelationKeyException(key, exchange);
             }
-            doAggregation(key, copy, aggregated);
+            doAggregation(key, copy, spooled, aggregated);
         } catch (CamelExchangeException e) {
             exchange.setException(e);
         } finally {
             lock.unlock();
+            if (spooled != null) {
+                // release the reference unless it was handed over to the 
group or the aggregated exchange: the copy
+                // failed, is retried with a new copy due to optimistic 
locking, or is discarded, or the repository has
+                // read the body when the copy was added (see 
isKeepingReferences)
+                spooled.onDone(copy);
+            }
             // we are completed so submit to completion outside the lock. This 
must also be done when the aggregation
             // failed, or must be retried due to optimistic locking, after a 
group was completed (such as a group
             // completed by pre-completion), as that group has already been 
removed from the repository
@@ -524,10 +564,13 @@ public class AggregateProcessor extends 
BaseProcessorSupport
      *
      * @param  key                                     the correlation key
      * @param  newExchange                             the exchange
+     * @param  spooled                                 the references that the 
exchange holds to spooled stream caches,
+     *                                                 which are taken when 
they are handed over, or <tt>null</tt>
      * @param  list                                    the list to add the 
aggregated exchange(s) which are complete to
      * @throws org.apache.camel.CamelExchangeException is thrown if error 
aggregating
      */
-    private void doAggregation(String key, Exchange newExchange, 
List<Exchange> list) throws CamelExchangeException {
+    private void doAggregation(String key, Exchange newExchange, 
SpooledStreamCaches spooled, List<Exchange> list)
+            throws CamelExchangeException {
         LOG.trace("onAggregation +++ start +++ with correlation key: {}", key);
 
         String complete = null;
@@ -667,9 +710,19 @@ public class AggregateProcessor extends 
BaseProcessorSupport
             }
             // only need to update aggregation repository if we are not 
complete
             doAggregationRepositoryAdd(newExchange.getContext(), key, 
originalExchange, answer);
+            if (spooled != null && isKeepingReferences()) {
+                // the group keeps the references until it is completed (see 
doOnCompletion)
+                groupSpooledStreamCaches.computeIfAbsent(key, k -> new 
SpooledStreamCaches()).add(spooled.take());
+            }
         } else {
             // if we are complete then add the answer to the list
             doAggregationComplete(complete, list, key, originalExchange, 
answer, aggregateFailed);
+            if (spooled != null && containsInstance(list, answer)) {
+                // the aggregated exchange takes over the references 
(onCompletion has handed over those of the group)
+                SpooledStreamCaches release = new SpooledStreamCaches();
+                release.add(spooled.take());
+                answer.getExchangeExtension().addOnCompletion(release);
+            }
         }
 
         LOG.trace("onAggregation +++  end  +++ with correlation key: {}", key);
@@ -711,6 +764,18 @@ public class AggregateProcessor extends 
BaseProcessorSupport
         }
     }
 
+    /**
+     * Whether the groups in the repository keep the references to the spooled 
stream caches of their exchanges until
+     * they are completed. This is only the case with the memory repository, 
which keeps the exchange instances (and
+     * their bodies), and without optimistic locking, as all access to the 
groups is then under the lock. A persistent
+     * repository has read the body when the exchange is added, so the 
reference is released after the add. With
+     * optimistic locking the reference is released after the add as well, and 
only the body of the exchange that
+     * completes the group is kept.
+     */
+    private boolean isKeepingReferences() {
+        return !optimisticLocking && aggregationRepository instanceof 
MemoryAggregationRepository;
+    }
+
     protected void doAggregationRepositoryAdd(
             CamelContext camelContext, String key, Exchange oldExchange, 
Exchange newExchange) {
         LOG.trace("In progress aggregated oldExchange: {}, newExchange: {} 
with correlation key: {}", oldExchange, newExchange,
@@ -894,6 +959,11 @@ public class AggregateProcessor extends 
BaseProcessorSupport
         if (original != null) {
             // remove from repository as its completed, we do this first as to 
trigger any OptimisticLockingException's
             aggregationRepository.remove(aggregated.getContext(), key, 
original);
+            // the aggregated exchange takes over the references that the 
group holds to spooled stream caches
+            SpooledStreamCaches spooled = takeGroupSpooledStreamCaches(key);
+            if (spooled != null) {
+                aggregated.getExchangeExtension().addOnCompletion(spooled);
+            }
         }
 
         // cleanup timeout map if it was a incoming exchange which triggered 
the timeout (and not the timeout checker)
@@ -949,9 +1019,82 @@ public class AggregateProcessor extends 
BaseProcessorSupport
         aggregationRepository.confirm(aggregated.getContext(), 
aggregated.getExchangeId());
         // and remove redelivery state as well
         redeliveryState.remove(aggregated.getExchangeId());
+        // and release the references that the discarded exchange holds to 
spooled stream caches (only those, as the
+        // other on completions of the exchange are left as they are when it 
is not discarded)
+        releaseSpooledStreamCaches(aggregated);
         // the completion was from timeout and we should just discard it
     }
 
+    /**
+     * Takes the references that the group holds to spooled stream caches, 
when the group is removed from the
+     * repository.
+     */
+    private SpooledStreamCaches takeGroupSpooledStreamCaches(String key) {
+        return isKeepingReferences() ? groupSpooledStreamCaches.remove(key) : 
null;
+    }
+
+    /**
+     * Releases the references that the aggregator handed over to a discarded 
exchange (see doOnCompletion).
+     */
+    private static void releaseSpooledStreamCaches(Exchange exchange) {
+        List<Synchronization> completions = 
exchange.getExchangeExtension().handoverCompletions();
+        if (completions != null) {
+            for (Synchronization completion : completions) {
+                if (completion instanceof SpooledStreamCaches spooled) {
+                    spooled.onDone(exchange);
+                } else {
+                    
exchange.getExchangeExtension().addOnCompletion(completion);
+                }
+            }
+        }
+    }
+
+    /**
+     * The on completions that release references that the aggregator took to 
spooled stream caches, so the bodies can
+     * be read until the aggregator is done with them. They are kept apart 
from the other on completions of the
+     * exchanges, which the aggregator leaves as they are. When added to an 
exchange as on completion, they are released
+     * when the exchange is done.
+     */
+    private static final class SpooledStreamCaches extends 
SynchronizationAdapter {
+        private List<Synchronization> releases;
+
+        void add(List<Synchronization> synchronizations) {
+            if (synchronizations != null && !synchronizations.isEmpty()) {
+                if (releases == null) {
+                    releases = new ArrayList<>(synchronizations);
+                } else {
+                    releases.addAll(synchronizations);
+                }
+            }
+        }
+
+        List<Synchronization> take() {
+            List<Synchronization> answer = releases;
+            releases = null;
+            return answer;
+        }
+
+        @Override
+        public void onDone(Exchange exchange) {
+            // release only once
+            UnitOfWorkHelper.doneSynchronizations(exchange, take());
+        }
+
+        @Override
+        public String toString() {
+            return "AggregateOnCompletion[SpooledStreamCaches]";
+        }
+    }
+
+    private static boolean containsInstance(List<Exchange> list, Exchange 
exchange) {
+        for (Exchange e : list) {
+            if (e == exchange) {
+                return true;
+            }
+        }
+        return false;
+    }
+
     private void onSubmitCompletion(final String key, final Exchange exchange) 
{
         LOG.debug("Aggregation complete for correlation key {} sending 
aggregated exchange: {}", key, exchange);
 
@@ -1902,6 +2045,18 @@ public class AggregateProcessor extends 
BaseProcessorSupport
         // shutdown aggregation repository and the strategy
         ServiceHelper.stopAndShutdownServices(aggregationRepository, 
aggregationStrategy);
 
+        // the memory repository has dropped the groups, so release the 
references they held to spooled stream caches
+        if (!groupSpooledStreamCaches.isEmpty()) {
+            lock.lock();
+            try {
+                Exchange dummy = new DefaultExchange(camelContext);
+                groupSpooledStreamCaches.values().forEach(spooled -> 
spooled.onDone(dummy));
+                groupSpooledStreamCaches.clear();
+            } finally {
+                lock.unlock();
+            }
+        }
+
         // cleanup when shutting down
         inProgressCompleteExchanges.clear();
         completingExchanges.clear();
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateStreamCachingSpoolTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateStreamCachingSpoolTest.java
new file mode 100644
index 000000000000..dc334841518b
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/aggregator/AggregateStreamCachingSpoolTest.java
@@ -0,0 +1,458 @@
+/*
+ * 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.aggregator;
+
+import java.io.BufferedInputStream;
+import java.io.ByteArrayInputStream;
+import java.io.File;
+import java.io.InputStream;
+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.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicInteger;
+
+import org.apache.camel.AggregationStrategy;
+import org.apache.camel.CamelContext;
+import org.apache.camel.CamelExecutionException;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.apache.camel.processor.aggregate.AggregateController;
+import org.apache.camel.processor.aggregate.ClosedCorrelationKeyException;
+import org.apache.camel.processor.aggregate.DefaultAggregateController;
+import org.apache.camel.processor.aggregate.GroupedBodyAggregationStrategy;
+import org.apache.camel.processor.aggregate.MemoryAggregationRepository;
+import org.apache.camel.processor.aggregate.UseLatestAggregationStrategy;
+import org.apache.camel.spi.AggregationRepository;
+import org.apache.camel.spi.OptimisticLockingAggregationRepository;
+import org.apache.camel.support.DefaultExchange;
+import org.apache.camel.support.SynchronizationAdapter;
+import org.apache.camel.support.service.ServiceSupport;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertArrayEquals;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertInstanceOf;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * The Aggregate EIP must keep the bodies that stream caching spooled to disk 
until the aggregated exchange is done with
+ * them, and must not leave the spool files behind, also when a group or an 
incoming exchange is discarded.
+ */
+public class AggregateStreamCachingSpoolTest extends ContextTestSupport {
+
+    private static final byte[] DATA = createData(16 * 1024);
+
+    // the aggregated exchange is processed only when the incoming exchanges 
are done, so the spool files would already
+    // be deleted if the aggregator did not hold its own references to them
+    private final CountDownLatch sendersDone = new CountDownLatch(1);
+    private final AggregateController controller = new 
DefaultAggregateController();
+    private final FailOnceRepository failOnce = new FailOnceRepository();
+    private final FailOnceRepository failOnceAndClose = new 
FailOnceRepository();
+    // whether the on completion that the aggregation strategy added to the 
first exchange of the group has run
+    private final AtomicBoolean strategyCompletionDone = new AtomicBoolean();
+
+    @Test
+    public void testGroupedBodies() throws Exception {
+        MockEndpoint result = getMockEndpoint("mock:result");
+        result.expectedMessageCount(1);
+
+        template.sendBody("direct:grouped", stream());
+        template.sendBody("direct:grouped", stream());
+        template.sendBody("direct:grouped", stream());
+        sendersDone.countDown();
+
+        assertMockEndpointsSatisfied();
+        List<?> bodies = 
result.getReceivedExchanges().get(0).getMessage().getBody(List.class);
+        assertEquals(3, bodies.size());
+        for (Object body : bodies) {
+            assertArrayEquals(DATA, (byte[]) body);
+        }
+        assertSpoolDirectoryEmpty();
+    }
+
+    @Test
+    public void testCompletionTimeout() throws Exception {
+        sendAndAssertBody("direct:timeout");
+    }
+
+    @Test
+    public void testCompletionSizeOne() throws Exception {
+        // the aggregated exchange is sent on the aggregator's own thread, so 
even the body of the exchange that
+        // completes the group must not depend on the incoming exchange
+        sendAndAssertBody("direct:size");
+    }
+
+    @Test
+    public void testOptimisticLockingRetry() throws Exception {
+        MockEndpoint result = getMockEndpoint("mock:result");
+        result.expectedMessageCount(1);
+
+        // the first attempt to add the first exchange fails, and the retry 
takes a new copy
+        template.sendBody("direct:optimistic", stream());
+        template.sendBody("direct:optimistic", stream());
+        sendersDone.countDown();
+
+        assertMockEndpointsSatisfied();
+        assertEquals(2, failOnce.calls.get());
+        assertArrayEquals(DATA, 
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+        assertSpoolDirectoryEmpty();
+    }
+
+    @Test
+    public void testClosedCorrelationKeyOnRetry() throws Exception {
+        MockEndpoint result = getMockEndpoint("mock:result");
+        result.expectedMessageCount(1);
+        sendersDone.countDown();
+
+        // while the first attempt to add this exchange fails, another 
exchange completes and closes the group, so the
+        // retry finds the correlation key closed
+        CamelExecutionException e = assertThrows(CamelExecutionException.class,
+                () -> template.sendBody("direct:closed", stream()));
+        assertInstanceOf(ClosedCorrelationKeyException.class, e.getCause());
+
+        assertMockEndpointsSatisfied();
+        assertArrayEquals(DATA, 
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+        assertSpoolDirectoryEmpty();
+    }
+
+    @Test
+    public void testDiscardOnAggregationFailure() throws Exception {
+        getMockEndpoint("mock:result").expectedMessageCount(0);
+
+        // the first exchange fails to aggregate, and is discarded
+        template.sendBodyAndHeader("direct:failure", stream(), "fail", true);
+        // the second exchange starts a group, and the third fails to 
aggregate, which discards the group
+        template.sendBody("direct:failure", stream());
+        template.sendBodyAndHeader("direct:failure", stream(), "fail", true);
+        sendersDone.countDown();
+
+        assertMockEndpointsSatisfied();
+        assertSpoolDirectoryEmpty();
+    }
+
+    @Test
+    public void testDiscardOnCompletionTimeout() throws Exception {
+        getMockEndpoint("mock:result").expectedMessageCount(0);
+
+        template.sendBody("direct:discardTimeout", stream());
+        sendersDone.countDown();
+
+        assertSpoolDirectoryEmpty();
+        assertMockEndpointsSatisfied();
+    }
+
+    @Test
+    public void testForceDiscardingOfGroup() throws Exception {
+        getMockEndpoint("mock:result").expectedMessageCount(0);
+
+        template.sendBody("direct:forceDiscard", stream());
+        template.sendBody("direct:forceDiscard", stream());
+        sendersDone.countDown();
+        assertEquals(1, controller.forceDiscardingOfGroup("group"));
+
+        assertMockEndpointsSatisfied();
+        assertSpoolDirectoryEmpty();
+    }
+
+    @Test
+    public void testRepositoryStoringCopies() throws Exception {
+        // like the persistent repositories, the repository reads the body 
when the exchange is added
+        sendAndAssertBody("direct:copies");
+    }
+
+    @Test
+    public void testMemoryRepositoryStoringCopies() throws Exception {
+        // a subclass of the memory repository that does not keep the exchange 
instance
+        sendAndAssertBody("direct:memoryCopies");
+    }
+
+    @Test
+    public void testStrategyCompletionWithOptimisticLocking() throws Exception 
{
+        sendAndAssertStrategyCompletion("direct:strategyCompletionOptimistic");
+    }
+
+    @Test
+    public void testStrategyCompletionWithRepositoryStoringCopies() throws 
Exception {
+        sendAndAssertStrategyCompletion("direct:strategyCompletionCopies");
+    }
+
+    @Test
+    public void testStrategyCompletionWithMemoryRepository() throws Exception {
+        sendAndAssertStrategyCompletion("direct:strategyCompletionMemory");
+    }
+
+    private void sendAndAssertStrategyCompletion(String uri) throws Exception {
+        // the aggregator must only release its own references to the spooled 
stream caches, and leave the on
+        // completions that the aggregation strategy adds to the exchanges 
alone (such as the ZipAggregationStrategy
+        // that deletes its zip file when the aggregated exchange is done)
+        MockEndpoint result = getMockEndpoint("mock:result");
+        result.expectedMessageCount(1);
+
+        template.sendBody(uri, stream());
+        template.sendBody(uri, stream());
+        template.sendBodyAndHeader(uri, stream(), "last", true);
+        sendersDone.countDown();
+
+        assertMockEndpointsSatisfied();
+        assertArrayEquals(DATA, 
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+        assertSpoolDirectoryEmpty();
+        if (uri.endsWith("Memory")) {
+            // the memory repository keeps the exchange, so the on completion 
runs when the aggregated exchange is done
+            Awaitility.await().atMost(5, 
TimeUnit.SECONDS).untilTrue(strategyCompletionDone);
+        }
+    }
+
+    private void sendAndAssertBody(String uri) throws Exception {
+        MockEndpoint result = getMockEndpoint("mock:result");
+        result.expectedMessageCount(1);
+
+        template.sendBody(uri, stream());
+        sendersDone.countDown();
+
+        assertMockEndpointsSatisfied();
+        assertArrayEquals(DATA, 
result.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+        assertSpoolDirectoryEmpty();
+    }
+
+    private void assertSpoolDirectoryEmpty() {
+        File spoolDir = testDirectory().toFile();
+        Awaitility.await().atMost(5, TimeUnit.SECONDS).untilAsserted(() -> {
+            String[] files = spoolDir.list();
+            assertNotNull(files);
+            assertEquals(0, files.length, "Spool files left behind: " + 
List.of(files));
+        });
+    }
+
+    private static InputStream stream() {
+        // a stream that is not converted to an in-memory cache, so it is 
spooled to disk
+        return new BufferedInputStream(new ByteArrayInputStream(DATA));
+    }
+
+    private static byte[] createData(int size) {
+        byte[] data = new byte[size];
+        for (int i = 0; i < size; i++) {
+            data[i] = (byte) ('a' + i % 26);
+        }
+        return data;
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                
context.getStreamCachingStrategy().setSpoolDirectory(testDirectory().toFile());
+                context.getStreamCachingStrategy().setSpoolEnabled(true);
+                context.getStreamCachingStrategy().setSpoolThreshold(1024);
+                
context.getStreamCachingStrategy().setRemoveSpoolDirectoryWhenStopping(false);
+                context.setStreamCaching(true);
+
+                // reads the bodies after the incoming exchanges are done
+                Processor read = e -> {
+                    sendersDone.await(20, TimeUnit.SECONDS);
+                    Object body = e.getMessage().getBody();
+                    if (body instanceof List<?> list) {
+                        List<byte[]> answer = new ArrayList<>();
+                        for (Object o : list) {
+                            
answer.add(e.getContext().getTypeConverter().mandatoryConvertTo(byte[].class, 
e, o));
+                        }
+                        e.getMessage().setBody(answer);
+                    } else {
+                        
e.getMessage().setBody(e.getMessage().getMandatoryBody(byte[].class));
+                    }
+                };
+
+                // adds an on completion to the first exchange of the group, 
which must not run before the group completes
+                AggregationStrategy addCompletion = (oldExchange, newExchange) 
-> {
+                    if (oldExchange == null) {
+                        strategyCompletionDone.set(false);
+                        newExchange.getExchangeExtension().addOnCompletion(new 
SynchronizationAdapter() {
+                            @Override
+                            public void onDone(Exchange exchange) {
+                                strategyCompletionDone.set(true);
+                            }
+                        });
+                        return newExchange;
+                    }
+                    if (strategyCompletionDone.get()) {
+                        throw new IllegalStateException("The on completion of 
the aggregation strategy has already run");
+                    }
+                    
oldExchange.getMessage().setBody(newExchange.getMessage().getBody());
+                    return oldExchange;
+                };
+
+                AggregationStrategy failOnHeader = (oldExchange, newExchange) 
-> {
+                    if (newExchange.getMessage().getHeader("fail") != null) {
+                        throw new IllegalArgumentException("Forced");
+                    }
+                    return newExchange;
+                };
+
+                from("direct:grouped")
+                        .aggregate(constant("group"), new 
GroupedBodyAggregationStrategy()).completionSize(3)
+                        .process(read).to("mock:result");
+
+                from("direct:timeout")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy())
+                        
.completionTimeout(100).completionTimeoutCheckerInterval(10)
+                        .process(read).to("mock:result");
+
+                from("direct:size")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy()).completionSize(1)
+                        .process(read).to("mock:result");
+
+                from("direct:optimistic")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy())
+                        
.aggregationRepository(failOnce).optimisticLocking().completionSize(2)
+                        .process(read).to("mock:result");
+
+                from("direct:closed")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy())
+                        
.aggregationRepository(failOnceAndClose).optimisticLocking()
+                        
.completionPredicate(header("complete").isNotNull()).closeCorrelationKeyOnCompletion(100)
+                        .process(read).to("mock:result");
+
+                from("direct:failure")
+                        .aggregate(constant("group"), 
failOnHeader).discardOnAggregationFailure().completionSize(5)
+                        .to("mock:result");
+
+                from("direct:discardTimeout")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy())
+                        
.completionTimeout(100).completionTimeoutCheckerInterval(10).discardOnCompletionTimeout()
+                        .to("mock:result");
+
+                from("direct:forceDiscard")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy()).aggregateController(controller)
+                        .completionSize(5)
+                        .to("mock:result");
+
+                from("direct:strategyCompletionOptimistic")
+                        .aggregate(constant("group"), addCompletion)
+                        .aggregationRepository(new 
MemoryAggregationRepository(true)).optimisticLocking()
+                        
.completionPredicate(header("last").isNotNull()).eagerCheckCompletion()
+                        .process(read).to("mock:result");
+
+                from("direct:strategyCompletionCopies")
+                        .aggregate(constant("group"), addCompletion)
+                        .aggregationRepository(new CopyingRepository())
+                        
.completionPredicate(header("last").isNotNull()).eagerCheckCompletion()
+                        .process(read).to("mock:result");
+
+                from("direct:strategyCompletionMemory")
+                        .aggregate(constant("group"), addCompletion)
+                        
.completionPredicate(header("last").isNotNull()).eagerCheckCompletion()
+                        .process(read).to("mock:result");
+
+                from("direct:memoryCopies")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy())
+                        .aggregationRepository(new 
MemoryAggregationRepository() {
+                            @Override
+                            public Exchange add(CamelContext camelContext, 
String key, Exchange exchange) {
+                                Exchange copy = exchange.copy();
+                                
copy.getMessage().setBody(exchange.getMessage().getBody(byte[].class));
+                                return super.add(camelContext, key, copy);
+                            }
+                        })
+                        
.completionTimeout(100).completionTimeoutCheckerInterval(10)
+                        .process(read).to("mock:result");
+
+                from("direct:copies")
+                        .aggregate(constant("group"), new 
UseLatestAggregationStrategy())
+                        .aggregationRepository(new CopyingRepository())
+                        
.completionTimeout(100).completionTimeoutCheckerInterval(10)
+                        .process(read).to("mock:result");
+            }
+        };
+    }
+
+    /**
+     * An optimistic locking repository whose first add fails. When it closes 
the group, it first sends an exchange that
+     * completes the group, which closes the correlation key.
+     */
+    private final class FailOnceRepository extends MemoryAggregationRepository 
{
+        private final AtomicInteger calls = new AtomicInteger();
+
+        private FailOnceRepository() {
+            super(true);
+        }
+
+        @Override
+        public Exchange add(CamelContext camelContext, String key, Exchange 
oldExchange, Exchange newExchange) {
+            if (calls.incrementAndGet() == 1) {
+                if (this == failOnceAndClose) {
+                    template.sendBodyAndHeader("direct:closed", stream(), 
"complete", true);
+                }
+                throw new 
OptimisticLockingAggregationRepository.OptimisticLockingException();
+            }
+            return super.add(camelContext, key, oldExchange, newExchange);
+        }
+    }
+
+    /**
+     * A repository that stores a copy of the exchange with the body read into 
memory, as the persistent repositories
+     * store a serialized copy, and returns a new exchange instance on every 
get.
+     */
+    private static final class CopyingRepository extends ServiceSupport 
implements AggregationRepository {
+        private final Map<String, Exchange> exchanges = new 
ConcurrentHashMap<>();
+
+        @Override
+        public Exchange add(CamelContext camelContext, String key, Exchange 
exchange) {
+            Exchange old = exchanges.put(key, copy(exchange, 
exchange.getMessage().getBody(byte[].class)));
+            return old != null ? copy(old, old.getMessage().getBody()) : null;
+        }
+
+        @Override
+        public Exchange get(CamelContext camelContext, String key) {
+            Exchange exchange = exchanges.get(key);
+            return exchange != null ? copy(exchange, 
exchange.getMessage().getBody()) : null;
+        }
+
+        @Override
+        public void remove(CamelContext camelContext, String key, Exchange 
exchange) {
+            exchanges.remove(key);
+        }
+
+        @Override
+        public void confirm(CamelContext camelContext, String exchangeId) {
+            // noop
+        }
+
+        @Override
+        public Set<String> getKeys() {
+            return exchanges.keySet();
+        }
+
+        private static Exchange copy(Exchange exchange, Object body) {
+            Exchange copy = new DefaultExchange(exchange.getContext());
+            copy.setExchangeId(exchange.getExchangeId());
+            copy.getProperties().putAll(exchange.getProperties());
+            copy.getMessage().setHeaders(exchange.getMessage().getHeaders());
+            copy.getMessage().setBody(body);
+            return copy;
+        }
+    }
+}
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 7e4a2cd8edcb..89c1cda4e272 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
@@ -211,6 +211,22 @@ not complete by timeout until another exchange that has a 
completion timeout arr
 `AggregationRepository` that does not keep exchange properties never completes 
a group by timeout. This only affects
 optimistic locking combined with `completionTimeoutExpression`.
 
+=== Aggregate EIP - stream caching with spooling to disk
+
+When stream caching spools message bodies to disk, the aggregator now keeps 
the spool files of the bodies it
+aggregates until the aggregated exchange is done. Prior to Camel 4.23 a spool 
file was deleted as soon as the incoming
+exchange was done, so the aggregated exchange could fail to read the bodies 
with a `NoSuchFileException`.
+
+With the default `MemoryAggregationRepository` (without optimistic locking) 
the spool files of all exchanges of a group
+are now kept until the group completes and the aggregated exchange is done, so 
more disk space can be used in the spool
+directory than before. For example `UseLatestAggregationStrategy` with a large 
`completionSize` now keeps a spool file
+for every exchange of the group, although only the last body is used. A group 
that never completes keeps its spool
+files until the aggregator is shut down, such as when the route is removed or 
`CamelContext` is stopped.
+
+With optimistic locking, or a persistent aggregation repository, the spool 
file of an exchange is released when the
+exchange has been added to the repository (a persistent repository has read 
the body at that time), and only the
+spool file of the exchange that completes a group is kept until the aggregated 
exchange is done.
+
 === Variable Receive
 
 When an EIP with `variableReceive` stores a message into a variable that 
already holds a message, the header variables

Reply via email to