[ 
https://issues.apache.org/jira/browse/CAMEL-25130?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
 ]

Claus Ibsen updated CAMEL-25130:
--------------------------------
    Fix Version/s: 4.23.0

> 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
>            Priority: Minor
>             Fix For: 4.23.0
>
>
> 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