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`