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 b0fed420c974 CAMEL-25100: camel-core - OnCompletion EIP: read a body 
that stream caching spooled to disk (#27002)
b0fed420c974 is described below

commit b0fed420c974fc25033d3edc9ec55700a96a3114
Author: allthingssecurity <[email protected]>
AuthorDate: Thu Oct 1 11:17:07 2026 +0530

    CAMEL-25100: camel-core - OnCompletion EIP: read a body that stream caching 
spooled to disk (#27002)
    
    The spool file of a FileInputStreamCache is deleted by an on completion
    with the default order 0. The onCompletion EIP in the default
    (after consumer) mode is an on completion with the order LOWEST, and the
    unit of work runs them sorted by order. So the file was always deleted
    before the onCompletion copied the exchange: reading the body failed
    with NoSuchFileException, and a parallel onCompletion stopped at the
    first step that read the body without logging anything.
    
    - The stream cache clean-up now has the order LOWEST, so it runs after
      the other on completions of the exchange. The after consumer
      onCompletion uses LOWEST - 1, just before it.
    - prepareExchange gives the onCompletion's copy its own reference to the
      stream cache (as the Wire Tap EIP does), so the file lives until the
      onCompletion route is done, also with parallelProcessing.
    - With parallelProcessing, a task that never runs (rejected, discarded
      by a pool that is shut down, or dropped by shutdownNow) releases the
      copy's reference and is no longer counted as pending for the graceful
      shutdown (CAMEL-25012). The pool is checked before submitting, so a
      task accepted before a graceful shutdown still runs.
    
    Adds OnCompletionStreamCachingSpoolTest, and a 4.23 upgrade guide note
    about the new order of the stream cache clean-up.
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../camel/processor/OnCompletionProcessor.java     | 112 ++++++--
 .../OnCompletionStreamCachingSpoolTest.java        | 292 +++++++++++++++++++++
 .../converter/stream/FileInputStreamCache.java     |   8 +
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  15 ++
 4 files changed, 406 insertions(+), 21 deletions(-)

diff --git 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
index ac2d5de89da2..63fb2f3b01a8 100644
--- 
a/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
+++ 
b/core/camel-core-processor/src/main/java/org/apache/camel/processor/OnCompletionProcessor.java
@@ -16,8 +16,11 @@
  */
 package org.apache.camel.processor;
 
+import java.io.IOException;
 import java.util.ArrayList;
 import java.util.List;
+import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ExecutorService;
 import java.util.concurrent.atomic.LongAdder;
 
@@ -32,6 +35,7 @@ import org.apache.camel.Predicate;
 import org.apache.camel.Processor;
 import org.apache.camel.Route;
 import org.apache.camel.ShutdownRunningTask;
+import org.apache.camel.StreamCache;
 import org.apache.camel.Traceable;
 import org.apache.camel.spi.IdAware;
 import org.apache.camel.spi.RouteIdAware;
@@ -40,6 +44,7 @@ import org.apache.camel.spi.StepIdAware;
 import org.apache.camel.spi.SynchronizationRouteAware;
 import org.apache.camel.support.ExchangeHelper;
 import org.apache.camel.support.SynchronizationAdapter;
+import org.apache.camel.support.UnitOfWorkHelper;
 import org.apache.camel.support.service.ServiceHelper;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
@@ -69,6 +74,8 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
     private final boolean afterConsumer;
     private final boolean routeScoped;
     private final LongAdder taskCount = new LongAdder();
+    // the parallel onCompletion tasks that have been submitted but have not 
started yet
+    private final Set<ParallelTask> pendingTasks = 
ConcurrentHashMap.newKeySet();
 
     public OnCompletionProcessor(CamelContext camelContext, Processor 
processor, ExecutorService executorService,
                                  boolean shutdownExecutorService,
@@ -112,10 +119,12 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
     protected void doShutdown() throws Exception {
         ServiceHelper.stopAndShutdownService(processor);
         if (shutdownExecutorService) {
-            List<Runnable> dropped = 
getCamelContext().getExecutorServiceManager().shutdownNow(executorService);
-            if (dropped != null && !dropped.isEmpty()) {
-                // the tasks still queued in the thread pool will never run, 
so they are no longer pending
-                taskCount.add(-dropped.size());
+            
getCamelContext().getExecutorServiceManager().shutdownNow(executorService);
+            // the tasks that have not started (the tasks still queued in the 
thread pool, which shutdownNow dropped, and
+            // a task a thread has taken from the queue but not started) will 
never run, so they are no longer pending,
+            // and what their copies hold is released
+            for (ParallelTask task : pendingTasks) {
+                task.discard();
             }
         }
     }
@@ -188,27 +197,74 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
     }
 
     /**
-     * Submits the onCompletion task to the thread pool (parallel processing). 
The task is counted as pending from when
-     * it is submitted until it is done, so a graceful shutdown waits for it.
+     * Submits the onCompletion task of the given copy to the thread pool 
(parallel processing). The task is counted as
+     * pending from when it is submitted until it is done, so a graceful 
shutdown waits for it.
+     * <p>
+     * The copy may hold its own reference to a stream cache (see {@link 
#prepareExchange(Exchange)}), which is released
+     * when the copy is done. When the task never runs (the thread pool 
rejects it, discards it because it is shut down,
+     * or drops it when the processor shuts its thread pool down), it is no 
longer counted as pending and the reference
+     * is released instead, as otherwise a spooled file would be kept until 
the stream caching strategy is stopped.
+     * <p>
+     * A thread pool that is shut down rejects every task, and may do so 
without failing (such as with the CallerRuns or
+     * Discard policy), so the task is discarded when the thread pool was 
already shut down when the task was submitted.
+     * A thread pool that is shut down (gracefully) after it has accepted the 
task still runs it, so the task must then
+     * not be discarded. In the rare case that the thread pool is shut down 
while the task is being submitted, and
+     * silently rejects it, the task is left as pending, and is discarded when 
the processor shuts its thread pool down.
      */
     @SuppressWarnings("deprecation")
-    private void submitTask(Runnable task) {
+    private void submitTask(Exchange copy, Runnable task) {
+        ParallelTask parallelTask = new ParallelTask(copy, task);
         taskCount.increment();
-        Runnable counted = () -> {
-            try {
-                task.run();
-            } finally {
-                taskCount.decrement();
-            }
-        };
+        pendingTasks.add(parallelTask);
+        // check before submitting, as the thread pool may be shut down by 
another thread after it accepted the task
+        boolean shutdown = executorService.isShutdown();
         try {
             // Deprecated since 4.19.0
-            executorService.submit(prepareMDCParallelTask(camelContext, 
counted));
+            executorService.submit(prepareMDCParallelTask(camelContext, 
parallelTask));
         } catch (RuntimeException e) {
             // the task will not run
-            taskCount.decrement();
+            parallelTask.discard();
             throw e;
         }
+        if (shutdown) {
+            // the task was rejected, maybe without failing (unless the 
rejection policy has run it, and then this is a
+            // noop)
+            parallelTask.discard();
+        }
+    }
+
+    /**
+     * A parallel onCompletion task. Either the thread pool runs it, or it is 
discarded because it will not run, never
+     * both: whichever comes first removes it from the pending tasks. Both 
stop counting it as pending.
+     */
+    private final class ParallelTask implements Runnable {
+
+        private final Exchange copy;
+        private final Runnable task;
+
+        private ParallelTask(Exchange copy, Runnable task) {
+            this.copy = copy;
+            this.task = task;
+        }
+
+        @Override
+        public void run() {
+            // a task that was discarded (when the thread pool was shut down) 
must not be processed
+            if (pendingTasks.remove(this)) {
+                try {
+                    task.run();
+                } finally {
+                    taskCount.decrement();
+                }
+            }
+        }
+
+        void discard() {
+            if (pendingTasks.remove(this)) {
+                taskCount.decrement();
+                UnitOfWorkHelper.doneSynchronizations(copy, 
copy.getExchangeExtension().handoverCompletions());
+            }
+        }
     }
 
     protected boolean isCreateCopy() {
@@ -309,6 +365,20 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
             }
             // set MEP to InOnly as this onCompletion is a fire and forget
             answer.setPattern(ExchangePattern.InOnly);
+            // the copy is routed when the original exchange is done (or in 
parallel with its completion), so it must
+            // hold its own reference to a stream cache, as a spooled file is 
deleted when the original exchange is
+            // done (same as the Wire Tap EIP does)
+            
answer.removeProperty(ExchangePropertyKey.STREAM_CACHE_UNIT_OF_WORK);
+            if (answer.getIn().getBody() instanceof StreamCache sc) {
+                try {
+                    StreamCache copied = sc.copy(answer);
+                    if (copied != null) {
+                        answer.getIn().setBody(copied);
+                    }
+                } catch (IOException e) {
+                    answer.setException(e);
+                }
+            }
         } else {
             // use the exchange as-is
             answer = exchange;
@@ -339,8 +409,8 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
 
         @Override
         public int getOrder() {
-            // we want to be last
-            return Ordered.LOWEST;
+            // we want to be last, but before the stream cache clean up 
(Ordered.LOWEST), so we can read a spooled body
+            return Ordered.LOWEST - 1;
         }
 
         @Override
@@ -381,7 +451,7 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
                     LOG.debug("Processing onComplete: {}", copy);
                     doProcess(processor, copy);
                 };
-                submitTask(task);
+                submitTask(copy, task);
             } else {
                 // run without thread-pool
                 LOG.debug("Processing onComplete: {}", copy);
@@ -411,7 +481,7 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
                     // restore exception after processing
                     copy.setException(original);
                 };
-                submitTask(task);
+                submitTask(copy, task);
             } else {
                 // run without thread-pool
                 LOG.debug("Processing onFailure: {}", copy);
@@ -538,7 +608,7 @@ public class OnCompletionProcessor extends 
BaseProcessorSupport
                             LOG.debug("Processing onAfterRoute: {}", copy);
                             doProcess(processor, copy);
                         };
