shashank created CAMEL-25164:
--------------------------------
Summary: camel-infinispan - the aggregation repositories recover a
completed exchange without the message that completed the group
Key: CAMEL-25164
URL: https://issues.apache.org/jira/browse/CAMEL-25164
Project: Camel
Issue Type: Bug
Components: camel-infinispan
Reporter: shashank
Since CAMEL-24622, {{InfinispanAggregationRepository}} (the shared base of
{{InfinispanEmbeddedAggregationRepository}} and
{{InfinispanRemoteAggregationRepository}}) keeps a completed exchange for
recovery under a {{"camel-recovery:" + exchangeId}} key in the same cache.
{{remove(ctx, key, exchange)}} stores the entry it removed from the cache
there, and only uses the exchange passed to {{remove}} when the cache had no
entry:
{code:java}
DefaultExchangeHolder holder = getCache().remove(key);
if (useRecovery) {
if (holder == null) {
holder = DefaultExchangeHolder.marshal(exchange, true,
allowSerializedHeaders);
}
getCache().put(recoveryKey(exchange.getExchangeId()), holder);
}
{code}
When an incoming message completes a group, the Aggregate EIP aggregates it
into the group and calls {{remove}} without adding the final state to the
repository first ({{AggregateProcessor}} only calls {{add}} for a group that is
not complete). The entry in the cache is therefore the group before the last
message. If the processing after the aggregator fails, the recover task sends
the stored entry, and the recovered exchange is missing the message that
completed the group.
This is the defect CAMEL-24946 fixed for {{KeyValueAggregationRepository}}, and
the Caffeine and Ehcache repositories store the given exchange as well
(CAMEL-25153). Claus Ibsen pointed out the Infinispan case in the review of
apache/camel#27108.
h3. Reproduction
Route {{from("direct:start").aggregate(header("id"),
strategy).aggregationRepository(repo).completionSize(3).to("mock:aggregated").process(failOnce)}}
with an {{InfinispanEmbeddedAggregationRepository}} and
{{recoveryInterval=100}} (only to make the test fast). The strategy appends the
body to the old exchange and returns it. After sending "a", "b" and "c" for one
group, the mock receives "a+b+c", then the recovered exchange "a+b" with
{{CamelRedelivered=true}}. The message "c" is lost.
h3. Proposed fix
In {{remove}}, store the marshalled exchange passed to {{remove}} under the
recovery key, as {{KeyValueAggregationRepository}} and the Caffeine and Ehcache
repositories do. The embedded and the remote repository share this code, so
both are fixed. With the fix the route above receives "a+b+c" twice (the second
time as redelivered).
Tests: a route test in camel-infinispan-embedded (the one above) and an
operations test that the exchange given to {{remove}} is recovered. Both fail
without the fix ({{expected: <a+b+c> but was: <a+b>}}). The remote repository
tests are integration tests that need a container.
Affected: the versions with the CAMEL-24622 recovery store. The same code is on
main and on the camel-4.18.x and camel-4.22.x branches.
_Filed with Claude Code on behalf of allthingssecurity._
--
This message was sent by Atlassian Jira
(v8.20.10#820010)