[ 
https://issues.apache.org/jira/browse/CAMEL-25153?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Claus Ibsen reassigned CAMEL-25153:
-----------------------------------

    Assignee: shashank

> camel-caffeine, camel-ehcache - the aggregation repositories hand 
> aggregations that are still open to the recover task, so incomplete groups 
> are sent on as redelivered exchanges, and a completed exchange that failed is 
> never recovered
> ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------
>
>                 Key: CAMEL-25153
>                 URL: https://issues.apache.org/jira/browse/CAMEL-25153
>             Project: Camel
>          Issue Type: Bug
>          Components: camel-caffeine, camel-ehcache
>            Reporter: shashank
>            Assignee: shashank
>            Priority: Minor
>             Fix For: 4.23.0
>
>
> {{CaffeineAggregationRepository}} and {{EhcacheAggregationRepository}} 
> implement {{RecoverableAggregationRepository}}, and recovery is enabled by 
> default ({{useRecovery=true}}, {{recoveryInterval=5000}}). Both keep a single 
> cache keyed by the correlation key and have no store for completed exchanges, 
> so the recovery methods work on the wrong key space. This is the defect fixed 
> for camel-infinispan in CAMEL-24622.
> 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 when 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 both repositories (line numbers of main, Caffeine / Ehcache):
> * {{remove}} (:149 / :173) deletes the entry of the correlation key. The 
> completed exchange is gone, nothing is left to recover.
> * {{confirm}} (:155 / :179) deletes the entry {{exchangeId}}, which is never 
> a key of the cache: a no-op.
> * {{scan}} (:168 / :193) returns {{getKeys()}}, the correlation keys of the 
> groups that are still open.
> * {{recover}} (:176 / :201) is called with such a key and returns the open 
> group.
> {{AggregateProcessor.RecoverTask}} compares the scanned ids with the exchange 
> ids in progress, which never match a correlation key. So:
> * every group that stays open longer than about one second (the first run of 
> the recover task, then every {{recoveryInterval}}) is sent to the route after 
> the aggregator, incomplete, with {{CamelRedelivered=true}} and 
> {{CamelRedeliveryCounter}}; it is sent again at every run until 
> {{maximumRedeliveries}} (3) is reached, then the recover task tries to move 
> it to the dead letter channel ({{deadLetterUri}}, none by default). The group 
> itself stays open and is sent once more when it really completes, so its 
> messages are delivered several times;
> * a completed group whose processing fails after the aggregator is not 
> recovered: the repository deleted it in {{remove}}.
> Both repositories are the in-memory (or Ehcache configured) choices for the 
> Aggregate EIP; the open-group case hits any aggregation with a completion 
> size, predicate or timeout that is not reached within a second, which is the 
> normal use.
> h3. Reproduction
> Route {{from("direct:in").aggregate(header("id"), 
> strategy).aggregationRepository(repo).completionSize(3).process(record)}} 
> with each repository and {{recoveryInterval=100}} (only to make the test 
> fast; the default is 5000 ms). The recover task is observed through a 
> subclass that counts the calls of {{scan}} and {{recover}} (no sleeps):
> * two of the three messages of a group: after two runs of the recover task 
> the route has received the open group "a+b" three times, each with 
> {{CamelRedelivered=true}}; with {{useRecovery=false}} it receives nothing 
> (control). 3 of 3 runs, and 20 of 20 in a loop, for both repositories;
> * {{completionSize(2)}}, the step after the aggregator fails the first time: 
> the completed group is received once and never again after four runs of the 
> recover task; the repository is empty. 3 of 3 runs for both.
> A TLA+ model of the group, the repository, the completion, the downstream 
> processing and the recover task shows the same: the current key space 
> delivers a partial group (a violation of "only complete groups leave the 
> aggregator") and never re-delivers a failed completed group; with the fix 
> below both properties hold, and the control without the recover task holds on 
> the current code.
> h3. Proposed fix
> The same as CAMEL-24622 for Infinispan: keep completed exchanges under a 
> prefixed exchange-id key in the same cache.
> * {{remove}}: delete the correlation key, and when {{useRecovery}} is enabled 
> put the marshalled exchange passed to {{remove}} under {{"camel-recovery:" + 
> exchangeId}}. (Infinispan keeps the entry it removed instead, which does not 
> contain the message that completed the group; storing the given exchange is 
> what CAMEL-24946 did for {{KeyValueAggregationRepository}} and what 
> {{JdbcAggregationRepository}} does.)
> * {{confirm}}: delete {{"camel-recovery:" + exchangeId}} (when 
> {{useRecovery}} is enabled).
> * {{scan}}: the exchange ids of the prefixed keys (empty when recovery is 
> disabled); {{getKeys}}: only the correlation keys.
> * {{recover}}: read the prefixed key.
> With this the harness above receives nothing for the open group, and the 
> failed completed group is re-delivered once with {{CamelRedelivered=true}}, 
> then confirmed (3 of 3 for both repositories). No configuration change is 
> needed (one cache, as for Infinispan).
> The operation tests of both repositories ({{testConfirmExist}}, {{testScan}}, 
> {{testRecover}}) encode the correlation-key behaviour (they confirm, scan and 
> recover by correlation key) and are reworked with the fix, as for Infinispan. 
> A route test per module (an open group and a completed group whose downstream 
> step fails once) receives the open group first without the fix, and with the 
> fix only the completed group, then once more as redelivered, then nothing; 
> with the fix camel-caffeine passes 84 tests and camel-ehcache 66 tests. The 
> upgrade guide should mention that open groups are no longer returned by 
> {{scan()}} and that completed exchanges are now kept until confirmed.
> Affected: all versions that have these repositories (the same code at 
> camel-3.20.0, 4.0.0, 4.10.0, 4.14.0, 4.18.0 and 4.22.0).
> Duplicate check (2026-09-30): JIRA text "CaffeineAggregationRepository" 
> (CAMEL-23411, deserialization filter), "EhcacheAggregationRepository" 
> (CAMEL-19096, OversizeMappingException), "aggregation repository" with 
> "recover" since 2023 (CAMEL-24622 Infinispan, CAMEL-24946 
> KeyValueAggregationRepository, CAMEL-24943, CAMEL-24141, CAMEL-24991); GitHub 
> pull requests "CaffeineAggregationRepository", 
> "EhcacheAggregationRepository", "caffeine aggregation", "ehcache 
> aggregation", "aggregation repository recover": none for these two 
> repositories. No open pull request touches them.
> _Filed with Claude Code on behalf of allthingssecurity._



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

Reply via email to