allthingssecurity opened a new pull request, #27033: URL: https://github.com/apache/camel/pull/27033
# Description [CAMEL-25124](https://issues.apache.org/jira/browse/CAMEL-25124) With stream caching spooled to disk, the Aggregate EIP keeps bodies whose spool files have already been deleted, so the aggregated exchange cannot read them: ``` from("direct:start") .aggregate(constant(true), new GroupedBodyAggregationStrategy()).completionSize(3) .to("direct:process"); // reading the bodies fails WARN AggregateProcessor - Error processing aggregated exchange ... NoSuchFileException: cos123.tmp ``` `AggregateProcessor.doProcess` stores `ExchangeHelper.createCorrelatedCopy(exchange, false)`. The copy shares the `FileInputStreamCache` of the incoming exchange, but it takes no reference to the file. The incoming exchange completes as soon as it has been aggregated, and `FileInputStreamCache.TempFileManager` deletes the file. This hits every body that the group keeps as a stream (`GroupedBodyAggregationStrategy`, `UseLatestAggregationStrategy`, `flexible()` into a collection, ...): all grouped bodies, the body of a timeout completion, and even the body of the exchange that completes the group, because the aggregated exchange is sent on the aggregator's own thread and races with the completion of that exchange. This change: - `doProcess` gives the copy its own reference when its body is a stream cache that is not in memory: it removes `CamelStreamCacheUnitOfWork` from the copy and replaces the body with `sc.copy(copy)`, as the Wire Tap EIP does (CAMEL-12108). The reference is released by an on completion of the copy. - When the copy is aggregated, its on completions, and those of the group it is merged into, are handed over to the exchange that now holds the group (after the repository add, which can fail with optimistic locking), and when the group completes, to the aggregated exchange (in `doOnCompletion`, after the repository remove). So the files are deleted when the aggregated exchange is done. - The reference is released on every path where a copy or a group is dropped: an aggregation failure (with or without `discardOnAggregationFailure`, for the first exchange of a group and for a group), `discardOnCompletionTimeout` and `forceDiscardingOfGroup` (in `discard`), an optimistic locking retry (each attempt makes a new copy), and a correlation key found closed under the lock. - The group keeps the references only with `MemoryAggregationRepository` (the default, which keeps the exchange instance) and without optimistic locking: - A repository that stores a copy of the exchange (the persistent ones, and `KeyValueAggregationRepository`) has read the body in `add`, so the reference is released right after the add. A subclass of the memory repository that stores a copy is detected after the add (`get(key) != answer`). - With optimistic locking, other threads read and replace the exchanges of the memory repository without a lock, and handing references over into an exchange that is already in the repository would race with them. So the references are released after the add there too: the stored bodies behave as before, and only the body of the exchange that completes the group is kept. I preferred this to adding locking around the repository calls (a first attempt with a lock made `AggregatePreCompleteLostGroupTest` hang, as its repository blocks inside `remove`). - `handoverCompletions`/`releaseCompletions` use the exchange's own on completion list (`handoverCompletions()`), never a unit of work: the copies have none. One consequence to call out: the spool files of a group now live until the aggregated exchange is done. For example with `UseLatestAggregationStrategy` and a `completionTimeout`, the file of every message of the group is kept until the group completes, as the aggregator cannot know which bodies a strategy keeps. Before, they were deleted early, which is the bug. Tests: - New `AggregateStreamCachingSpoolTest` (camel-core). Spooling is enabled with a 1 KB threshold into a test directory, and the bodies are 16 KB streams. The route of the aggregated exchange first waits on a latch that the test releases after the `sendBody` calls returned, so the incoming exchanges are done before the bodies are read, and the tests are deterministic. Every test asserts that the spool directory is empty at the end (Awaitility). - Output tests, which read the full bodies: grouped bodies with `completionSize(3)`, `UseLatest` with `completionTimeout`, `UseLatest` with `completionSize(1)`, optimistic locking with a repository whose first add fails, a repository that stores copies (not a `MemoryAggregationRepository`), and a `MemoryAggregationRepository` subclass that stores copies. - Drop paths: `discardOnAggregationFailure` (the first exchange of a group, and a stored group), `discardOnCompletionTimeout`, `forceDiscardingOfGroup`, and a retry that finds the correlation key closed (the repository's first add completes and closes the group, then fails). - Without the change, the grouped, timeout, `completionSize(1)` and optimistic tests fail (`mock://result Received message count. Expected: <1> but was: <0>`, with the `NoSuchFileException` above in the log). The drop path tests pass on main, as no reference was taken there; they fail with a version of this fix that hands the references over without releasing them on the drop paths (6 of them leave spool files behind). - With the change: camel-core `*Aggregat*,*StreamCach*`: 334 tests, 0 failures, 0 errors, 5 skipped. I also ran the whole camel-core suite once (7954 tests): the only failures are in 20 test classes (bean parameters, file cluster, properties, rest, saga, ...) that fail the same way without this change in my local build, because it uses older snapshots of other modules; none of them uses the aggregator or stream caching. This merges cleanly with main and with the related open #27002 (CAMEL-25100, OnCompletion EIP). Found while checking the stream cache reference counting across asynchronous hand-offs (a TLA+ model of the counter for the seda case, CAMEL-25094). I reproduced the aggregator failure with the real classes, and checked every drop path with a harness that counts the spool files left behind 1.5 s after everything is done: main leaves none (and fails to read the bodies), this change leaves none and reads the full bodies. # 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]
