[
https://issues.apache.org/jira/browse/CAMEL-25211?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen reassigned CAMEL-25211:
-----------------------------------
Assignee: shashank
> 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
>
> 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)