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)