[ 
https://issues.apache.org/jira/browse/CAMEL-25211?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Claus Ibsen updated CAMEL-25211:
--------------------------------
    Fix Version/s: 4.23.0

> camel-hazelcast - the aggregation repositories recover a completed exchange 
> without the message that completed the group, and 
> ReplicatedHazelcastAggregationRepository fails to complete any group with 
> recovery enabled
> ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25211
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25211
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-hazelcast
>            Reporter: shashank
>            Assignee: shashank
>            Priority: Major
>             Fix For: 4.23.0
>
>
> Both defects are in the non-optimistic {{remove(ctx, key, exchange)}} with 
> recovery enabled, which is the default configuration of both repositories 
> ({{optimistic=false}}, {{useRecovery=true}}).
> h3. 1. HazelcastAggregationRepository stores the previous state of the group 
> for recovery
> {{remove}} (line numbers of main, :390-396) removes the group in a Hazelcast 
> transaction and puts the entry it removed into the completed map:
> {code:java}
> DefaultExchangeHolder removedHolder = tCache.remove(key);
> tPersistentCache.put(exchange.getExchangeId(), removedHolder);
> {code}
> When an incoming message completes a group ({{completionSize}}, 
> {{completionPredicate}}), 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 map is therefore the group before the last message. If the 
> processing after the aggregator fails, the recover task sends that entry, and 
> the recovered exchange is missing the message that completed the group.
> This is the defect fixed for {{KeyValueAggregationRepository}} in CAMEL-24946 
> and for Infinispan in CAMEL-25164. The optimistic branch of the same method 
> already stores the given exchange.
> h3. 2. ReplicatedHazelcastAggregationRepository removes from the wrong map
> {{ReplicatedHazelcastAggregationRepository}} keeps its groups and completed 
> exchanges in {{ReplicatedMap}}s. Its {{remove}} (:292-298) runs the same 
> transaction as the IMap repository, on {{tCtx.getMap(mapName)}} and 
> {{tCtx.getMap(persistenceMapName)}}, which are *IMaps* with the same names, 
> not the replicated maps. The group is never in that IMap, so 
> {{removedHolder}} is {{null}}, the {{put}} fails with {{NullPointerException: 
> value can't be null}}, the transaction is rolled back and {{remove}} throws 
> {{RuntimeCamelException("Transaction ... was rolled back for remove operation 
> ...")}}. Every group of two or more messages fails to complete and stays in 
> the replicated map. (A group of one message is not removed, so it is not 
> affected.) The code is the same since the repository was added (Camel 3.4).
> h3. Reproduction
> Route {{from("direct:x").aggregate(header("id"), 
> strategy).aggregationRepository(repo).completionSize(3).to("mock:x").process(failOnce)}}
>  with an embedded two-member Hazelcast cluster (the module's test support), 
> {{recoveryInterval=100}} (only to make the test fast). The strategy appends 
> the body to the old exchange and returns it. Messages "a", "b", "c" for one 
> group:
> * {{HazelcastAggregationRepository}}: the mock receives "a+b+c", then the 
> recovered exchange "a+b" ({{CamelRedelivered=true}}); the message "c" is 
> lost. 3 of 3 runs.
> * {{ReplicatedHazelcastAggregationRepository}}: sending "c" fails with the 
> rolled back transaction caused by "value can't be null"; the group is not 
> sent. 3 of 3 runs.
> h3. Proposed fix
> * {{HazelcastAggregationRepository.remove}}: put the holder marshalled from 
> the given exchange (already computed at the top of the method) into the 
> completed map, inside the same transaction.
> * {{ReplicatedHazelcastAggregationRepository.remove}}: a {{ReplicatedMap}} 
> cannot take part in a Hazelcast transaction, so under the per-key lock that 
> {{add}} already uses ({{lockMap}}), put the given exchange into the 
> replicated completed map, then remove the group from the replicated map. 
> Writing the completed exchange first means a failure in between leaves it to 
> recovery rather than losing it.
> A route test ({{HazelcastAggregationRepositoryRecoverTest}}) runs the 
> reproduction for both repositories: without the fix it fails as above, with 
> the fix both recover "a+b+c" once and leave no group in the repository. With 
> the fix the camel-hazelcast tests pass (233 tests). 
> {{HazelcastAggregationRepositoryRoutesTest.checkAggregationFromTwoRoutes}} is 
> flaky on main as well (it receives a second exchange in the first run and 
> passes on rerun; {{HazelcastAggregationRepositoryRecoverableRoutesTest}} uses 
> the same repository name), with and without this change.
> Affected: all versions (the same code at camel-3.20.0, 4.0.0, 4.10.0, 4.14.0, 
> 4.18.0, 4.22.0 and main).
> Duplicate check (2026-09-30): JIRA text "HazelcastAggregationRepository" and 
> "ReplicatedHazelcastAggregationRepository" (CAMEL-24413 and CAMEL-23414 
> serialization, CAMEL-9017 confirm without recovery, CAMEL-8971 redelivery 
> with a strategy that returns the new exchange (Won't Fix), CAMEL-8438 
> optimistic locking); GitHub pull requests "HazelcastAggregationRepository", 
> "ReplicatedHazelcastAggregationRepository", "hazelcast aggregation recovery": 
> none for these defects. No open pull request touches these files.
> _Filed with Claude Code on behalf of allthingssecurity._



--
This message was sent by Atlassian Jira
(v8.20.10#820010)

Reply via email to