shashank created CAMEL-25130:
--------------------------------

             Summary: camel-sql - JdbcAggregationRepository with optimistic 
locking: add() ignores the old exchange, so a group completed meanwhile is 
stored again (duplicates) and a new group can be overwritten (lost messages)
                 Key: CAMEL-25130
                 URL: https://issues.apache.org/jira/browse/CAMEL-25130
             Project: Camel
          Issue Type: Bug
          Components: camel-sql
            Reporter: shashank


The Aggregate EIP with {{optimisticLocking()}} and 
{{JdbcAggregationRepository}} is the documented set-up for several Camel 
instances that share the aggregation table ("Optimistic locking" in the 
camel-sql documentation). {{AggregateProcessor.doAggregation}} reads the group 
with {{repository.get}} (which puts the row version in the 
{{CamelOptimisticLockVersion}} property), aggregates, and then calls 
{{OptimisticLockingAggregationRepository.add(camelContext, key, oldExchange, 
newExchange)}}. The contract of that method is a compare-and-set against 
{{oldExchange}}: when {{oldExchange}} is null there must be no stored exchange, 
otherwise the stored exchange must still be the one that was read. With 
optimistic locking the aggregator does not lock the key, so another thread or 
node can complete the group between {{get}} and {{add}}.

{{JdbcAggregationRepository.add(camelContext, key, oldExchange, newExchange)}} 
ignores {{oldExchange}} and calls {{add(camelContext, key, newExchange)}}, 
which does in one transaction:
* {{SELECT COUNT(1) ... WHERE id = ?}};
* if a row exists: {{UPDATE ... WHERE id = ? AND version = ?}} with the version 
of the new exchange (0 rows: {{OptimisticLockingException}});
* if no row exists: {{INSERT}} with version 1.

Two cases are not detected:

*1. The group was completed while this thread aggregated on it.* Thread A reads 
group [x] and aggregates its message a. Meanwhile the group completes on 
another thread or node: a message that satisfies the completion predicate or 
size, or the completion timeout checker. {{remove()}} deletes the row (the 
version matches) and [x, b] is sent. Thread A then calls {{add}}: the row is 
gone, so it inserts [x, a] as a new group. The next completion sends [x, a, 
...]: x is delivered twice. Every message of the completed group is sent again 
with the next group.

*2. Every group starts at version 1 (ABA).* As above, but a message c starts a 
new group for the key (INSERT, version 1) before thread A calls {{add}}. Thread 
A's {{UPDATE ... WHERE version = 1}} matches the new group's row and replaces 
[c] with [x, a]: c is lost without an error, and x is delivered twice. 
{{remove()}} has the same problem: a completion that read version 1 of the old 
group can delete the new group.

The window is the time between {{get}} and {{add}} (the aggregation strategy 
plus two queries). It is hit when messages for a key keep arriving while its 
group completes, which is normal under load.

h3. Reproduction

Embedded H2, a {{DataSourceTransactionManager}}, the real 
{{JdbcAggregationRepository}} (tables as in the documentation), and a route
{{from("direct:in").aggregate(header("id"), 
concat).aggregationRepository(repo).optimisticLocking().eagerCheckCompletion().completionPredicate(header("last").isEqualTo("true"))}}.
 The aggregation strategy returns the old exchange (the recommended style) and 
holds message "a" on a latch after {{get}} has read the group; no sleeps are 
needed for the ordering.
* completed by another message: send x; thread A sends a (held); send b with 
last=true (completes [x,b]); release A; send z with last=true. Sent groups: 
[x,b] and [x,a,z]. x is delivered twice, 3 of 3 runs.
* completed by the timeout: {{completionTimeout(300)}}, send x, thread A sends 
a (held), the timeout sends [x], release A. Sent groups: [x] and [x,a], 3 of 3 
runs.
* ABA: as the first case, plus send c (a new group) before releasing A. Sent 
groups: [x,b] and [x,a,z]. c is lost and x is delivered twice, 3 of 3 runs.
* Control: the same messages without overlap: [x,a,b], [z].

A TLA+ model of the repository (each repository call one atomic transaction, 
which is stronger than READ COMMITTED) with aggregating threads, a completing 
message and the timeout checker checks that no message is sent twice and none 
is lost. It finds case 1 in 4 steps (Get(a), Get(b), Complete(b), Add(a)) and 
the timeout variant, and case 2 when only the check for a removed row is added. 
With both parts of the fix below the properties hold for three threads plus the 
timeout checker, also starting from an empty table.

h3. Proposed fix

Make {{add(camelContext, key, oldExchange, newExchange)}} a compare-and-set 
against {{oldExchange}}:
* the expected version is the one of {{oldExchange}} (its 
{{CamelOptimisticLockVersion}}), not the one of {{newExchange}};
* {{oldExchange == null}} and a row exists: {{OptimisticLockingException}} (as 
today);
* {{oldExchange != null}} and no row: {{OptimisticLockingException}} (new), so 
the aggregator retries and starts a new group with only its own message;
* a new group starts with a random positive version instead of 1, so a version 
read from an earlier group of the same key never matches the new group. This 
also stops a stale {{remove()}} from deleting a new group.

The schema is unchanged: the documented column is {{version BIGINT NOT NULL}}. 
Groups stored before the upgrade keep their version. The non-optimistic 
{{add(camelContext, key, exchange)}} keeps its behaviour (the aggregator's lock 
serializes it); only the version stored for a new group changes. 
{{ClusteredJdbcAggregationRepository}}, {{PostgresAggregationRepository}} and 
{{ClusteredPostgresAggregationRepository}} inherit {{add}}. Recovery does not 
compare versions and is not affected. A side effect of taking the version from 
the old exchange: an aggregation strategy that returns the new exchange no 
longer fails every optimistic {{add}}.

Tests: the three cases above as unit tests with a latch in the aggregation 
strategy, and repository-level tests for add after a removal, add after a new 
group was started, and a stale remove of a new group. A note for the 4.23 
upgrade guide (new groups no longer start at version 1).

Affected: all versions with optimistic locking for this repository (the same 
{{add}} code is at the camel-3.4.0, 4.0.0, 4.14.0 and 4.18.0 tags). Before 
4.22.0 {{remove()}} also ignored its delete count (CAMEL-24048).

Duplicate check (2026-09-29): JIRA text "JdbcAggregationRepository" (31 issues, 
including CAMEL-24048 stale remove, CAMEL-17613 confirm order, CAMEL-7810 and 
CAMEL-6144 optimistic locking), "aggregation repository" with "version", and 
component camel-sql since June 2026: none covers add() after a completion or 
the version restarting at 1. GitHub pull requests for 
"JdbcAggregationRepository" and "aggregation repository optimistic": 
apache/camel#24710 (CAMEL-24048), #24655, #24707, #26791, not this.

_Filed with Claude Code on behalf of allthingssecurity._




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

Reply via email to