-                        submitTask(task);
+                        submitTask(copy, task);
                     } else {
                         // run without thread-pool
                         LOG.debug("Processing onAfterRoute: {}", copy);
diff --git 
a/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionStreamCachingSpoolTest.java
 
b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionStreamCachingSpoolTest.java
new file mode 100644
index 000000000000..3a80bbab3907
--- /dev/null
+++ 
b/core/camel-core/src/test/java/org/apache/camel/processor/OnCompletionStreamCachingSpoolTest.java
@@ -0,0 +1,292 @@
+/*
+ * 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;
+
+import java.io.BufferedInputStream;
+import java.io.ByteArrayInputStream;
+import java.io.File;
+import java.io.InputStream;
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.LinkedBlockingQueue;
+import java.util.concurrent.ThreadPoolExecutor;
+import java.util.concurrent.TimeUnit;
+
+import org.apache.camel.CamelExecutionException;
+import org.apache.camel.ContextTestSupport;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.builder.ThreadPoolProfileBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.awaitility.Awaitility;
+import org.junit.jupiter.api.AfterEach;
+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.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * The onCompletion EIP must be able to read a body that stream caching 
spooled to disk, and the spool file must be
+ * deleted once the onCompletion is done.
+ */
+public class OnCompletionStreamCachingSpoolTest extends ContextTestSupport {
+
+    private static final byte[] DATA = createData(16 * 1024);
+
+    // the parallel onCompletion waits until the original exchange is done
+    private final CountDownLatch originalDone = new CountDownLatch(1);
+    private final CountDownLatch blockingTaskStarted = new CountDownLatch(1);
+    private final CountDownLatch releaseBlockingTask = new CountDownLatch(1);
+
+    private final ExecutorService rejecting = 
Executors.newSingleThreadExecutor();
+    private final ExecutorService discarding = new ThreadPoolExecutor(
+            1, 1, 0, TimeUnit.SECONDS, new LinkedBlockingQueue<>(), new 
ThreadPoolExecutor.DiscardPolicy());
+    // a thread pool that another thread shuts down (gracefully) right after 
it has accepted a task
+    private final ExecutorService shutdownAfterSubmit = new ThreadPoolExecutor(
+            1, 1, 0, TimeUnit.SECONDS, new LinkedBlockingQueue<>(), new 
ThreadPoolExecutor.DiscardPolicy()) {
+        @Override
+        public Future<?> submit(Runnable task) {
+            Future<?> answer = super.submit(task);
+            shutdown();
+            return answer;
+        }
+    };
+
+    @Override
+    @AfterEach
+    public void tearDown() throws Exception {
+        releaseBlockingTask.countDown();
+        super.tearDown();
+        rejecting.shutdownNow();
+        discarding.shutdownNow();
+        shutdownAfterSubmit.shutdownNow();
+    }
+
+    @Test
+    public void testOnCompletion() throws Exception {
+        sendAndAssertReadByOnCompletion("direct:after");
+    }
+
+    @Test
+    public void testOnCompletionParallel() throws Exception {
+        sendAndAssertReadByOnCompletion("direct:parallel");
+    }
+
+    @Test
+    public void testOnCompletionBeforeConsumer() throws Exception {
+        // runs before the original exchange is done, so this has always worked
+        sendAndAssertReadByOnCompletion("direct:before");
+    }
+
+    @Test
+    public void testOnCompletionOnFailureOnly() throws Exception {
+        MockEndpoint done = getMockEndpoint("mock:done");
+        done.expectedMessageCount(1);
+
+        assertThrows(CamelExecutionException.class, () -> 
template.sendBody("direct:failure", stream()));
+
+        assertMockEndpointsSatisfied();
+        assertArrayEquals(DATA, 
done.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+        assertNoSpoolFiles();
+    }
+
+    @Test
+    public void testNoOnCompletion() throws Exception {
+        // the spool file is still deleted when the exchange is done
+        getMockEndpoint("mock:result").expectedMessageCount(1);
+
+        template.sendBody("direct:plain", stream());
+
+        assertMockEndpointsSatisfied();
+        assertNoSpoolFiles();
+    }
+
+    @Test
+    public void testParallelTaskRejected() throws Exception {
+        // the thread pool is shut down and rejects the onCompletion task with 
an exception
+        rejecting.shutdown();
+        getMockEndpoint("mock:result").expectedMessageCount(1);
+        getMockEndpoint("mock:done").expectedMessageCount(0);
+
+        template.sendBody("direct:rejected", stream());
+
+        assertMockEndpointsSatisfied();
+        assertNoSpoolFiles();
+        assertNoPendingTasks("rejected");
+    }
+
+    @Test
+    public void testParallelTaskDiscarded() throws Exception {
+        // the thread pool is shut down and discards the onCompletion task 
without an exception
+        discarding.shutdown();
+        getMockEndpoint("mock:result").expectedMessageCount(1);
+        getMockEndpoint("mock:done").expectedMessageCount(0);
+
+        template.sendBody("direct:discarded", stream());
+
+        assertMockEndpointsSatisfied();
+        assertNoSpoolFiles();
+        assertNoPendingTasks("discarded");
+    }
+
+    @Test
+    public void testParallelTaskAcceptedBeforeShutdown() throws Exception {
+        // the thread pool is shut down after it accepted the onCompletion 
task, and a graceful shutdown still runs the
+        // task, so it must not be discarded
+        MockEndpoint done = getMockEndpoint("mock:done");
+        done.expectedMessageCount(1);
+        getMockEndpoint("mock:result").expectedMessageCount(1);
+
+        template.sendBody("direct:shutdownAfterSubmit", stream());
+
+        assertMockEndpointsSatisfied();
+        assertArrayEquals(DATA, 
done.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+        assertTrue(shutdownAfterSubmit.isShutdown());
+        assertNoSpoolFiles();
+        Awaitility.await().atMost(5, TimeUnit.SECONDS)
+                .untilAsserted(() -> 
assertNoPendingTasks("shutdownAfterSubmit"));
+    }
+
+    @Test
+    public void testParallelTaskDroppedAtShutdown() throws Exception {
+        getMockEndpoint("mock:result").expectedMessageCount(2);
+
+        // the first onCompletion task blocks the only thread of the pool, the 
second one waits in its queue
+        template.sendBody("direct:dropped", stream());
+        assertTrue(blockingTaskStarted.await(10, TimeUnit.SECONDS));
+        template.sendBody("direct:dropped", stream());
+        assertMockEndpointsSatisfied();
+        OnCompletionProcessor onCompletion = context.getProcessor("dropped", 
OnCompletionProcessor.class);
+        assertEquals(2, onCompletion.getPendingExchangesSize());
+
+        // the graceful shutdown times out, and the thread pool is shut down 
with shutdownNow: the queued task never
+        // runs
+        context.getShutdownStrategy().setTimeout(1);
+        context.stop();
+
+        assertNoSpoolFiles();
+        Awaitility.await().atMost(5, TimeUnit.SECONDS)
+                .untilAsserted(() -> assertEquals(0, 
onCompletion.getPendingExchangesSize()));
+        assertEquals(0, getMockEndpoint("mock:done").getReceivedCounter(), 
"The queued onCompletion should not have run");
+    }
+
+    private void sendAndAssertReadByOnCompletion(String uri) throws Exception {
+        MockEndpoint done = getMockEndpoint("mock:done");
+        done.expectedMessageCount(1);
+        getMockEndpoint("mock:result").expectedMessageCount(1);
+
+        template.sendBody(uri, stream());
+        originalDone.countDown();
+
+        assertMockEndpointsSatisfied();
+        assertArrayEquals(DATA, 
done.getReceivedExchanges().get(0).getMessage().getBody(byte[].class));
+        assertNoSpoolFiles();
+    }
+
+    private void assertNoPendingTasks(String id) {
+        // the task that never runs is no longer counted as pending, so a 
graceful shutdown does not wait for it
+        assertEquals(0, context.getProcessor(id, 
OnCompletionProcessor.class).getPendingExchangesSize());
+    }
+
+    private void assertNoSpoolFiles() {
+        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);
+
+                context.getExecutorServiceManager().registerThreadPoolProfile(
+                        new 
ThreadPoolProfileBuilder("singleThread").poolSize(1).maxPoolSize(1).maxQueueSize(10).build());
+
+                from("direct:after")
+                        
.onCompletion().convertBodyTo(byte[].class).to("mock:done").end()
+                        .to("mock:result");
+
+                from("direct:parallel")
+                        .onCompletion().parallelProcessing()
+                        .process(e -> originalDone.await(20, TimeUnit.SECONDS))
+                        .convertBodyTo(byte[].class).to("mock:done").end()
+                        .to("mock:result");
+
+                from("direct:before")
+                        
.onCompletion().modeBeforeConsumer().convertBodyTo(byte[].class).to("mock:done").end()
+                        .to("mock:result");
+
+                from("direct:failure")
+                        
.onCompletion().onFailureOnly().convertBodyTo(byte[].class).to("mock:done").end()
+                        .to("mock:result")
+                        .throwException(new 
IllegalArgumentException("Forced"));
+
+                from("direct:plain")
+                        .to("mock:result");
+
+                from("direct:rejected")
+                        
.onCompletion().id("rejected").parallelProcessing().executorService(rejecting).to("mock:done").end()
+                        .to("mock:result");
+
+                from("direct:discarded")
+                        
.onCompletion().id("discarded").parallelProcessing().executorService(discarding).to("mock:done").end()
+                        .to("mock:result");
+
+                from("direct:shutdownAfterSubmit")
+                        
.onCompletion().id("shutdownAfterSubmit").parallelProcessing().executorService(shutdownAfterSubmit)
+                        .convertBodyTo(byte[].class).to("mock:done").end()
+                        .to("mock:result");
+
+                from("direct:dropped")
+                        
.onCompletion().id("dropped").parallelProcessing().executorService("singleThread")
+                        .process(e -> {
+                            blockingTaskStarted.countDown();
+                            releaseBlockingTask.await(20, TimeUnit.SECONDS);
+                        })
+                        .to("mock:done").end()
+                        .to("mock:result");
+            }
+        };
+    }
+}
diff --git 
a/core/camel-support/src/main/java/org/apache/camel/converter/stream/FileInputStreamCache.java
 
