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)