[
https://issues.apache.org/jira/browse/CAMEL-25124?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-25124.
---------------------------------
Resolution: Fixed
Fixed by https://github.com/apache/camel/pull/27033 (merged to main, 4.23.0).
_Claude Code on behalf of davsclaus_
> 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
> Priority: Minor
> Fix For: 4.23.0
>
>
> 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)