allthingssecurity opened a new pull request, #27002:
URL: https://github.com/apache/camel/pull/27002

   # Description
   
   [CAMEL-25100](https://issues.apache.org/jira/browse/CAMEL-25100)
   
   With stream caching spooled to disk, an `onCompletion()` in the default mode 
(`modeAfterConsumer`) cannot read the body: the spool file is always deleted 
before the onCompletion runs.
   
   ```
   from("direct:start")
       .onCompletion().convertBodyTo(byte[].class).to("mock:done").end()
       .to("mock:result");
   
   onCompletion: NoSuchFileException for the spool file
   onCompletion().parallelProcessing(): the onCompletion route stops at the 
first step that reads the body, nothing is logged
   onCompletion().modeBeforeConsumer(): works
   ```
   
   The spool file is deleted by the on completion that 
`FileInputStreamCache.TempFileManager.addExchange` registers. It is a 
`SynchronizationAdapter` with the default order `0`. The after-consumer 
onCompletion is also an on completion of the unit of work 
(`OnCompletionSynchronizationAfterConsumer`), with the order `Ordered.LOWEST`, 
and `UnitOfWorkHelper` sorts the on completions by order. So the file is 
deleted first, and the onCompletion then copies an exchange whose body points 
to a deleted file. This is long-standing (the orders have not changed in 4.x).
   
   This change:
   - The stream cache clean-up on completion now has the order 
`Ordered.LOWEST`, so it runs after the other on completions of the exchange.
   - `OnCompletionSynchronizationAfterConsumer` now has the order 
`Ordered.LOWEST - 1`, just before the clean-up.
   - `prepareExchange` gives the onCompletion's copy its own reference to the 
stream cache: it removes `CamelStreamCacheUnitOfWork` from the copy and 
replaces the body with `sc.copy(copy)`, as the Wire Tap EIP does (CAMEL-12108). 
The copy then keeps the file until the onCompletion route is done. That is what 
makes `parallelProcessing` work, because the task can run after the original 
exchange is done. This only happens when a copy is made (after-consumer mode, 
or parallel processing). The before-consumer mode without parallel processing 
uses the exchange as is, and is unchanged.
   - With `parallelProcessing`, the copy's reference is released when its task 
never runs. That is the case when the thread pool rejects the task with an 
exception, when a thread pool that is shut down discards it silently (for 
example the default `CallerRuns` policy), and when the processor's own thread 
pool is shut down with `shutdownNow` and the task was still queued. Without 
this, the file would be kept until the spool directory is removed when the 
context stops.
   
   This builds on CAMEL-25012 (#26870), which made the processor count its 
parallel tasks as pending for the graceful shutdown. A task that never runs has 
to stop being counted as well as release its copy, so both are now done in one 
place. `submitTask(copy, task)` wraps the task in a `ParallelTask`, counts it 
and adds it to a set of the tasks that have not started yet. When the thread 
pool runs it, it removes itself from the set, runs the onCompletion and then 
stops being counted. When it will not run, `discard()` removes it from the set, 
stops counting it and releases the copy. Whichever removes the task from the 
set first wins, so a task is either run or discarded, never both. `discard()` 
is called when `submit` throws, when the thread pool is shut down after 
`submit` returns, and, after `shutdownNow` in `doShutdown`, for every task 
still in the set. The last one replaces subtracting the number of tasks that 
`shutdownNow` returns: the set also has a task that a thread took from th
 e queue but did not start before `shutdownNow`, which then does not run 
either. As a side effect this also fixes a gap in CAMEL-25012: a task that a 
thread pool that is shut down discarded without an exception stayed counted as 
pending.
   
   **The first point changes the order of every stream cache clean-up**, not 
only in routes with an onCompletion:
   - A spooled file is now deleted after all the other on completions of the 
exchange. Before, it was deleted among the other order `0` on completions, in 
reverse registration order. It is only ever deleted later than before, never 
earlier, so nothing that worked before can now find the file deleted. A custom 
`Synchronization` with an order above `0` can now also read a spooled body.
   - On completions with the same order `Ordered.LOWEST` keep their reverse 
registration order among themselves. These are the FTP and SMB consumers' 
disconnect, which is registered before the stream cache is created and 
therefore still runs after the clean-up, and the close of a JPA `EntityManager`.
   - The onCompletion now runs before all `Ordered.LOWEST` on completions. That 
was already the case for the FTP and SMB disconnect, because of the 
registration order. A JPA `EntityManager` close registered after the 
onCompletion (a JPA endpoint later in the route) used to run before the 
onCompletion and now runs after it.
   - The stream caches that Multicast and Split register on the parent's unit 
of work are released in the same way. So an onCompletion can now also read a 
spooled result of a Multicast or Split (I did not write a test for this).
   
   The upgrade guide for 4.23 describes the new order.
   
   Not covered: with `useOriginalMessage`, the onCompletion's copy gets the 
unit of work's original message (shared, not copied), so the copy does not take 
its own reference to that body. Without `parallelProcessing` this should now 
work because of the new order (not tested). With `parallelProcessing` the task 
can still read the original body after its file has been deleted. Fixing that 
would mean copying the original message, which I left out to keep this change 
small.
   
   Tests (new `OnCompletionStreamCachingSpoolTest`, spooling enabled with a 1 
KB threshold into a test directory, a 16 KB stream body):
   - `testOnCompletion`, `testOnCompletionParallel`, 
`testOnCompletionOnFailureOnly` (a failing route): the onCompletion reads the 
full body. The parallel onCompletion first waits on a latch that the test 
releases after `sendBody` returns, so the original exchange is done before the 
body is read.
   - `testOnCompletionBeforeConsumer` and `testNoOnCompletion`: controls, and 
the spool file is still deleted when the exchange is done.
   - `testParallelTaskRejected` (a shut down pool that throws 
`RejectedExecutionException`), `testParallelTaskDiscarded` (a shut down pool 
with a discard policy) and `testParallelTaskDroppedAtShutdown` (a single-thread 
pool whose queued task is dropped by `shutdownNow` when the graceful shutdown 
of `CamelContext` times out): no spool file is left behind, and no onCompletion 
task is still counted as pending.
   - Every test asserts that the spool directory is empty at the end 
(Awaitility).
   
   Without the change (with the `OnCompletionProcessor` and 
`FileInputStreamCache` of `main`), `testOnCompletion`, 
`testOnCompletionParallel` and `testOnCompletionOnFailureOnly` fail 
(`mock://done Received message count. Expected: <1> but was: <0>`), 
`testParallelTaskDiscarded` fails because the discarded task is still counted 
as pending (`expected: <0> but was: <1>`), and the other four pass. With the 
change, but with the release of the copy disabled, exactly the three 
`testParallelTask*` tests fail with spool files left behind. With the change, 
but with a discarded task still counted as pending, the same three tests fail 
on the pending count, and so does 
`OnCompletionParallelProcessingForcedShutdownTest`. So each of the paths is 
covered for both.
   
   Because the clean-up order changes for every exchange, I ran the full test 
suite of camel-core with this change on top of the current `main`, together 
with the core modules and the seda module (whose tests live in camel-core): 
camel-core 7791 tests, 0 failures, 45 skipped, and 6 errors that are not 
related to this change. Five are `BeanParameterValueWithCommaTest`, which needs 
the camel-bean change of CAMEL-25081 that my offline build did not include; it 
passes with camel-bean built. The sixth is `ValidatorExternalResourceTest`, 
which timed out downloading an XSD from GitHub and passes when run again. The 
targeted tests (`*OnCompletion*`, `*StreamCach*`, `*Shutdown*` and `*WireTap*`, 
225 tests, including the CAMEL-25012 tests) pass.
   
   # Target
   
   - [x] I checked that the commit is targeting the correct branch (Camel 4 
uses the `main` branch)
   
   # Tracking
   - [x] If this is a large change, bug fix, or code improvement, I checked 
there is a [JIRA issue](https://issues.apache.org/jira/browse/CAMEL) filed for 
the change (usually before you start working on it).
   
   # Apache Camel coding standards and style
   
   - [x] I checked that each commit in the pull request has a meaningful 
subject line and body.
   - [ ] I have run `mvn clean install -DskipTests` locally from root folder 
and I have committed all auto-generated changes.
     (I built and tested the affected modules, including the formatter and 
import-sort plugins. I did not run the full root build.)
   
   # AI-assisted contributions
   
   - [x] If this PR includes AI-generated code, commits have proper 
co-authorship attribution (e.g., `Co-authored-by` trailers) and the PR 
description identifies the AI tool used.
     This PR was prepared with Claude Code (Claude Opus 5.5). The commit 
carries a `Co-Authored-By` trailer.
   
   _Claude Code on behalf of allthingssecurity_
   
   🤖 Generated with [Claude Code](https://claude.com/claude-code)
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to