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 80fd91fbef34 CAMEL-25169, CAMEL-25170, CAMEL-25171, CAMEL-25172: 
camel-couchbase - close the connection, apply connectTimeout, report route 
failures, accept persistTo=2 (#27123)
80fd91fbef34 is described below

commit 80fd91fbef3404f3aa814b3fdc6afb03c451941e
Author: Andrea Cosentino <[email protected]>
AuthorDate: Thu Oct 1 17:41:51 2026 +0200

    CAMEL-25169, CAMEL-25170, CAMEL-25171, CAMEL-25172: camel-couchbase - close 
the connection, apply connectTimeout, report route failures, accept persistTo=2 
(#27123)
    
    Co-Authored-By: Claude Opus 5 <[email protected]>
---
 .../component/couchbase/CouchbaseConsumer.java     |  78 +++++++---
 .../component/couchbase/CouchbaseEndpoint.java     |  90 +++++++++---
 .../component/couchbase/CouchbaseProducer.java     |  51 ++++---
 .../component/couchbase/CouchbaseConsumerTest.java | 163 +++++++++++++++++++++
 .../component/couchbase/CouchbaseEndpointTest.java | 140 ++++++++++++++++++
 .../component/couchbase/CouchbaseProducerTest.java |  13 ++
 .../ROOT/pages/camel-4x-upgrade-guide-4_23.adoc    |  21 +++
 7 files changed, 489 insertions(+), 67 deletions(-)

diff --git 
a/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseConsumer.java
 
b/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseConsumer.java
index d7068391d38a..8200004c8cf7 100644
--- 
a/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseConsumer.java
+++ 
b/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseConsumer.java
@@ -75,19 +75,7 @@ public class CouchbaseConsumer extends 
ScheduledBatchPollingConsumer implements
     protected void doInit() throws Exception {
         super.doInit();
 
-        if (endpoint.getScope() != null) {
-            this.scope = bucket.scope(endpoint.getScope());
-        } else {
-            this.scope = bucket.defaultScope();
-        }
-
-        if (endpoint.getCollection() != null) {
-            this.collection = scope.collection(endpoint.getCollection());
-        } else {
-            this.collection = bucket.defaultCollection();
-        }
-
-        // Determine query mode
+        // Determine query mode. This reads only endpoint options, so it does 
not need a connection
         if (endpoint.getStatement() != null) {
             // Explicit SQL++ statement provided
             useSqlQuery = true;
@@ -131,16 +119,25 @@ public class CouchbaseConsumer extends 
ScheduledBatchPollingConsumer implements
 
     @Override
     protected void doStart() throws Exception {
-        super.doStart();
-        ResumeStrategyHelper.resume(getEndpoint().getCamelContext(), this, 
resumeStrategy, COUCHBASE_RESUME_ACTION);
-    }
+        // Take the bucket, scope and collection again on every start. The 
endpoint owns the cluster and
+        // disconnects it when it stops, so handles resolved once at init 
would point at a dead cluster after a
+        // restart in place. They used to survive only because the old close 
was a no-op
+        bucket = endpoint.createClient();
 
-    @Override
-    protected void doStop() throws Exception {
-        super.doStop();
-        if (bucket != null) {
-            bucket.core().shutdown();
+        if (endpoint.getScope() != null) {
+            this.scope = bucket.scope(endpoint.getScope());
+        } else {
+            this.scope = bucket.defaultScope();
+        }
+
+        if (endpoint.getCollection() != null) {
+            this.collection = scope.collection(endpoint.getCollection());
+        } else {
+            this.collection = bucket.defaultCollection();
         }
+
+        super.doStart();
+        ResumeStrategyHelper.resume(getEndpoint().getCamelContext(), this, 
resumeStrategy, COUCHBASE_RESUME_ACTION);
     }
 
     @Override
@@ -296,12 +293,49 @@ public class CouchbaseConsumer extends 
ScheduledBatchPollingConsumer implements
             exchange.setProperty(ExchangePropertyKey.BATCH_SIZE, total);
             exchange.setProperty(ExchangePropertyKey.BATCH_COMPLETE, index == 
total - 1);
             this.pendingExchanges = total - index - 1;
-            getProcessor().process(exchange);
+            processExchange(exchange);
+        }
+
+        // Anything still queued was never handed to the route - the batch was 
cut short because the consumer is
+        // stopping, or the poll returned more rows than maxMessagesPerPoll. 
Nothing else will release these, and a
+        // pooled exchange that is never released never returns to the pool.
+        //
+        // Releasing them is not the same as handling them: with 
consumerProcessedStrategy=delete the document was
+        // already removed during the poll, above, so these rows are lost 
rather than redelivered. That predates
+        // this method and is not fixed here - see CAMEL-25221
+        Exchange remaining;
+        while ((remaining = (Exchange) exchanges.poll()) != null) {
+            releaseExchange(remaining, false);
         }
 
         return answer;
     }
 
+    /**
+     * Hands the exchange to the route and reports a failure through the 
consumer's exception handler.
+     * <p/>
+     * A failing route does not throw out of {@code process()} - the consumer 
processor is asynchronous, so the failure
+     * is left on the exchange instead. That is why the exception is read back 
afterwards rather than only caught: a
+     * try/catch on its own never sees the common case, and the consumer would 
go on to the next poll as though the
+     * exchange had been delivered.
+     * <p/>
+     * The route's own error handler has already logged the exhausted failure 
by this point, so this is not the only
+     * record of it; what it adds is that the consumer no longer treats a 
failed exchange as a delivered one. It matters
+     * most with {@code consumerProcessedStrategy=delete}, where the document 
is removed during the poll, before the
+     * route runs, so a failure means the document is gone.
+     */
+    private void processExchange(Exchange exchange) {
+        try {
+            getProcessor().process(exchange);
+        } catch (Exception e) {
+            exchange.setException(e);
+        }
+        Exception cause = exchange.getException();
+        if (cause != null) {
+            getExceptionHandler().handleException("Error processing exchange", 
exchange, cause);
+        }
+    }
+
     private void logDetails(String id, Object doc, String key, String 
designDocumentName, String viewName, Exchange exchange) {
         if (LOG.isTraceEnabled()) {
             LOG.trace("Created exchange = {}", exchange);
diff --git 
a/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseEndpoint.java
 
b/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseEndpoint.java
index 801b8427036c..5ca1a7f9b9e3 100644
--- 
a/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseEndpoint.java
+++ 
b/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseEndpoint.java
@@ -25,6 +25,8 @@ import java.util.LinkedHashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
+import java.util.concurrent.locks.Lock;
+import java.util.concurrent.locks.ReentrantLock;
 import java.util.stream.Collectors;
 
 import com.couchbase.client.java.Bucket;
@@ -164,6 +166,14 @@ public class CouchbaseEndpoint extends 
ScheduledPollEndpoint implements Endpoint
     @UriParam(label = "advanced", defaultValue = "30000", javaType = 
"java.time.Duration")
     private long connectTimeout = DEFAULT_CONNECT_TIMEOUT;
 
+    /**
+     * Guards the lazily created {@link #cluster}, so that a producer and a 
consumer starting concurrently on the same
+     * endpoint share one connection instead of racing to open two.
+     */
+    private final Lock clusterLock = new ReentrantLock();
+    private Cluster cluster;
+    private ClusterEnvironment clusterEnvironment;
+
     public CouchbaseEndpoint() {
     }
 
@@ -668,30 +678,66 @@ public class CouchbaseEndpoint extends 
ScheduledPollEndpoint implements Endpoint
         return uriArray;
     }
 
-    //create from couchbase-client
-    private Bucket createClient() throws Exception {
+    /**
+     * The bucket on this endpoint's cluster, creating the cluster on first 
use.
+     * <p/>
+     * Package-private because the consumer and the producer take their 
handles again on every start: this endpoint
+     * disconnects the cluster when it stops, so handles cached across a 
restart would point at a dead cluster.
+     */
+    Bucket createClient() throws Exception {
         if (bucket == null || bucket.isEmpty()) {
             throw new CamelException(COUCHBASE_URI_ERROR);
         }
 
-        ClusterEnvironment env = createClusterEnvironment();
-
-        String connStr;
-        if (connectionString != null && !connectionString.isEmpty()) {
-            connStr = connectionString;
-        } else {
-            List<URI> hosts = Arrays.asList(makeBootstrapURI());
-            String addHosts = hosts.stream()
-                    .map(URI::getHost)
-                    .collect(Collectors.joining(","));
-            connStr = !addHosts.isEmpty() ? addHosts : hostname;
+        clusterLock.lock();
+        try {
+            if (cluster == null) {
+                clusterEnvironment = createClusterEnvironment();
+
+                String connStr;
+                if (connectionString != null && !connectionString.isEmpty()) {
+                    connStr = connectionString;
+                } else {
+                    List<URI> hosts = Arrays.asList(makeBootstrapURI());
+                    String addHosts = hosts.stream()
+                            .map(URI::getHost)
+                            .collect(Collectors.joining(","));
+                    connStr = !addHosts.isEmpty() ? addHosts : hostname;
+                }
+
+                cluster = Cluster.connect(connStr, ClusterOptions
+                        .clusterOptions(username, password)
+                        .environment(clusterEnvironment));
+            }
+            return cluster.bucket(bucket);
+        } finally {
+            clusterLock.unlock();
         }
+    }
 
-        Cluster cluster = Cluster.connect(connStr, ClusterOptions
-                .clusterOptions(username, password)
-                .environment(env));
-
-        return cluster.bucket(bucket);
+    @Override
+    protected void doStop() throws Exception {
+        super.doStop();
+
+        clusterLock.lock();
+        try {
+            if (cluster != null) {
+                // disconnect() is the blocking close. The consumer and the 
producer used to call
+                // bucket.core().shutdown() instead, which returns a cold Mono 
nobody subscribed to and so
+                // closed nothing at all
+                cluster.disconnect();
+                cluster = null;
+            }
+            if (clusterEnvironment != null) {
+                // the environment is built here and handed to the SDK, which 
records it as *external* and
+                // therefore never shuts it down on disconnect - its event 
loops and schedulers outlive the
+                // cluster unless they are stopped explicitly
+                clusterEnvironment.shutdown();
+                clusterEnvironment = null;
+            }
+        } finally {
+            clusterLock.unlock();
+        }
     }
 
     /**
@@ -709,10 +755,12 @@ public class CouchbaseEndpoint extends 
ScheduledPollEndpoint implements Endpoint
     ClusterEnvironment createClusterEnvironment() {
         ClusterEnvironment.Builder cfb = ClusterEnvironment.builder();
         cfb.jsonSerializer(DefaultJsonSerializer.create());
+        // connectTimeout is always applied, so that the documented default of 
30s holds rather than the SDK's
+        // own 10s. queryTimeout keeps its guard on purpose: the documented 
2500ms default is far shorter than
+        // the SDK's 75s, and applying it unconditionally would cut short 
every query that is slower than that
+        cfb.timeoutConfig().connectTimeout(Duration.ofMillis(connectTimeout));
         if (queryTimeout != DEFAULT_QUERY_TIMEOUT) {
-            cfb.timeoutConfig()
-                    .connectTimeout(Duration.ofMillis(connectTimeout))
-                    .queryTimeout(Duration.ofMillis(queryTimeout));
+            cfb.timeoutConfig().queryTimeout(Duration.ofMillis(queryTimeout));
         }
         return cfb.build();
     }
diff --git 
a/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseProducer.java
 
b/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseProducer.java
index 11a65b576b87..fbd1521a3fa4 100644
--- 
a/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseProducer.java
+++ 
b/components/camel-couchbase/src/main/java/org/apache/camel/component/couchbase/CouchbaseProducer.java
@@ -47,8 +47,7 @@ public class CouchbaseProducer extends DefaultProducer {
 
     private final AtomicLong startId = new AtomicLong();
     private final CouchbaseEndpoint endpoint;
-    private final Bucket client;
-    private final Collection collection;
+    private Collection collection;
     private final PersistTo persistTo;
     private final ReplicateTo replicateTo;
     private final int producerRetryPause;
@@ -58,20 +57,7 @@ public class CouchbaseProducer extends DefaultProducer {
     public CouchbaseProducer(CouchbaseEndpoint endpoint, Bucket client, int 
persistTo, int replicateTo) {
         super(endpoint);
         this.endpoint = endpoint;
-        this.client = client;
-        Scope scope;
-
-        if (endpoint.getScope() != null) {
-            scope = client.scope(endpoint.getScope());
-        } else {
-            scope = client.defaultScope();
-        }
-
-        if (endpoint.getCollection() != null) {
-            this.collection = scope.collection(endpoint.getCollection());
-        } else {
-            this.collection = client.defaultCollection();
-        }
+        this.collection = resolveCollection(client);
 
         if (endpoint.isAutoStartIdForInserts()) {
             this.startId.set(endpoint.getStartingIdForInsertsFrom());
@@ -89,6 +75,9 @@ public class CouchbaseProducer extends DefaultProducer {
             case 1:
                 this.persistTo = PersistTo.ACTIVE;
                 break;
+            case 2:
+                this.persistTo = PersistTo.TWO;
+                break;
             case 3:
                 this.persistTo = PersistTo.THREE;
                 break;
@@ -120,6 +109,28 @@ public class CouchbaseProducer extends DefaultProducer {
 
     }
 
+    private Collection resolveCollection(Bucket client) {
+        Scope scope;
+        if (endpoint.getScope() != null) {
+            scope = client.scope(endpoint.getScope());
+        } else {
+            scope = client.defaultScope();
+        }
+
+        if (endpoint.getCollection() != null) {
+            return scope.collection(endpoint.getCollection());
+        }
+        return client.defaultCollection();
+    }
+
+    @Override
+    protected void doStart() throws Exception {
+        super.doStart();
+        // Take the collection again on every start. The endpoint owns the 
cluster and disconnects it when it
+        // stops, so a handle kept from construction would point at a dead 
cluster after a restart in place
+        this.collection = resolveCollection(endpoint.createClient());
+    }
+
     @Override
     public void process(Exchange exchange) throws Exception {
         Map<String, Object> headers = exchange.getIn().getHeaders();
@@ -155,12 +166,4 @@ public class CouchbaseProducer extends DefaultProducer {
         exchange.getIn().removeHeader(HEADER_ID);
     }
 
-    @Override
-    protected void doShutdown() throws Exception {
-        super.doShutdown();
-        if (client != null) {
-            client.core().shutdown();
-        }
-    }
-
 }
diff --git 
a/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseConsumerTest.java
 
b/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseConsumerTest.java
new file mode 100644
index 000000000000..bf3ffe78e972
--- /dev/null
+++ 
b/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseConsumerTest.java
@@ -0,0 +1,163 @@
+/*
+ * 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.couchbase;
+
+import java.util.ArrayDeque;
+import java.util.Queue;
+import java.util.concurrent.atomic.AtomicReference;
+
+import com.couchbase.client.java.Bucket;
+import com.couchbase.client.java.Cluster;
+import com.couchbase.client.java.ClusterOptions;
+import com.couchbase.client.java.Collection;
+import com.couchbase.client.java.Scope;
+import org.apache.camel.Exchange;
+import org.apache.camel.Processor;
+import org.apache.camel.impl.DefaultCamelContext;
+import org.apache.camel.spi.ExceptionHandler;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
+import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.when;
+
+class CouchbaseConsumerTest {
+
+    private DefaultCamelContext context;
+    private CouchbaseEndpoint endpoint;
+    private CouchbaseConsumer consumer;
+    private MockedStatic<Cluster> clusters;
+
+    @BeforeEach
+    void setUp() throws Exception {
+        context = new DefaultCamelContext();
+        endpoint = new CouchbaseEndpoint(
+                "couchbase:http://localhost:8091";, "http://localhost:8091";,
+                new CouchbaseComponent(context));
+        endpoint.setBucket("bucket");
+        endpoint.setUsername("user");
+        endpoint.setPassword("secret");
+        context.start();
+
+        // the consumer takes its handles from the endpoint on every start, so 
the endpoint has to hand out a
+        // bucket without a server behind it
+        Bucket bucket = mock(Bucket.class);
+        Scope scope = mock(Scope.class);
+        when(bucket.defaultScope()).thenReturn(scope);
+        when(bucket.defaultCollection()).thenReturn(mock(Collection.class));
+        Cluster cluster = mock(Cluster.class);
+        when(cluster.bucket(anyString())).thenReturn(bucket);
+        clusters = mockStatic(Cluster.class);
+        clusters.when(() -> Cluster.connect(anyString(), 
any(ClusterOptions.class))).thenReturn(cluster);
+    }
+
+    @AfterEach
+    void tearDown() throws Exception {
+        if (consumer != null) {
+            consumer.stop();
+        }
+        context.stop();
+        if (clusters != null) {
+            clusters.close();
+        }
+    }
+
+    /**
+     * {@code isBatchAllowed()} is false until the consumer is started, so the 
batch has to run against a started
+     * consumer to exercise anything at all. The initial delay keeps the 
scheduler from ever polling the mocked bucket.
+     */
+    private CouchbaseConsumer startedConsumer(Processor processor) throws 
Exception {
+        consumer = new CouchbaseConsumer(endpoint, endpoint.createClient(), 
processor);
+        consumer.setInitialDelay(Long.MAX_VALUE / 2);
+        consumer.start();
+        return consumer;
+    }
+
+    /**
+     * A failing route does not throw out of {@code process()} - the failure 
is left on the exchange - so a consumer
+     * that only catches never learns about the common case. With {@code 
consumerProcessedStrategy=delete} the document
+     * is already gone by then, so an unreported failure loses the message 
outright.
+     */
+    @Test
+    void aRouteFailureIsReportedToTheExceptionHandler() throws Exception {
+        Exception failure = new IllegalStateException("the route blew up");
+        CouchbaseConsumer consumer = startedConsumer(ex -> 
ex.setException(failure));
+
+        AtomicReference<Exception> reported = new AtomicReference<>();
+        consumer.setExceptionHandler(new CapturingExceptionHandler(reported));
+
+        Queue<Object> exchanges = new ArrayDeque<>();
+        exchanges.add(endpoint.createExchange());
+
+        consumer.processBatch(exchanges);
+
+        assertNotNull(reported.get(), "the route failure should have been 
handed to the exception handler");
+        assertEquals(failure, reported.get());
+    }
+
+    /**
+     * A poll that produces more exchanges than the batch hands to the route 
must not simply drop the remainder: nothing
+     * else releases them, and a pooled exchange that is never released never 
returns to the pool.
+     */
+    @Test
+    void exchangesBeyondTheBatchAreReleasedRatherThanDropped() throws 
Exception {
+        CouchbaseConsumer consumer = startedConsumer(ex -> {
+        });
+        consumer.setMaxMessagesPerPoll(1);
+
+        Queue<Object> exchanges = new ArrayDeque<>();
+        for (int i = 0; i < 3; i++) {
+            exchanges.add(endpoint.createExchange());
+        }
+
+        consumer.processBatch(exchanges);
+
+        assertTrue(exchanges.isEmpty(), "the exchanges not handed to the route 
should have been released, not left behind");
+    }
+
+    private static final class CapturingExceptionHandler implements 
ExceptionHandler {
+
+        private final AtomicReference<Exception> captured;
+
+        private CapturingExceptionHandler(AtomicReference<Exception> captured) 
{
+            this.captured = captured;
+        }
+
+        @Override
+        public void handleException(Throwable exception) {
+            captured.set((Exception) exception);
+        }
+
+        @Override
+        public void handleException(String message, Throwable exception) {
+            captured.set((Exception) exception);
+        }
+
+        @Override
+        public void handleException(String message, Exchange exchange, 
Throwable exception) {
+            captured.set((Exception) exception);
+        }
+    }
+}
diff --git 
a/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseEndpointTest.java
 
b/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseEndpointTest.java
index 6d6f6ee8fd62..6f5716b670c2 100644
--- 
a/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseEndpointTest.java
+++ 
b/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseEndpointTest.java
@@ -16,20 +16,39 @@
  */
 package org.apache.camel.component.couchbase;
 
+import java.time.Duration;
 import java.util.HashMap;
 import java.util.Map;
 
+import com.couchbase.client.java.Bucket;
+import com.couchbase.client.java.Cluster;
+import com.couchbase.client.java.ClusterOptions;
+import com.couchbase.client.java.Collection;
+import com.couchbase.client.java.Scope;
 import com.couchbase.client.java.codec.DefaultJsonSerializer;
 import com.couchbase.client.java.codec.JsonSerializer;
 import com.couchbase.client.java.env.ClusterEnvironment;
+import org.apache.camel.CamelContext;
 import org.apache.camel.Processor;
+import org.apache.camel.Producer;
+import org.apache.camel.impl.DefaultCamelContext;
 import org.junit.jupiter.api.Test;
+import org.mockito.MockedStatic;
 
 import static 
org.apache.camel.component.couchbase.CouchbaseConstants.DEFAULT_COUCHBASE_PORT;
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.atLeastOnce;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.mockStatic;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
 
 public class CouchbaseEndpointTest {
 
@@ -271,4 +290,125 @@ public class CouchbaseEndpointTest {
             env.shutdown();
         }
     }
+
+    /**
+     * connectTimeout used to sit inside a guard that tested queryTimeout, so 
setting it on its own did nothing and the
+     * documented 30s default never applied - the SDK's own 10s did.
+     */
+    @Test
+    void connectTimeoutIsAppliedOnItsOwn() {
+        CouchbaseEndpoint endpoint = new CouchbaseEndpoint();
+        endpoint.setConnectTimeout(1234);
+        ClusterEnvironment env = endpoint.createClusterEnvironment();
+        try {
+            assertEquals(Duration.ofMillis(1234), 
env.timeoutConfig().connectTimeout());
+        } finally {
+            env.shutdown();
+        }
+    }
+
+    @Test
+    void connectTimeoutDefaultIsApplied() {
+        CouchbaseEndpoint endpoint = new CouchbaseEndpoint();
+        ClusterEnvironment env = endpoint.createClusterEnvironment();
+        try {
+            
assertEquals(Duration.ofMillis(CouchbaseConstants.DEFAULT_CONNECT_TIMEOUT),
+                    env.timeoutConfig().connectTimeout());
+        } finally {
+            env.shutdown();
+        }
+    }
+
+    @Test
+    void queryTimeoutIsAppliedWhenSet() {
+        CouchbaseEndpoint endpoint = new CouchbaseEndpoint();
+        endpoint.setQueryTimeout(9999);
+        ClusterEnvironment env = endpoint.createClusterEnvironment();
+        try {
+            assertEquals(Duration.ofMillis(9999), 
env.timeoutConfig().queryTimeout());
+        } finally {
+            env.shutdown();
+        }
+    }
+
+    /**
+     * The connection used to be opened once per producer and once per 
consumer and never closed at all: both sides
+     * called {@code bucket.core().shutdown()}, which returns a cold Mono 
nobody subscribed to.
+     */
+    @Test
+    void oneConnectionPerEndpointAndItIsClosedOnStop() throws Exception {
+        Cluster cluster = mock(Cluster.class);
+        Bucket bucket = mock(Bucket.class);
+        Scope scope = mock(Scope.class);
+        when(cluster.bucket(anyString())).thenReturn(bucket);
+        when(bucket.defaultScope()).thenReturn(scope);
+        when(bucket.defaultCollection()).thenReturn(mock(Collection.class));
+
+        try (MockedStatic<Cluster> clusters = mockStatic(Cluster.class)) {
+            clusters.when(() -> Cluster.connect(anyString(), 
any(ClusterOptions.class))).thenReturn(cluster);
+
+            CamelContext context = new DefaultCamelContext();
+            CouchbaseEndpoint endpoint = new CouchbaseEndpoint(
+                    "couchbase:http://localhost:8091";,
+                    "http://localhost:8091";, new CouchbaseComponent(context));
+            endpoint.setBucket("bucket");
+            endpoint.setUsername("user");
+            endpoint.setPassword("secret");
+            endpoint.start();
+
+            endpoint.createProducer();
+            endpoint.createProducer();
+
+            clusters.verify(() -> Cluster.connect(anyString(), 
any(ClusterOptions.class)), times(1));
+            verify(cluster, never()).disconnect();
+
+            endpoint.stop();
+            verify(cluster).disconnect();
+        }
+    }
+
+    /**
+     * The endpoint disconnects its cluster on stop, so a consumer or producer 
that kept a handle from before the
+     * restart would be talking to a dead cluster. Both take their handles 
again on start.
+     */
+    @Test
+    void aRestartedProducerTakesItsCollectionFromTheNewConnection() throws 
Exception {
+        Bucket first = mock(Bucket.class);
+        Bucket second = mock(Bucket.class);
+        when(first.defaultCollection()).thenReturn(mock(Collection.class));
+        when(second.defaultCollection()).thenReturn(mock(Collection.class));
+
+        Cluster one = mock(Cluster.class);
+        Cluster two = mock(Cluster.class);
+        when(one.bucket(anyString())).thenReturn(first);
+        when(two.bucket(anyString())).thenReturn(second);
+
+        try (MockedStatic<Cluster> clusters = mockStatic(Cluster.class)) {
+            clusters.when(() -> Cluster.connect(anyString(), 
any(ClusterOptions.class))).thenReturn(one, two);
+
+            CamelContext context = new DefaultCamelContext();
+            CouchbaseEndpoint endpoint = new CouchbaseEndpoint(
+                    "couchbase:http://localhost:8091";,
+                    "http://localhost:8091";, new CouchbaseComponent(context));
+            endpoint.setBucket("bucket");
+            endpoint.setUsername("user");
+            endpoint.setPassword("secret");
+
+            endpoint.start();
+            Producer producer = endpoint.createProducer();
+            producer.start();
+            // resolved twice against the first connection: once when 
constructed, once on start
+            verify(first, atLeastOnce()).defaultCollection();
+
+            endpoint.stop();
+            verify(one).disconnect();
+
+            // restart in place: the producer object survives, the connection 
behind it does not
+            endpoint.start();
+            producer.stop();
+            producer.start();
+
+            verify(second).defaultCollection();
+        }
+    }
 }
diff --git 
a/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseProducerTest.java
 
b/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseProducerTest.java
index 10c36953751b..7b79fd8cd181 100644
--- 
a/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseProducerTest.java
+++ 
b/components/camel-couchbase/src/test/java/org/apache/camel/component/couchbase/CouchbaseProducerTest.java
@@ -140,4 +140,17 @@ public class CouchbaseProducerTest {
 
         verify(collection).upsert(anyString(), any(), options.capture());
     }
+
+    /**
+     * PersistTo.TWO exists in the SDK and 2 sits inside the range the failure 
message advertises, but it was the one
+     * value in 0..4 the switch did not map.
+     */
+    @Test
+    void everyPersistToValueInTheAdvertisedRangeIsAccepted() {
+        for (int persistTo = 0; persistTo <= 4; persistTo++) {
+            int value = persistTo;
+            assertDoesNotThrow(() -> new CouchbaseProducer(endpoint, client, 
value, 0),
+                    "persistTo=" + value + " is within the documented range 
and should be accepted");
+        }
+    }
 }
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 181ec80ab258..ca4f926e4b21 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
@@ -528,6 +528,27 @@ A header that carries a literal `${...}` string is now 
used as the object key ve
 A configured `keyName` or `bucketName` whose Simple expression resolves to 
`null` now fails with an
 `IllegalArgumentException` at the producer, instead of passing a `null` 
key/bucket on to the AWS SDK.
 
+=== camel-couchbase
+
+The `connectTimeout` option is now applied. It previously sat inside a 
condition that tested `queryTimeout`, so
+setting it on its own had no effect and the documented default of 30000 ms 
never reached the SDK - connections used
+the Couchbase SDK's own default of 10 seconds instead. Deployments that relied 
on that 10 second behaviour without
+setting the option should now set `connectTimeout=10000` explicitly.
+
+`queryTimeout` is unchanged: it is still only applied when set to something 
other than its default, because the
+documented 2500 ms default is much shorter than the SDK's 75 seconds and 
applying it unconditionally would cut short
+queries that currently succeed.
+
+The consumer now reports a failed exchange through its exception handler, 
which logs a warning naming the endpoint,
+as `TimerConsumer` does. Previously the consumer discarded the outcome of 
routing entirely and went on to the next
+poll as if the exchange had been delivered.
+
+This is not a change to error handling. An exchange that failed while being 
routed has already passed through the
+route's error handler, which logs the exhausted failure, and 
`BridgeExceptionHandlerToErrorHandler` deliberately
+falls back to its logging handler for such an exchange rather than bridging it 
a second time. What changes is that
+the consumer itself no longer stays silent, and that with 
`consumerProcessedStrategy=delete` the loss is now visible:
+the document is removed during the poll, before the route runs.
+
 === camel-console
 
 The `context` developer console no longer counts routes created by Kamelets in 
its `routesTotal` and `routesStarted`

Reply via email to