shashank created CAMEL-25211:
--------------------------------

             Summary: 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


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