allthingssecurity commented on code in PR #27033:
URL: https://github.com/apache/camel/pull/27033#discussion_r4131020402
##########
core/camel-core-processor/src/main/java/org/apache/camel/processor/aggregate/AggregateProcessor.java:
##########
@@ -711,6 +747,53 @@ protected void doAggregationComplete(
}
}
+ /**
+ * Adds the answer to the repository. If the repository keeps the answer,
the answer takes over what the aggregated
+ * exchanges hold (such as a reference to a spooled stream cache), so it
is released when the group is done with.
+ */
+ private void doAggregationRepositoryAddAndHandover(
+ String key, Exchange originalExchange, Exchange newExchange,
Exchange answer) {
+ // the answer is not yet visible to other threads
+ handoverCompletions(newExchange, answer);
+ boolean keeping = isKeepingReferences();
+ // otherwise release after the add, as a persistent repository reads
the body when it stores the exchange
+ List<Synchronization> release = keeping ? null :
answer.getExchangeExtension().handoverCompletions();
Review Comment:
Agreed, fixed in 77279e8780bd using the design you suggested.
Right after `sc.copy(copy)`, the aggregator takes
`copy.getExchangeExtension().handoverCompletions()` (the copy has no
completions of its own) and keeps those syncs in an aggregator-owned
`SpooledStreamCaches` holder. From then on it releases or hands over only that
holder. It never moves or runs the exchanges' own on-completions, so
`DeleteZipFileOnCompletion` / `DeleteTarFileOnCompletion` behave as on main.
- Non-keeping mode (optimistic locking, or any repository other than the
memory one): the holder of the new exchange is released after the add. If the
exchange completes the group, its holder is added to the aggregated exchange as
a single on-completion.
- Keeping mode (`MemoryAggregationRepository` without optimistic locking):
the group's holders live in a map by correlation key, guarded by the
aggregation lock. `doOnCompletion` takes the group's holder and adds it to the
aggregated exchange. The discard paths (`discard()`, which also covers
`forceDiscarding`) release only `SpooledStreamCaches` instances.
New `ZipAggregationStrategyRepositoryTest` and
`TarAggregationStrategyRepositoryTest` cover `optimisticLocking()` and a
repository that stores copies. Both fail on the previous head, where the
archive is deleted after the first add, and pass now.
_Claude Code on behalf of allthingssecurity_
--
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]