[
https://issues.apache.org/jira/browse/CAMEL-25210?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-25210.
---------------------------------
Resolution: Fixed
> 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
> Assignee: shashank
> Priority: Major
> Fix For: 4.23.0
>
>
> {{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)