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

oscerd 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 c3632125413b CAMEL-25159, CAMEL-25160: camel-mongodb - stop running 
queries twice, and harden tail tracking (#27116)
c3632125413b is described below

commit c3632125413bec6fe6de557e0898b592b3ff684b
Author: Andrea Cosentino <[email protected]>
AuthorDate: Wed Sep 30 11:46:35 2026 +0200

    CAMEL-25159, CAMEL-25160: camel-mongodb - stop running queries twice, and 
harden tail tracking (#27116)
    
    CAMEL-25159: findAll, distinct and aggregate read their results with
    ret.iterator().forEachRemaining(...) and then, in a finally block meant to 
release
    the cursor, called ret.iterator().close(). MongoIterable.iterator() 
executes the
    operation, so the finally opened and closed a second cursor (a second round 
trip,
    and for aggregate a second run of the pipeline) and left the cursor that 
was read
    unclosed. Hold the cursor with try-with-resources, so the query runs once 
and the
    cursor that is read is the one that is closed.
    
    CAMEL-25160: persistToStore stored the full document returned by 
findOneAndUpdate
    back into trackingObj, and from then on its filter also matched the previous
    value, so another writer of that field (two routes sharing a persistentId) 
left
    the update matching nothing and the next persist throwing. Filter on the id 
only.
    recoverFromStore now handles a missing tracking document by warning and 
starting
    from the beginning.
    
    Co-Authored-By: Claude Opus 5 <[email protected]>
    Signed-off-by: Andrea Cosentino <[email protected]>
---
 components/camel-mongodb/pom.xml                   |   6 +
 .../camel/component/mongodb/MongoDbProducer.java   |  26 ++--
 .../mongodb/MongoDbTailTrackingManager.java        |  15 ++-
 .../mongodb/MongoDbProducerCursorTest.java         | 140 +++++++++++++++++++++
 4 files changed, 168 insertions(+), 19 deletions(-)

diff --git a/components/camel-mongodb/pom.xml b/components/camel-mongodb/pom.xml
index faec97739249..64002f9b2a1c 100644
--- a/components/camel-mongodb/pom.xml
+++ b/components/camel-mongodb/pom.xml
@@ -92,6 +92,12 @@
             <scope>test</scope>
         </dependency>
 
+        <dependency>
+            <groupId>org.mockito</groupId>
+            <artifactId>mockito-core</artifactId>
+            <version>${mockito-version}</version>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.apache.commons</groupId>
             <artifactId>commons-lang3</artifactId>
diff --git 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
index 9383d8e3ae00..0cc7af5dc8ea 100644
--- 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
+++ 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbProducer.java
@@ -30,6 +30,7 @@ import com.mongodb.client.AggregateIterable;
 import com.mongodb.client.DistinctIterable;
 import com.mongodb.client.FindIterable;
 import com.mongodb.client.MongoCollection;
+import com.mongodb.client.MongoCursor;
 import com.mongodb.client.MongoDatabase;
 import com.mongodb.client.model.BulkWriteOptions;
 import com.mongodb.client.model.Filters;
@@ -362,11 +363,10 @@ public class MongoDbProducer extends DefaultProducer {
             } else {
                 ret = dbCol.distinct(distinctFieldName, String.class);
             }
-            try {
-                ret.iterator().forEachRemaining(result::add);
+            // hold the cursor: every call to iterator() executes the query 
again
+            try (MongoCursor<String> cursor = ret.iterator()) {
+                cursor.forEachRemaining(result::add);
                 
exchange.getMessage().setHeader(MongoDbConstants.RESULT_PAGE_SIZE, 
result.size());
-            } finally {
-                ret.iterator().close();
             }
             return result;
         };
@@ -420,12 +420,11 @@ public class MongoDbProducer extends DefaultProducer {
 
             
ret.allowDiskUse(exchange.getIn().getHeader(MongoDbConstants.ALLOW_DISK_USE, 
Boolean.class));
             if 
(!MongoDbOutputType.MongoIterable.equals(endpoint.getOutputType())) {
-                try {
-                    result = new ArrayList<>();
-                    ret.iterator().forEachRemaining(((List<Document>) 
result)::add);
+                result = new ArrayList<>();
+                // hold the cursor: every call to iterator() executes the 
query again
+                try (MongoCursor<Document> cursor = ret.iterator()) {
+                    cursor.forEachRemaining(((List<Document>) result)::add);
                     exchange.getMessage().setHeader(RESULT_PAGE_SIZE, 
((List<Document>) result).size());
-                } finally {
-                    ret.iterator().close();
                 }
             } else {
                 result = ret;
@@ -572,12 +571,11 @@ public class MongoDbProducer extends DefaultProducer {
 
                 Iterable<Document> result;
                 if 
(!MongoDbOutputType.MongoIterable.equals(endpoint.getOutputType())) {
-                    try {
-                        result = new ArrayList<>();
-                        
aggregationResult.iterator().forEachRemaining(((List<Document>) result)::add);
+                    result = new ArrayList<>();
+                    // hold the cursor: every call to iterator() executes the 
pipeline again
+                    try (MongoCursor<Document> cursor = 
aggregationResult.iterator()) {
+                        cursor.forEachRemaining(((List<Document>) 
result)::add);
                         
exchange.getMessage().setHeader(MongoDbConstants.RESULT_PAGE_SIZE, 
((List<Document>) result).size());
-                    } finally {
-                        aggregationResult.iterator().close();
                     }
                 } else {
                     result = aggregationResult;
diff --git 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
index 653d32b03ea3..074f7602b56d 100644
--- 
a/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
+++ 
b/components/camel-mongodb/src/main/java/org/apache/camel/component/mongodb/MongoDbTailTrackingManager.java
@@ -21,8 +21,6 @@ import java.util.concurrent.locks.ReentrantLock;
 
 import com.mongodb.client.MongoClient;
 import com.mongodb.client.MongoCollection;
-import com.mongodb.client.model.FindOneAndUpdateOptions;
-import com.mongodb.client.model.ReturnDocument;
 import com.mongodb.client.model.Updates;
 import org.bson.Document;
 import org.bson.conversions.Bson;
@@ -76,8 +74,9 @@ public class MongoDbTailTrackingManager {
             }
 
             Bson updateObj = Updates.set(config.field, lastVal);
-            FindOneAndUpdateOptions options = new 
FindOneAndUpdateOptions().returnDocument(ReturnDocument.AFTER);
-            trackingObj = dbCol.findOneAndUpdate(trackingObj, updateObj, 
options);
+            // filter by the id only: storing the returned document back would 
make every later update
+            // also match on the previous value, and a filter that stops 
matching leaves trackingObj null
+            dbCol.updateOne(trackingObj, updateObj);
         } finally {
             lock.unlock();
         }
@@ -90,7 +89,13 @@ public class MongoDbTailTrackingManager {
                 return null;
             }
 
-            lastVal = dbCol.find(trackingObj).first().get(config.field);
+            Document tracked = dbCol.find(trackingObj).first();
+            if (tracked == null) {
+                LOG.warn("Tail tracking document {} is gone from collection 
{}, starting from the beginning",
+                        trackingObj, config.collection);
+                return null;
+            }
+            lastVal = tracked.get(config.field);
 
             if (LOG.isDebugEnabled()) {
                 LOG.debug("Recovered lastVal={} from store, collection: {}", 
lastVal, config.collection);
diff --git 
a/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbProducerCursorTest.java
 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbProducerCursorTest.java
new file mode 100644
index 000000000000..b10caf79ca5a
--- /dev/null
+++ 
b/components/camel-mongodb/src/test/java/org/apache/camel/component/mongodb/MongoDbProducerCursorTest.java
@@ -0,0 +1,140 @@
+/*
+ * 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.component.mongodb;
+
+import com.mongodb.client.AggregateIterable;
+import com.mongodb.client.DistinctIterable;
+import com.mongodb.client.FindIterable;
+import com.mongodb.client.MongoClient;
+import com.mongodb.client.MongoCollection;
+import com.mongodb.client.MongoCursor;
+import com.mongodb.client.MongoDatabase;
+import org.apache.camel.CamelContext;
+import org.apache.camel.RoutesBuilder;
+import org.apache.camel.builder.RouteBuilder;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.bson.Document;
+import org.junit.jupiter.api.Test;
+
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * Every call to MongoIterable.iterator() executes the operation again - it is 
not an accessor for an already open
+ * cursor - so closing "the" cursor in a finally block by calling iterator() a 
second time ran the whole query twice and
+ * left the cursor that had actually been read unclosed.
+ */
+class MongoDbProducerCursorTest extends CamelTestSupport {
+
+    private final MongoCollection<Document> collection = 
mock(MongoCollection.class);
+    private final DistinctIterable<String> distinctIterable = 
mock(DistinctIterable.class);
+    private final FindIterable<Document> findIterable = 
mock(FindIterable.class);
+    private final AggregateIterable<Document> aggregateIterable = 
mock(AggregateIterable.class);
+
+    @Test
+    void testDistinctOpensOneCursor() {
+        // findDistinct needs the field name, otherwise the operation has 
nothing to ask for
+        template.requestBodyAndHeader("direct:distinct", null, 
MongoDbConstants.DISTINCT_QUERY_FIELD, "name");
+
+        verify(distinctIterable, times(1)).iterator();
+    }
+
+    @Test
+    void testFindAllOpensOneCursor() {
+        template.requestBody("direct:findAll", (Object) null);
+
+        verify(findIterable, times(1)).iterator();
+    }
+
+    @Test
+    void testAggregateOpensOneCursor() {
+        template.requestBody("direct:aggregate", "[{$match: {}}]");
+
+        verify(aggregateIterable, times(1)).iterator();
+    }
+
+    @Override
+    protected CamelContext createCamelContext() throws Exception {
+        final CamelContext context = super.createCamelContext();
+
+        when(collection.withWriteConcern(any())).thenReturn(collection);
+        when(collection.distinct(anyString(), 
eq(String.class))).thenReturn(distinctIterable);
+        when(collection.distinct(anyString(), any(), 
eq(String.class))).thenReturn(distinctIterable);
+        
when(collection.find(any(org.bson.conversions.Bson.class))).thenReturn(findIterable);
+        when(collection.find()).thenReturn(findIterable);
+        when(collection.aggregate(any())).thenReturn(aggregateIterable);
+
+        stubCursor(distinctIterable);
+        stubFind(findIterable);
+        stubAggregate(aggregateIterable);
+
+        final MongoDatabase database = mock(MongoDatabase.class);
+        when(database.getCollection(anyString(), 
eq(Document.class))).thenReturn(collection);
+
+        final MongoClient client = mock(MongoClient.class);
+        when(client.getDatabase(anyString())).thenReturn(database);
+        
when(client.listDatabaseNames()).thenReturn(mock(com.mongodb.client.MongoIterable.class));
+
+        context.getRegistry().bind("mongoClient", client);
+        return context;
+    }
+
+    private void stubCursor(DistinctIterable<String> iterable) {
+        final MongoCursor<String> cursor = mock(MongoCursor.class);
+        when(cursor.hasNext()).thenReturn(false);
+        when(iterable.iterator()).thenReturn(cursor);
+    }
+
+    private void stubFind(FindIterable<Document> iterable) {
+        final MongoCursor<Document> cursor = mock(MongoCursor.class);
+        when(cursor.hasNext()).thenReturn(false);
+        when(iterable.iterator()).thenReturn(cursor);
+        when(iterable.projection(any())).thenReturn(iterable);
+        when(iterable.sort(any())).thenReturn(iterable);
+        
when(iterable.skip(org.mockito.ArgumentMatchers.anyInt())).thenReturn(iterable);
+        
when(iterable.limit(org.mockito.ArgumentMatchers.anyInt())).thenReturn(iterable);
+        when(iterable.allowDiskUse(any())).thenReturn(iterable);
+    }
+
+    private void stubAggregate(AggregateIterable<Document> iterable) {
+        final MongoCursor<Document> cursor = mock(MongoCursor.class);
+        when(cursor.hasNext()).thenReturn(false);
+        when(iterable.iterator()).thenReturn(cursor);
+        when(iterable.allowDiskUse(any())).thenReturn(iterable);
+        
when(iterable.batchSize(org.mockito.ArgumentMatchers.anyInt())).thenReturn(iterable);
+    }
+
+    @Override
+    protected RoutesBuilder createRouteBuilder() {
+        return new RouteBuilder() {
+            @Override
+            public void configure() {
+                from("direct:distinct")
+                        
.to("mongodb:mongoClient?database=test&collection=c&operation=findDistinct&dynamicity=false");
+                from("direct:findAll")
+                        
.to("mongodb:mongoClient?database=test&collection=c&operation=findAll&dynamicity=false");
+                from("direct:aggregate")
+                        
.to("mongodb:mongoClient?database=test&collection=c&operation=aggregate&dynamicity=false");
+            }
+        };
+    }
+}

Reply via email to