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)