b/core/camel-support/src/main/java/org/apache/camel/converter/stream/FileInputStreamCache.java
index d1d16a403aa7..c125669e5a36 100644
--- 
a/core/camel-support/src/main/java/org/apache/camel/converter/stream/FileInputStreamCache.java
+++ 
b/core/camel-support/src/main/java/org/apache/camel/converter/stream/FileInputStreamCache.java
@@ -36,6 +36,7 @@ import javax.crypto.CipherOutputStream;
 
 import org.apache.camel.Exchange;
 import org.apache.camel.ExchangePropertyKey;
+import org.apache.camel.Ordered;
 import org.apache.camel.RuntimeCamelException;
 import org.apache.camel.StreamCache;
 import org.apache.camel.spi.StreamCachingStrategy;
@@ -257,6 +258,13 @@ public final class FileInputStreamCache extends 
InputStream implements StreamCac
                 exchangeCounter.incrementAndGet();
                 // add on completion so we can cleanup after the exchange is 
done such as deleting temporary files
                 Synchronization onCompletion = new SynchronizationAdapter() {
+                    @Override
+                    public int getOrder() {
+                        // delete the file after the other on completions 
(such as the onCompletion EIP), which may
+                        // still read the body
+                        return Ordered.LOWEST;
+                    }
+
                     @Override
                     public void onDone(Exchange exchange) {
                         int actualExchanges = 
exchangeCounter.decrementAndGet();
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 fa1e19bcdef1..0f7c35cf6615 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
@@ -155,6 +155,21 @@ the first name a header is set with is the one that is 
kept. Code that iterates
 case-sensitively, such as `Exchange.CONTENT_TYPE.equals(key)` or 
`key.startsWith("Camel")`, no longer matches
 `content-type` or `camelfilename`, as before Camel 4.21.
 
+=== camel-core - stream caching deletes a spooled file after the other on 
completions
+
+When stream caching spools a message body to disk, the temporary file is 
deleted by an on completion
+(`Synchronization`) of the exchange. This on completion now has the order 
`Ordered.LOWEST`, instead of the default
+order `0`, so it runs after the other on completions of the exchange. The file 
is therefore deleted later than
+before, never earlier. On completions with the same order (such as the 
disconnect of the FTP and SMB consumers, and
+the close of a JPA `EntityManager`) still run in the reverse order in which 
they were added.
+
+The `onCompletion` EIP in the default mode (`modeAfterConsumer`) now uses the 
order `Ordered.LOWEST - 1`, so it runs
+just before the file is deleted, and its copy of the exchange holds its own 
reference to a spooled body. An
+onCompletion can therefore read a spooled body, also with 
`parallelProcessing`. Previously it failed with a
+`NoSuchFileException`, and a parallel onCompletion stopped at the first step 
that read the body. As a side effect,
+an on completion with the order `Ordered.LOWEST` that was added after the 
onCompletion EIP (for example the close
+of a JPA `EntityManager` used later in the route) used to run before the 
onCompletion, and now runs after it.
+
 === @PropertyInject - an invalid property value is an error (Breaking change)
 
 When a property injected with `@PropertyInject` has a value that cannot be 
converted to the type of the field or

Reply via email to