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]
