shashank created CAMEL-25124:
--------------------------------

             Summary: camel-core - Aggregate EIP keeps bodies that stream 
caching spooled to disk, but their spool files are deleted when the incoming 
exchanges complete, so the aggregated message cannot be read 
(NoSuchFileException)
                 Key: CAMEL-25124
                 URL: https://issues.apache.org/jira/browse/CAMEL-25124
             Project: Camel
          Issue Type: Bug
          Components: camel-core
            Reporter: shashank


With stream caching spooled to disk ({{spoolEnabled=true}}, a body over the 
spool threshold), the body is a {{FileInputStreamCache}}, and its temporary 
file is deleted when the exchange that created it is done 
({{FileInputStreamCache.TempFileManager}}).

{{AggregateProcessor.doProcess}} stores a correlated copy of the incoming 
exchange ({{ExchangeHelper.createCorrelatedCopy(exchange, false)}}, no 
hand-over of on completions). The copy shares the {{FileInputStreamCache}} of 
the incoming exchange, but takes no reference to its file. The incoming 
exchange completes as soon as the aggregator has added it to the group, and 
deletes the file. Every body that the group keeps as a stream (one that the 
aggregation strategy does not read while the incoming exchange is still 
running) then points to a deleted file.

This is not limited to the bodies of earlier exchanges: the aggregated exchange 
is always sent on the aggregator's own thread (a single thread executor when 
{{parallelProcessing}} is off), so even the body of the exchange that completes 
the group races with the completion of that exchange.

Results (a streamed body over the threshold, the route of the aggregated 
exchange reads the bodies after the incoming exchanges are done, 3 runs each):
* {{GroupedBodyAggregationStrategy}} with {{completionSize(3)}}: all three 
bodies fail with {{NoSuchFileException}}. Without that wait, all bodies failed 
in 9 of 13 runs, and in the others all but the last one.
* {{UseLatestAggregationStrategy}} with {{completionTimeout}}: the body fails 
with {{NoSuchFileException}}.
* {{UseLatestAggregationStrategy}} with {{completionSize(1)}}: the body fails 
with {{NoSuchFileException}}.
* The same with optimistic locking and a retry of the repository add.

The aggregator's exception handler logs the failure at WARN ("Error processing 
aggregated exchange"), but the incoming exchanges have completed (a consumer 
has committed their messages) and the aggregated output cannot be processed, so 
the data is lost. It affects the strategies that keep the bodies as they are: 
{{GroupedBodyAggregationStrategy}}, {{GroupedExchangeAggregationStrategy}} and 
{{GroupedMessageAggregationStrategy}} (by code reading), {{flexible()}} into a 
collection, {{UseLatestAggregationStrategy}} and 
{{UseOriginalAggregationStrategy}}. A strategy that reads the body in 
{{aggregate()}} (such as the string or zip strategies) is not affected.

h3. Proposed fix

* In {{doProcess}}, if the body of the copy is a stream cache that is not in 
memory, the copy takes its own reference: it removes 
{{CamelStreamCacheUnitOfWork}} from the copy (as the Wire Tap EIP does) and 
replaces the body with {{sc.copy(copy)}}. 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 holds the group, and when the 
group completes, to the aggregated exchange. So the files are deleted when the 
aggregated exchange is done.
* The reference is released on every path where the copy or the group is 
dropped: an aggregation failure (with and without 
{{discardOnAggregationFailure}}, for the first exchange of a group and for a 
group), {{discardOnCompletionTimeout}}, {{forceDiscardingOfGroup}}, 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 the {{MemoryAggregationRepository}} 
(the default, which keeps the exchange instance) and without optimistic 
locking. A repository that stores a copy of the exchange (every persistent 
repository, and {{KeyValueAggregationRepository}}) has read the body when the 
exchange is added, so the reference is released right after the add. With 
optimistic locking, other threads use the exchanges of the memory repository 
without a lock, so the references are also released after the add there: only 
the body of the exchange that completes the group is kept, and the stored 
bodies behave as before. A subclass of the memory repository that stores a copy 
is detected after the add.

A consequence of the fix: 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), as the 
aggregator cannot know which bodies a strategy keeps.

With the fix, the cases above read the full bodies, and the spool directory is 
empty afterwards on every drop path (checked with a harness per path, and with 
unit tests; before the fix no file was left behind either, as no reference was 
taken).

Affected: all 4.x versions.

Duplicate check (2026-09-29): JIRA "stream caching" with "aggregator" or 
"aggregate" (CAMEL-24844, CAMEL-11497: unrelated), "spool" with "aggregat*" 
(CAMEL-7787, 2014, the multicast unit of work), "FileInputStreamCache" (17 
issues, none about the aggregator). CAMEL-24991 (Aggregate EIP deep review) did 
not touch stream caching. GitHub pull requests "StreamCache", "spool": nothing 
about the aggregator.

_Filed with Claude Code on behalf of allthingssecurity._




--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to