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)

Reply via email to