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)