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 c866afdce9d6 CAMEL-25290: camel-couchdb - the consumer must apply 
updates=false, and move past the changes it skips (#27322)
c866afdce9d6 is described below

commit c866afdce9d6e227610dc0f9d92ed4ed0ac8acea
Author: allthingssecurity <[email protected]>
AuthorDate: Sun Oct 4 12:33:40 2026 +0530

    CAMEL-25290: camel-couchdb - the consumer must apply updates=false, and 
move past the changes it skips (#27322)
    
    Co-Authored-By: Claude Opus 5.5 <[email protected]>
---
 .../camel/component/couchdb/CouchDbConsumer.java   |  16 ++-
 .../couchdb/CouchDbConsumerFilterTest.java         | 110 +++++++++++++++++++++
 2 files changed, 117 insertions(+), 9 deletions(-)

diff --git 
a/components/camel-couchdb/src/main/java/org/apache/camel/component/couchdb/CouchDbConsumer.java
 
b/components/camel-couchdb/src/main/java/org/apache/camel/component/couchdb/CouchDbConsumer.java
index 53a4e379c011..6be0b6483665 100644
--- 
a/components/camel-couchdb/src/main/java/org/apache/camel/component/couchdb/CouchDbConsumer.java
+++ 
b/components/camel-couchdb/src/main/java/org/apache/camel/component/couchdb/CouchDbConsumer.java
@@ -78,19 +78,17 @@ public class CouchDbConsumer extends 
ScheduledBatchPollingConsumer implements Re
                 = couchClient.pollChanges(endpoint.getStyle(), since, 
endpoint.getHeartbeat(), getMaxMessagesPerPoll());
 
         for (ChangesResultItem changesResultItem : 
changesResultResponse.getResult().getResults()) {
-            if (changesResultItem.isDeleted() != null) {
-                if (changesResultItem.isDeleted() && !endpoint.isDeletes()) {
-                    continue;
-                }
-                if (!changesResultItem.isDeleted() && !endpoint.isUpdates()) {
-                    continue;
-                }
+            // CouchDB sets deleted only on the changes of deleted documents
+            boolean deleted = 
Boolean.TRUE.equals(changesResultItem.isDeleted());
+            if (deleted ? !endpoint.isDeletes() : !endpoint.isUpdates()) {
+                // move past the skipped change, or a page of skipped changes 
is polled again forever
+                since = changesResultItem.getSeq();
+                continue;
             }
 
             lastSequence = changesResultItem.getSeq();
 
-            Exchange exchange = this.createExchange(lastSequence, 
changesResultItem.getId(), changesResultItem,
-                    changesResultItem.isDeleted() == null ? false : 
changesResultItem.isDeleted());
+            Exchange exchange = this.createExchange(lastSequence, 
changesResultItem.getId(), changesResultItem, deleted);
 
             if (LOG.isTraceEnabled()) {
                 LOG.trace("Created exchange [exchange={}, _id={}, seq={}", 
exchange, changesResultItem.getId(), lastSequence);
diff --git 
a/components/camel-couchdb/src/test/java/org/apache/camel/component/couchdb/CouchDbConsumerFilterTest.java
 
b/components/camel-couchdb/src/test/java/org/apache/camel/component/couchdb/CouchDbConsumerFilterTest.java
new file mode 100644
index 000000000000..48091adaab5c
--- /dev/null
+++ 
b/components/camel-couchdb/src/test/java/org/apache/camel/component/couchdb/CouchDbConsumerFilterTest.java
@@ -0,0 +1,110 @@
+/*
+ * 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.couchdb;
+
+import java.util.List;
+import java.util.concurrent.CopyOnWriteArrayList;
+
+import com.ibm.cloud.cloudant.v1.model.ChangesResult;
+import com.ibm.cloud.sdk.core.http.Response;
+import com.ibm.cloud.sdk.core.util.GsonSingleton;
+import org.apache.camel.Exchange;
+import org.apache.camel.test.junit6.CamelTestSupport;
+import org.junit.jupiter.api.Test;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/**
+ * The consumer must skip the changes that its options filter out, and move 
past them.
+ */
+public class CouchDbConsumerFilterTest extends CamelTestSupport {
+
+    // CouchDB sets "deleted" only on the changes of deleted documents
+    private static final String UPDATE = 
"{\"seq\":\"%s\",\"id\":\"%s\",\"changes\":[{\"rev\":\"1-a\"}]}";
+    private static final String DELETE = 
"{\"seq\":\"%s\",\"id\":\"%s\",\"changes\":[{\"rev\":\"2-b\"}],\"deleted\":true}";
+
+    private final CouchDbClientWrapper client = 
mock(CouchDbClientWrapper.class);
+    private final List<Exchange> received = new CopyOnWriteArrayList<>();
+
+    @Test
+    void testDeletesFalseMovesPastDeletedDocuments() throws Exception {
+        CouchDbConsumer consumer = createConsumer("deletes=false");
+        // a whole page of deleted documents, which this consumer does not 
publish
+        Response<ChangesResult> page = changes(DELETE.formatted("1-x", 
"doc1"), DELETE.formatted("2-x", "doc2"),
+                DELETE.formatted("3-x", "doc3"));
+        when(client.pollChanges(anyString(), anyString(), anyLong(), 
anyLong())).thenReturn(page);
+
+        consumer.poll();
+        consumer.poll();
+
+        assertEquals(0, received.size());
+        // the second poll must ask for the changes after the last one seen, 
not for the same page again
+        verify(client).pollChanges(anyString(), eq("3-x"), anyLong(), 
anyLong());
+    }
+
+    @Test
+    void testUpdatesFalseSkipsUpdatedDocuments() throws Exception {
+        CouchDbConsumer consumer = createConsumer("updates=false");
+        Response<ChangesResult> page = changes(UPDATE.formatted("1-x", 
"doc1"), DELETE.formatted("2-x", "doc2"));
+        when(client.pollChanges(anyString(), anyString(), anyLong(), 
anyLong())).thenReturn(page);
+
+        consumer.poll();
+
+        assertEquals(1, received.size());
+        assertEquals("doc2", 
received.get(0).getIn().getHeader(CouchDbConstants.HEADER_DOC_ID));
+        assertEquals("DELETE", 
received.get(0).getIn().getHeader(CouchDbConstants.HEADER_METHOD));
+    }
+
+    @Test
+    void testDefaultPublishesUpdatesAndDeletes() throws Exception {
+        CouchDbConsumer consumer = createConsumer("updates=true");
+        Response<ChangesResult> page = changes(UPDATE.formatted("1-x", 
"doc1"), DELETE.formatted("2-x", "doc2"));
+        when(client.pollChanges(anyString(), anyString(), anyLong(), 
anyLong())).thenReturn(page);
+
+        consumer.poll();
+        consumer.poll();
+
+        assertEquals(4, received.size());
+        assertEquals("UPDATE", 
received.get(0).getIn().getHeader(CouchDbConstants.HEADER_METHOD));
+        assertEquals("DELETE", 
received.get(1).getIn().getHeader(CouchDbConstants.HEADER_METHOD));
+        verify(client).pollChanges(anyString(), eq("2-x"), anyLong(), 
anyLong());
+    }
+
+    private CouchDbConsumer createConsumer(String options) throws Exception {
+        when(client.getLatestUpdateSequence()).thenReturn("0");
+        CouchDbEndpoint endpoint = 
context.getEndpoint("couchdb:http://localhost:5984/camel?"; + options, 
CouchDbEndpoint.class);
+        CouchDbConsumer consumer = new CouchDbConsumer(endpoint, client, 
received::add);
+        consumer.setMaxMessagesPerPoll(endpoint.getMaxMessagesPerPoll());
+        consumer.init();
+        return consumer;
+    }
+
+    @SuppressWarnings("unchecked")
+    private static Response<ChangesResult> changes(String... items) {
+        String json = "{\"last_seq\":\"9-x\",\"pending\":0,\"results\":[" + 
String.join(",", items) + "]}";
+        ChangesResult result = GsonSingleton.getGson().fromJson(json, 
ChangesResult.class);
+        Response<ChangesResult> response = mock(Response.class);
+        when(response.getResult()).thenReturn(result);
+        return response;
+    }
+}

Reply via email to