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;
+ }
+}