[
https://issues.apache.org/jira/browse/CAMEL-25153?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Claus Ibsen resolved CAMEL-25153.
---------------------------------
Resolution: Fixed
Fixed via https://github.com/apache/camel/pull/27108 (merged to main for
4.23.0).
_Claude Code on behalf of davsclaus_
> 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
> 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)