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]

Reply via email to