This is an automated email from the ASF dual-hosted git repository.

davsclaus pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/camel.git


The following commit(s) were added to refs/heads/main by this push:
     new 2c32c7517370 CAMEL-25130: camel-sql - JdbcAggregationRepository with 
optimistic locking should not store a completed group again (#27039)
2c32c7517370 is described below

commit 2c32c751737007a5d873c11525e717ac243c7141
Author: allthingssecurity <[email protected]>
AuthorDate: Wed Sep 30 16:27:32 2026 +0530

    CAMEL-25130: camel-sql - JdbcAggregationRepository with optimistic locking 
should not store a completed group again (#27039)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../aggregate/jdbc/JdbcAggregationRepository.java  |  39 ++++-
 .../JdbcAggregateOptimisticCompletedGroupTest.java | 182 +++++++++++++++++++++
 ...JdbcAggregationRepositoryOptimisticAddTest.java | 134 +++++++++++++++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  13 ++
 4 files changed, 363 insertions(+), 5 deletions(-)

diff --git 
a/components/camel-sql/src/main/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregationRepository.java
 
b/components/camel-sql/src/main/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregationRepository.java
index 7d401ad90490..6f1a07858852 100644
--- 
a/components/camel-sql/src/main/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregationRepository.java
+++ 
b/components/camel-sql/src/main/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregationRepository.java
@@ -26,6 +26,7 @@ import java.util.LinkedHashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.ThreadLocalRandom;
 import java.util.concurrent.TimeUnit;
 
 import javax.sql.DataSource;
@@ -184,7 +185,10 @@ public class JdbcAggregationRepository extends 
ServiceSupport
             throws OptimisticLockingException {
 
         try {
-            return add(camelContext, correlationId, newExchange);
+            // compare-and-set against the exchange that was read by get(): 
the stored group must still have the version
+            // of oldExchange, and when oldExchange is null there must be no 
stored group
+            Long expectedVersion = oldExchange != null ? 
oldExchange.getProperty(VERSION_PROPERTY, Long.class) : null;
+            return doAdd(camelContext, correlationId, newExchange, 
expectedVersion, true);
         } catch (Exception e) {
             if (jdbcOptimisticLockingExceptionMapper != null && 
jdbcOptimisticLockingExceptionMapper.isOptimisticLocking(e)) {
                 throw new OptimisticLockingException();
@@ -196,6 +200,20 @@ public class JdbcAggregationRepository extends 
ServiceSupport
 
     @Override
     public Exchange add(final CamelContext camelContext, final String 
correlationId, final Exchange exchange) {
+        return doAdd(camelContext, correlationId, exchange, 
exchange.getProperty(VERSION_PROPERTY, Long.class), false);
+    }
+
+    /**
+     * Stores the exchange: updates the stored group if it still has the 
expected version, or inserts a new group.
+     *
+     * @param expectedVersion the version of the stored group the exchange was 
aggregated on, or <tt>null</tt> for a new
+     *                        group
+     * @param mustExist       whether the stored group must still exist when 
an expected version is given (optimistic
+     *                        locking), as it may have been completed and 
removed meanwhile
+     */
+    private Exchange doAdd(
+            final CamelContext camelContext, final String correlationId, final 
Exchange exchange,
+            final Long expectedVersion, final boolean mustExist) {
         return transactionTemplate.execute(new TransactionCallback<Exchange>() 
{
 
             public Exchange doInTransaction(TransactionStatus status) {
@@ -215,18 +233,21 @@ public class JdbcAggregationRepository extends 
ServiceSupport
                     }
 
                     if (present) {
-                        Long versionLong = 
exchange.getProperty(VERSION_PROPERTY, Long.class);
-                        if (versionLong == null) {
+                        if (expectedVersion == null) {
                             LOG.debug("Race while inserting record with key 
{}", correlationId);
                             throw new OptimisticLockingException();
                         } else {
-                            long version = versionLong.longValue();
+                            long version = expectedVersion.longValue();
                             LOG.debug("Updating record with key {} and version 
{}", correlationId, version);
                             update(camelContext, correlationId, exchange, 
table, version);
                         }
+                    } else if (mustExist && expectedVersion != null) {
+                        // the group the exchange was aggregated on has been 
completed (removed) meanwhile
+                        LOG.debug("Race while updating record with key {} as 
it has been removed", correlationId);
+                        throw new OptimisticLockingException();
                     } else {
                         LOG.debug("Inserting record with key {}", 
correlationId);
-                        insert(camelContext, correlationId, exchange, table, 
1L);
+                        insert(camelContext, correlationId, exchange, table, 
newGroupVersion());
                     }
 
                 } catch (Exception e) {
@@ -239,6 +260,14 @@ public class JdbcAggregationRepository extends 
ServiceSupport
         });
     }
 
+    /**
+     * The version of a new group. It is not 1 for every group, so that an 
exchange that was read from a previous group
+     * with the same correlation key (which has been completed meanwhile) can 
never match the version of the new group.
+     */
+    protected long newGroupVersion() {
+        return ThreadLocalRandom.current().nextLong(1, Long.MAX_VALUE / 2);
+    }
+
     // Useful to verify if the table name does not contain invalid characters.
     // Allows simple names (my_table) and schema-qualified names 
(myschema.my_table).
     protected static void verifyTableName(String tableName) {
diff --git 
a/components/camel-sql/src/test/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregateOptimisticCompletedGroupTest.java
 
b/components/camel-sql/src/test/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregateOptimisticCompletedGroupTest.java
new file mode 100644
index 000000000000..0be6a5e515ec
--- /dev/null
+++ 
b/components/camel-sql/src/test/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregateOptimisticCompletedGroupTest.java
@@ -0,0 +1,182 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.processor.aggregate.jdbc;
+
+import java.util.List;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import org.apache.camel.AggregationStrategy;
+import org.apache.camel.Exchange;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.component.mock.MockEndpoint;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+
+import static org.awaitility.Awaitility.await;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/**
+ * Aggregate EIP with optimistic locking: a message that is aggregated on a 
group while the group is completed by
+ * another message or the completion timeout must start a new group, and must 
not store the completed group again
+ * (duplicates) or overwrite a new group of the same key (lost messages).
+ */
+public class JdbcAggregateOptimisticCompletedGroupTest extends 
AbstractJdbcAggregationTestSupport {
+
+    private final AtomicBoolean hold = new AtomicBoolean();
+    private volatile CountDownLatch aggregating;
+    private volatile CountDownLatch release;
+    private JdbcAggregationRepository repo2;
+
+    @BeforeEach
+    public void resetLatches() {
+        aggregating = new CountDownLatch(1);
+        release = new CountDownLatch(1);
+        hold.set(true);
+    }
+
+    @Test
+    public void testGroupCompletedByAnotherMessage() throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:aggregated");
+        mock.expectedMessageCount(2);
+
+        template.sendBodyAndHeader("direct:predicate", "x", "id", "k1");
+        Thread threadA = sendAndHold("direct:predicate", "k1");
+        // the group [x] is completed while thread A aggregates on it
+        sendLast("direct:predicate", "b", "k1");
+        await().atMost(5, TimeUnit.SECONDS).until(() -> 
repo.scan(context).isEmpty());
+
+        release.countDown();
+        threadA.join(10000);
+        assertFalse(threadA.isAlive(), "the aggregating thread should have 
finished");
+        sendLast("direct:predicate", "z", "k1");
+
+        MockEndpoint.assertIsSatisfied(context);
+        // x is sent once, a starts a new group
+        assertEquals(List.of("x,b", "a,z"), bodies(mock));
+    }
+
+    @Test
+    public void testGroupCompletedAndNewGroupStarted() throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:aggregated");
+        mock.expectedMessageCount(2);
+
+        template.sendBodyAndHeader("direct:predicate", "x", "id", "k2");
+        Thread threadA = sendAndHold("direct:predicate", "k2");
+        sendLast("direct:predicate", "b", "k2");
+        // a new group for the same key
+        template.sendBodyAndHeader("direct:predicate", "c", "id", "k2");
+        await().atMost(5, TimeUnit.SECONDS).until(() -> 
repo.scan(context).isEmpty());
+
+        release.countDown();
+        threadA.join(10000);
+        assertFalse(threadA.isAlive(), "the aggregating thread should have 
finished");
+        sendLast("direct:predicate", "z", "k2");
+
+        MockEndpoint.assertIsSatisfied(context);
+        // x is sent once, and c is not lost
+        assertEquals(List.of("x,b", "c,a,z"), bodies(mock));
+    }
+
+    @Test
+    public void testGroupCompletedByTimeout() throws Exception {
+        MockEndpoint mock = getMockEndpoint("mock:timeout");
+        mock.expectedMessageCount(1);
+
+        template.sendBodyAndHeader("direct:timeout", "x", "id", "k3");
+        Thread threadA = sendAndHold("direct:timeout", "k3");
+        // the group [x] is completed by the completion timeout while thread A 
aggregates on it
+        MockEndpoint.assertIsSatisfied(context);
+        await().atMost(5, TimeUnit.SECONDS).until(() -> 
repo2.scan(context).isEmpty());
+
+        release.countDown();
+        threadA.join(10000);
+        assertFalse(threadA.isAlive(), "the aggregating thread should have 
finished");
+
+        await().atMost(5, TimeUnit.SECONDS).until(() -> 
mock.getReceivedCounter() == 2);
+        // x is sent once, a starts a new group that completes by timeout too
+        assertEquals(List.of("x", "a"), bodies(mock));
+    }
+
+    private static List<String> bodies(MockEndpoint mock) {
+        return mock.getReceivedExchanges().stream().map(e -> 
e.getMessage().getBody(String.class)).toList();
+    }
+
+    // sends "a" on another thread and waits until the aggregation strategy 
got the group it read
+    private Thread sendAndHold(String uri, String key) throws 
InterruptedException {
+        Thread thread = new Thread(() -> template.sendBodyAndHeader(uri, "a", 
"id", key), "thread-A");
+        thread.start();
+        assertTrue(aggregating.await(5, TimeUnit.SECONDS), "message a should 
be aggregating");
+        return thread;
+    }
+
+    private void sendLast(String uri, String body, String key) {
+        template.send(uri, e -> {
+            e.getMessage().setBody(body);
+            e.getMessage().setHeader("id", key);
+            e.getMessage().setHeader("last", true);
+        });
+    }
+
+    @Override
+    protected RouteBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                configureJdbcAggregationRepository();
+                repo2 = applicationContext.getBean("repo2", 
JdbcAggregationRepository.class);
+
+                from("direct:predicate")
+                        .aggregate(header("id"), new 
HoldingAggregationStrategy())
+                        .aggregationRepository(repo).optimisticLocking()
+                        
.eagerCheckCompletion().completionPredicate(header("last").isEqualTo(true))
+                        .to("mock:aggregated");
+
+                from("direct:timeout")
+                        .aggregate(header("id"), new 
HoldingAggregationStrategy())
+                        .aggregationRepository(repo2).optimisticLocking()
+                        
.completionTimeout(300).completionTimeoutCheckerInterval(50)
+                        .to("mock:timeout");
+            }
+        };
+    }
+
+    private final class HoldingAggregationStrategy implements 
AggregationStrategy {
+
+        @Override
+        public Exchange aggregate(Exchange oldExchange, Exchange newExchange) {
+            String body = newExchange.getMessage().getBody(String.class);
+            if (oldExchange != null && "a".equals(body) && 
hold.compareAndSet(true, false)) {
+                // message a has read the group, hold it until the group has 
been completed
+                aggregating.countDown();
+                try {
+                    release.await(10, TimeUnit.SECONDS);
+                } catch (InterruptedException e) {
+                    Thread.currentThread().interrupt();
+                }
+            }
+            if (oldExchange == null) {
+                return newExchange;
+            }
+            
oldExchange.getMessage().setBody(oldExchange.getMessage().getBody(String.class) 
+ "," + body);
+            return oldExchange;
+        }
+    }
+}
diff --git 
a/components/camel-sql/src/test/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregationRepositoryOptimisticAddTest.java
 
b/components/camel-sql/src/test/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregationRepositoryOptimisticAddTest.java
new file mode 100644
index 000000000000..ea913d8d3d4d
--- /dev/null
+++ 
b/components/camel-sql/src/test/java/org/apache/camel/processor/aggregate/jdbc/JdbcAggregationRepositoryOptimisticAddTest.java
@@ -0,0 +1,134 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *      http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+package org.apache.camel.processor.aggregate.jdbc;
+
+import org.apache.camel.Exchange;
+import 
org.apache.camel.spi.OptimisticLockingAggregationRepository.OptimisticLockingException;
+import org.apache.camel.support.DefaultExchange;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/**
+ * The optimistic locking add(camelContext, key, oldExchange, newExchange) is 
a compare-and-set against the exchange
+ * that was read by get(), also when the group has been completed (removed) 
meanwhile, and a new group with the same key
+ * never gets the version of a previous group.
+ */
+public class JdbcAggregationRepositoryOptimisticAddTest extends 
AbstractJdbcAggregationTestSupport {
+
+    private static final String VERSION = "CamelOptimisticLockVersion";
+
+    @Test
+    public void testAddNewGroup() {
+        repo.add(context, "foo", null, newExchange("a"));
+
+        assertEquals("a", repo.get(context, 
"foo").getMessage().getBody(String.class));
+    }
+
+    @Test
+    public void testAddNewGroupWhenGroupExists() {
+        repo.add(context, "foo", null, newExchange("a"));
+
+        // another thread or node started the group first
+        assertThrows(OptimisticLockingException.class, () -> repo.add(context, 
"foo", null, newExchange("b")));
+        assertEquals("a", repo.get(context, 
"foo").getMessage().getBody(String.class));
+    }
+
+    @Test
+    public void testUpdateGroup() {
+        repo.add(context, "foo", null, newExchange("a"));
+        Exchange group = repo.get(context, "foo");
+        group.getMessage().setBody("a,b");
+
+        repo.add(context, "foo", group, group);
+        assertEquals("a,b", repo.get(context, 
"foo").getMessage().getBody(String.class));
+
+        // an update with the same (now stale) version fails
+        Exchange stale = group;
+        assertThrows(OptimisticLockingException.class, () -> repo.add(context, 
"foo", stale, stale));
+    }
+
+    @Test
+    public void testUpdateGroupWithNewExchange() {
+        repo.add(context, "foo", null, newExchange("a"));
+        Exchange group = repo.get(context, "foo");
+
+        // an aggregation strategy that returns the new exchange: the version 
is the one of the old exchange
+        repo.add(context, "foo", group, newExchange("a,b"));
+        assertEquals("a,b", repo.get(context, 
"foo").getMessage().getBody(String.class));
+    }
+
+    @Test
+    public void testAddAfterGroupCompleted() {
+        repo.add(context, "foo", null, newExchange("x"));
+        // thread A reads the group [x]
+        Exchange readByA = repo.get(context, "foo");
+        // meanwhile the group is completed by another thread
+        Exchange completed = repo.get(context, "foo");
+        repo.remove(context, "foo", completed);
+        repo.confirm(context, completed.getExchangeId());
+
+        // thread A aggregates and stores the group it read: it must not be 
stored again (x would be sent twice)
+        readByA.getMessage().setBody("x,a");
+        assertThrows(OptimisticLockingException.class, () -> repo.add(context, 
"foo", readByA, readByA));
+        assertNull(repo.get(context, "foo"));
+    }
+
+    @Test
+    public void testAddAfterGroupCompletedAndNewGroupStarted() {
+        repo.add(context, "foo", null, newExchange("x"));
+        Exchange readByA = repo.get(context, "foo");
+        Exchange completed = repo.get(context, "foo");
+        repo.remove(context, "foo", completed);
+        repo.confirm(context, completed.getExchangeId());
+        // a new group is started for the same key
+        repo.add(context, "foo", null, newExchange("c"));
+        Exchange newGroup = repo.get(context, "foo");
+        assertNotEquals(readByA.getProperty(VERSION, Long.class), 
newGroup.getProperty(VERSION, Long.class));
+
+        // thread A must not overwrite the new group (c would be lost)
+        readByA.getMessage().setBody("x,a");
+        assertThrows(OptimisticLockingException.class, () -> repo.add(context, 
"foo", readByA, readByA));
+        assertEquals("c", repo.get(context, 
"foo").getMessage().getBody(String.class));
+    }
+
+    @Test
+    public void testRemoveAfterGroupCompletedAndNewGroupStarted() {
+        repo.add(context, "foo", null, newExchange("x"));
+        Exchange readByA = repo.get(context, "foo");
+        Exchange completed = repo.get(context, "foo");
+        repo.remove(context, "foo", completed);
+        repo.confirm(context, completed.getExchangeId());
+        repo.add(context, "foo", null, newExchange("c"));
+
+        // a stale completion must not remove the new group
+        assertThrows(OptimisticLockingException.class, () -> 
repo.remove(context, "foo", readByA));
+        Exchange stored = repo.get(context, "foo");
+        assertNotNull(stored);
+        assertEquals("c", stored.getMessage().getBody(String.class));
+    }
+
+    private Exchange newExchange(String body) {
+        Exchange exchange = new DefaultExchange(context);
+        exchange.getMessage().setBody(body);
+        return exchange;
+    }
+}
diff --git 
a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc 
b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
index 9164fc70f888..1040b2966262 100644
--- a/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
+++ b/docs/user-manual/modules/ROOT/pages/camel-4x-upgrade-guide-4_23.adoc
@@ -244,6 +244,19 @@ With optimistic locking, or a persistent aggregation 
repository, the spool file
 exchange has been added to the repository (a persistent repository has read 
the body at that time), and only the
 spool file of the exchange that completes a group is kept until the aggregated 
exchange is done.
 
+=== camel-sql - JdbcAggregationRepository with optimistic locking
+
+With optimistic locking, `JdbcAggregationRepository` (and 
`ClusteredJdbcAggregationRepository`,
+`PostgresAggregationRepository` and `ClusteredPostgresAggregationRepository`) 
now fails with an
+`OptimisticLockingException`, so that the aggregator retries, when a message 
was aggregated on a group that another
+thread or Camel instance completed meanwhile. Prior to Camel 4.23 the 
completed group was stored again together with
+the new message, so its messages were sent a second time, or a new group of 
the same correlation key was overwritten
+and its messages were lost.
+
+A new group now starts with a random positive version instead of `1`, so that 
a version read from an earlier group of
+the same correlation key never matches it. The `version` column must therefore 
be a 64-bit integer, as in the
+documented `version BIGINT NOT NULL`. Groups stored before the upgrade keep 
their version.
+
 === Variable Receive
 
 When an EIP with `variableReceive` stores a message into a variable that 
already holds a message, the header variables

Reply via email to