shashank created CAMEL-25210:
--------------------------------

             Summary: camel-cassandraql - CassandraAggregationRepository hands 
aggregations that are still open to the recover task, which sends and then 
deletes them, and a completed exchange that failed is never recovered
                 Key: CAMEL-25210
                 URL: https://issues.apache.org/jira/browse/CAMEL-25210
             Project: Camel
          Issue Type: Bug
          Components: camel-cassandraql
            Reporter: shashank


{{CassandraAggregationRepository}} (and 
{{NamedCassandraAggregationRepository}}) implement 
{{RecoverableAggregationRepository}}, and recovery is enabled by default 
({{useRecovery=true}}, {{recoveryInterval=5000}}). The table only holds the 
aggregations in progress, one row per correlation key, so the recovery methods 
work on the wrong rows. This is the defect fixed for camel-infinispan in 
CAMEL-24622 and for camel-caffeine and camel-ehcache in CAMEL-25153.

The Aggregate EIP uses two key spaces: on completion it calls {{remove(ctx, 
key, exchange)}}, where a recoverable repository keeps the completed exchange 
under its *exchange id*; {{confirm(ctx, exchangeId)}} deletes it once the 
exchange was processed; the recover task calls {{scan(ctx)}} for the exchange 
ids of completed but unconfirmed exchanges and loads each with {{recover(ctx, 
exchangeId)}}.

In {{CassandraAggregationRepository}} (line numbers of main):
* {{remove}} (:279) deletes the row of the correlation key. The completed 
exchange is gone, nothing is left to recover.
* {{scan}} (:326) returns the {{EXCHANGE_ID}} column of every row, that is the 
exchange ids of the aggregations still in progress.
* {{recover}} (:339) loads the row with that exchange id, an aggregation in 
progress.
* {{confirm}} (:251) deletes the row whose {{EXCHANGE_ID}} is the confirmed id 
({{DELETE ... IF EXCHANGE_ID=?}}), which is again an aggregation in progress.

So:
* an aggregation that stays open longer than the first run of the recover task 
(about one second after start, then every {{recoveryInterval}}) is sent to the 
route after the aggregator, incomplete, with {{CamelRedelivered=true}}; when 
that exchange is processed, {{confirm}} deletes the open aggregation, so the 
messages that arrive later start a new group and the group is split;
* a completed group whose processing fails after the aggregator is not 
recovered: the repository deleted it in {{remove}}.

h3. Reproduction

The repository runs against an in-memory table behind a mocked {{CqlSession}} 
(Docker is not available here to run the module's Cassandra ITs), with the 
route {{from("direct:x").aggregate(header("id"), 
strategy).aggregationRepository(repo).completionSize(3).to("mock:x").process(failOnce)}}
 and {{recoveryInterval=100}} (only to make the test fast). The strategy 
appends the body to the old exchange and returns it. The runs of the recover 
task are counted through a subclass of the repository (no sleeps).
* two messages "a" and "b" of a group, then three runs of the recover task: the 
mock receives the open group "a+b" twice (the first delivery fails, the second 
is processed and confirmed, which deletes the group). Expected: nothing. 3 of 3 
runs.
* three messages, the step after the aggregator fails the first time: the mock 
receives "a+b+c" once and never again. Expected: "a+b+c" again with 
{{CamelRedelivered=true}}. 3 of 3 runs.

h3. Proposed fix

As for Infinispan, Caffeine and Ehcache, keep the completed exchange in the 
same table, under the aggregation key {{"camel-recovery:" + exchangeId}}, until 
it is confirmed. No change of the table is needed.
* {{remove}}: when {{useRecovery}} is enabled, insert the exchange passed to 
{{remove}} under the recovery key (the stored row does not contain the message 
that completed the group, see CAMEL-24946), then delete the row of the 
correlation key. Writing the completed exchange first means a failure in 
between leaves it to recovery rather than losing it.
* {{confirm}}: delete the recovery key (when {{useRecovery}} is enabled). The 
{{DELETE ... IF}} statement is no longer needed.
* {{scan}}: the exchange ids of the recovery keys of the repository's rows 
(none when recovery is disabled); {{getKeys}}: only the other keys.
* {{recover}}: read the recovery key.

With the fix the harness above receives nothing for the open group, and the 
failed completed group is recovered once with all three messages and confirmed 
(3 of 3). The test ({{CassandraAggregationRepositoryRecoveryTest}}, a unit test 
with the in-memory session) is added to the module; the camel-cassandraql unit 
tests pass (8 tests). {{CassandraAggregationRepositoryIT}} and 
{{NamedCassandraAggregationRepositoryIT}} encode the old key space 
({{testConfirmExist}}, {{testScan}}, {{testRecover}} confirm, scan and recover 
aggregations in progress) and are reworked; they compile but were not run here 
(they need Cassandra in Docker). The upgrade guide gets a note, as {{scan()}} 
no longer returns open groups and the table now also holds completed exchanges 
until they are confirmed ({{ttl}} applies to them as well).

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 "CassandraAggregationRepository" 
(CAMEL-20306, CAMEL-23372, CAMEL-23609, all deserialization filters), component 
camel-cassandraql with "aggregation" (the same), "cassandra" with "aggregation" 
and "recover" (none); GitHub pull requests "CassandraAggregationRepository", 
"cassandra aggregation": none for this defect. No open pull request touches the 
file.

_Filed with Claude Code on behalf of allthingssecurity._




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

Reply via email to