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)