shashank created CAMEL-25153:
--------------------------------

             Summary: 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


{{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