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

github-merge-queue[bot] pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new 09941a1781 [Improve][Connector-V2] Make Couchbase sink readiness 
timeout configurable (#12168)
09941a1781 is described below

commit 09941a178186ca004eb55af377b9757d8dab1871
Author: Goutam Adwant <[email protected]>
AuthorDate: Fri Sep 11 02:49:45 2026 +0000

    [Improve][Connector-V2] Make Couchbase sink readiness timeout configurable 
(#12168)
    
    Signed-off-by: Goutam Adwant <[email protected]>
---
 docs/en/connectors/sink/Couchbase.md               |  17 +++
 docs/zh/connectors/sink/Couchbase.md               |  14 ++
 .../couchbase/config/CouchbaseSinkOptions.java     |   9 ++
 .../couchbase/sink/CouchbaseSinkFactory.java       |   5 +
 .../seatunnel/couchbase/sink/CouchbaseWriter.java  |   7 +-
 .../couchbase/sink/CouchbaseWriterOptions.java     |  25 +++
 .../couchbase/sink/CouchbaseSinkFactoryTest.java   | 157 +++++++++++++++++++
 .../sink/CouchbaseWriterConstructorLeakTest.java   |  39 +++++
 .../couchbase/sink/CouchbaseWriterOptionsTest.java |  77 ++++++++++
 .../connector/couchbase/CouchbaseReadinessIT.java  | 167 +++++++++++++++++++++
 10 files changed, 516 insertions(+), 1 deletion(-)

diff --git a/docs/en/connectors/sink/Couchbase.md 
b/docs/en/connectors/sink/Couchbase.md
index 1304718d51..57767d6f7c 100644
--- a/docs/en/connectors/sink/Couchbase.md
+++ b/docs/en/connectors/sink/Couchbase.md
@@ -78,12 +78,29 @@ Couchbase stores JSON documents. The connector maps 
SeaTunnel types to JSON valu
 | bucket                | String        | Yes      | -          | Target 
bucket name. |
 | scope                 | String        | No       | `_default` | Target scope 
name within the bucket. |
 | collection            | String        | Yes      | -          | Target 
collection name. |
+| ready.timeout         | Integer       | No       | `30`       | Maximum 
seconds to wait for the target bucket to become ready during writer 
initialization. Must be greater than zero. |
 | primary-key           | `List<String>` | No       | -          | Field names 
used to build the document key (length-prefixed encoding: `<len>:<value>` 
components separated by `#`). A random UUID is used when not set. |
 | upsert-enable         | Boolean       | No       | `false`    | Enable 
upsert (insert-or-replace) mode. When `false`, duplicate keys will cause an 
error. |
 | buffer-flush.max-rows | Integer       | No       | `1000`     | Maximum rows 
to buffer before a batch write is triggered. Use `-1` to disable. |
 | retry.max             | Integer       | No       | `3`        | Maximum 
retry attempts on transient write failure. |
 | retry.interval        | Long          | No       | `1000`     | Base 
milliseconds for linear retry delay. Attempt `n` waits `retry.interval × n` ms. 
|
 
+### Startup readiness
+
+`ready.timeout` controls the bucket-readiness wait during writer 
initialization. The default
+remains 30 seconds. For a cluster that needs more time to become available, 
set a larger positive
+value, for example `ready.timeout = 60`.
+
+The value is in seconds, not milliseconds. No connector-specific upper limit 
is enforced;
+choose the smallest budget that covers the cluster's observed recovery time. 
Excessively large
+values can delay writer-initialization failure when the bucket remains 
unavailable.
+
+The Couchbase SDK handles connection attempts within this wait; the connector 
does not add an
+outer bootstrap retry loop. `retry.max` and `retry.interval` still apply only 
to writes. This
+option does not change individual SDK operation timeouts or the engine's 
job-startup timeout.
+An expired readiness wait still fails writer initialization and disconnects 
the client.
+Increasing the budget does not correct invalid credentials, incorrect 
addresses or a missing bucket.
+
 ## Security
 
 ### TLS / encrypted transport
diff --git a/docs/zh/connectors/sink/Couchbase.md 
b/docs/zh/connectors/sink/Couchbase.md
index 43883e9544..b19e806395 100644
--- a/docs/zh/connectors/sink/Couchbase.md
+++ b/docs/zh/connectors/sink/Couchbase.md
@@ -73,12 +73,26 @@ sh bin/install-plugin.sh ${version}
 | bucket                 | String         | 是       | -          | 目标 Bucket 
名称。 |
 | scope                  | String         | 否       | `_default` | Bucket 中的目标 
Scope 名称。 |
 | collection             | String         | 是       | -          | 目标 
Collection 名称。 |
+| ready.timeout          | Integer        | 否       | `30`       | 写入器初始化时等待目标 
bucket 就绪的最长时间(秒),必须大于零。 |
 | primary-key            | `List<String>`  | 否       | -          | 
用于构建文档键的字段名列表(长度前缀编码:`<长度>:<值>` 分量以 `#` 分隔)。未设置时使用随机 UUID。 |
 | upsert-enable          | Boolean        | 否       | `false`    | 是否启用 
Upsert(插入或替换)模式。为 `false` 时,重复键将报错。 |
 | buffer-flush.max-rows  | Integer        | 否       | `1000`     | 
触发批量写入的最大缓冲行数。设为 `-1` 禁用。 |
 | retry.max              | Integer        | 否       | `3`        | 
写入失败时的最大重试次数。 |
 | retry.interval         | Long           | 否       | `1000`     | 
线性退避基础间隔(毫秒)。第 n 次重试等待 `retry.interval × n` 毫秒。 |
 
+### 启动就绪等待
+
+`ready.timeout` 控制写入器初始化时等待目标 bucket 就绪的时间,默认仍为 30 秒。
+对于需要更长时间才能恢复可用的集群,可配置更大的正数,例如 `ready.timeout = 60`。
+
+该值的单位是秒,而不是毫秒。连接器不额外限制上限,应根据集群实际恢复时间选择足够的最小值。
+当 bucket 持续不可用时,过大的值会延迟写入器初始化失败的报告。
+
+Couchbase SDK 在此等待期间处理连接尝试,连接器不会额外添加启动重试循环。
+`retry.max` 和 `retry.interval` 仍仅用于写入重试。此选项不会修改 SDK 单次操作的超时时间,
+也不会修改引擎的作业启动超时时间。等待超时后,写入器初始化仍会失败并断开客户端连接。
+增加等待时间不能修复无效凭据、错误地址或不存在的 bucket。
+
 ## 安全性
 
 ### TLS / 加密传输
diff --git 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
index 4d09a802c5..9596c32dad 100644
--- 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
+++ 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/config/CouchbaseSinkOptions.java
@@ -25,6 +25,15 @@ import java.util.List;
 /** Configuration options specific to the Couchbase sink connector. */
 public class CouchbaseSinkOptions extends CouchbaseConfig {
 
+    /** Maximum time to wait for bucket readiness during writer 
initialization, in seconds. */
+    public static final Option<Integer> READY_TIMEOUT =
+            Options.key("ready.timeout")
+                    .intType()
+                    .defaultValue(30)
+                    .withDescription(
+                            "The timeout in seconds for waiting until the 
target bucket is ready"
+                                    + " during writer initialization. Must be 
greater than zero.");
+
     /**
      * Maximum number of rows buffered before a batch write is triggered.
      *
diff --git 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
index 630580311e..a4ebbc265b 100644
--- 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
+++ 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactory.java
@@ -18,6 +18,7 @@
 package org.apache.seatunnel.connectors.seatunnel.couchbase.sink;
 
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.Conditions;
 import org.apache.seatunnel.api.configuration.util.OptionRule;
 import org.apache.seatunnel.api.table.catalog.CatalogTable;
 import org.apache.seatunnel.api.table.catalog.TableIdentifier;
@@ -60,6 +61,9 @@ public class CouchbaseSinkFactory implements TableSinkFactory 
{
                         CouchbaseSinkOptions.RETRY_INTERVAL,
                         CouchbaseSinkOptions.UPSERT_ENABLE,
                         CouchbaseSinkOptions.PRIMARY_KEY)
+                .optional(
+                        CouchbaseSinkOptions.READY_TIMEOUT,
+                        
Conditions.greaterThan(CouchbaseSinkOptions.READY_TIMEOUT, 0))
                 .build();
     }
 
@@ -88,6 +92,7 @@ public class CouchbaseSinkFactory implements TableSinkFactory 
{
                         
.withUsername(config.get(CouchbaseSinkOptions.USERNAME))
                         
.withPassword(config.get(CouchbaseSinkOptions.PASSWORD))
                         .withBucket(config.get(CouchbaseSinkOptions.BUCKET))
+                        
.withReadyTimeout(config.get(CouchbaseSinkOptions.READY_TIMEOUT))
                         .withScope(config.get(CouchbaseSinkOptions.SCOPE))
                         
.withCollection(config.get(CouchbaseSinkOptions.COLLECTION));
 
diff --git 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
index 5da720e4ad..7b4b4e9c48 100644
--- 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
+++ 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriter.java
@@ -115,7 +115,12 @@ public class CouchbaseWriter implements 
SinkWriter<SeaTunnelRow, Void, Void> {
                         options.getPassword());
         Collection resolvedCollection;
         try {
-            
connectedCluster.bucket(options.getBucket()).waitUntilReady(Duration.ofSeconds(30));
+            log.debug(
+                    "Waiting up to {} seconds for Couchbase bucket readiness",
+                    options.getReadyTimeout());
+            connectedCluster
+                    .bucket(options.getBucket())
+                    
.waitUntilReady(Duration.ofSeconds(options.getReadyTimeout()));
             resolvedCollection =
                     connectedCluster
                             .bucket(options.getBucket())
diff --git 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
index 324176d276..fe75caf449 100644
--- 
a/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
+++ 
b/seatunnel-connectors-v2/connector-couchbase/src/main/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptions.java
@@ -17,6 +17,8 @@
 
 package org.apache.seatunnel.connectors.seatunnel.couchbase.sink;
 
+import 
org.apache.seatunnel.connectors.seatunnel.couchbase.config.CouchbaseSinkOptions;
+
 import lombok.Getter;
 
 import java.io.Serializable;
@@ -42,6 +44,7 @@ public class CouchbaseWriterOptions implements Serializable {
     private final String[] primaryKey;
     private final int retryMax;
     private final long retryInterval;
+    private final int readyTimeout;
 
     private CouchbaseWriterOptions(Builder builder) {
         this.connectionString = builder.connectionString;
@@ -55,12 +58,18 @@ public class CouchbaseWriterOptions implements Serializable 
{
         this.primaryKey = builder.primaryKey;
         this.retryMax = builder.retryMax;
         this.retryInterval = builder.retryInterval;
+        this.readyTimeout = builder.readyTimeout;
     }
 
     public static Builder builder() {
         return new Builder();
     }
 
+    /** Retains the previous readiness budget for options serialized before 
this field existed. */
+    public int getReadyTimeout() {
+        return readyTimeout == 0 ? 
CouchbaseSinkOptions.READY_TIMEOUT.defaultValue() : readyTimeout;
+    }
+
     /** Fluent builder for {@link CouchbaseWriterOptions}. */
     public static class Builder {
         private String connectionString;
@@ -74,6 +83,7 @@ public class CouchbaseWriterOptions implements Serializable {
         private String[] primaryKey = new String[0];
         private int retryMax = 3;
         private long retryInterval = 1000L;
+        private int readyTimeout = 
CouchbaseSinkOptions.READY_TIMEOUT.defaultValue();
 
         public Builder withConnectionString(String connectionString) {
             this.connectionString = connectionString;
@@ -130,6 +140,21 @@ public class CouchbaseWriterOptions implements 
Serializable {
             return this;
         }
 
+        /**
+         * Sets the bucket-readiness budget used during writer initialization.
+         *
+         * @param readyTimeout positive readiness timeout in seconds
+         * @return this builder
+         * @throws IllegalArgumentException if the timeout is zero or negative
+         */
+        public Builder withReadyTimeout(int readyTimeout) {
+            if (readyTimeout <= 0) {
+                throw new IllegalArgumentException("'ready.timeout' must be 
greater than zero.");
+            }
+            this.readyTimeout = readyTimeout;
+            return this;
+        }
+
         public CouchbaseWriterOptions build() {
             return new CouchbaseWriterOptions(this);
         }
diff --git 
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactoryTest.java
 
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactoryTest.java
new file mode 100644
index 0000000000..e0547924e8
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseSinkFactoryTest.java
@@ -0,0 +1,157 @@
+/*
+ * 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.seatunnel.connectors.seatunnel.couchbase.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.TableIdentifier;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+import org.mockito.MockedStatic;
+import org.mockito.Mockito;
+
+import com.couchbase.client.java.Bucket;
+import com.couchbase.client.java.Cluster;
+import com.couchbase.client.java.Collection;
+import com.couchbase.client.java.Scope;
+
+import java.time.Duration;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class CouchbaseSinkFactoryTest {
+
+    @Test
+    void testConfiguredReadinessTimeoutReachesWriter() throws Exception {
+        Map<String, Object> config = baseConfig();
+        config.put("ready.timeout", 60);
+        verifyReadinessTimeout(config, Duration.ofSeconds(60));
+    }
+
+    @Test
+    void testDefaultReadinessTimeoutRemainsThirtySeconds() throws Exception {
+        verifyReadinessTimeout(baseConfig(), Duration.ofSeconds(30));
+    }
+
+    @ParameterizedTest
+    @ValueSource(ints = {0, -1})
+    void testNonPositiveReadinessTimeoutRejectedBeforeConnecting(int timeout) {
+        Map<String, Object> config = baseConfig();
+        config.put("ready.timeout", timeout);
+        try (MockedStatic<Cluster> staticCluster = 
Mockito.mockStatic(Cluster.class)) {
+            OptionValidationException error =
+                    assertThrows(OptionValidationException.class, () -> 
createSink(config));
+            assertTrue(error.getMessage().contains("ready.timeout"));
+            staticCluster.verifyNoInteractions();
+        }
+    }
+
+    @ParameterizedTest
+    @ValueSource(ints = {0, -1})
+    void testDirectFactoryAlsoRejectsNonPositiveTimeout(int timeout) {
+        Map<String, Object> config = baseConfig();
+        config.put("ready.timeout", timeout);
+        try (MockedStatic<Cluster> staticCluster = 
Mockito.mockStatic(Cluster.class)) {
+            IllegalArgumentException error =
+                    assertThrows(
+                            IllegalArgumentException.class,
+                            () ->
+                                    new CouchbaseSinkFactory()
+                                            .createSink(
+                                                    new 
TableSinkFactoryContext(
+                                                            null,
+                                                            
ReadonlyConfig.fromMap(config),
+                                                            
getClass().getClassLoader())));
+            assertTrue(error.getMessage().contains("ready.timeout"));
+            staticCluster.verifyNoInteractions();
+        }
+    }
+
+    @Test
+    void testWriteRetriesDoNotChangeReadinessTimeout() throws Exception {
+        Map<String, Object> config = baseConfig();
+        config.put("retry.max", 10);
+        config.put("retry.interval", 5000L);
+        verifyReadinessTimeout(config, Duration.ofSeconds(30));
+    }
+
+    private void verifyReadinessTimeout(Map<String, Object> config, Duration 
expected)
+            throws Exception {
+        Cluster cluster = mock(Cluster.class);
+        Bucket bucket = mock(Bucket.class);
+        Scope scope = mock(Scope.class);
+        Collection collection = mock(Collection.class);
+        when(cluster.bucket("test_bucket")).thenReturn(bucket);
+        when(bucket.scope("_default")).thenReturn(scope);
+        when(scope.collection("_default")).thenReturn(collection);
+
+        try (MockedStatic<Cluster> staticCluster = 
Mockito.mockStatic(Cluster.class)) {
+            staticCluster
+                    .when(() -> Cluster.connect("couchbase://localhost", 
"user", "pass"))
+                    .thenReturn(cluster);
+            CouchbaseWriter writer = createSink(config).createWriter(null);
+            try {
+                verify(bucket).waitUntilReady(expected);
+            } finally {
+                writer.close();
+            }
+            verify(cluster).disconnect();
+        }
+    }
+
+    private CouchbaseSink createSink(Map<String, Object> config) {
+        CouchbaseSinkFactory factory = new CouchbaseSinkFactory();
+        ReadonlyConfig options = ReadonlyConfig.fromMap(config);
+        ConfigValidator.of(options).validate(factory.optionRule());
+        CatalogTable table =
+                CatalogTable.of(
+                        TableIdentifier.of("catalog", "database", "table"),
+                        TableSchema.builder().build(),
+                        Collections.emptyMap(),
+                        Collections.emptyList(),
+                        "");
+        return (CouchbaseSink)
+                factory.createSink(
+                                new TableSinkFactoryContext(
+                                        table, options, 
getClass().getClassLoader()))
+                        .createSink();
+    }
+
+    private Map<String, Object> baseConfig() {
+        Map<String, Object> config = new HashMap<>();
+        config.put("connection.string", "couchbase://localhost");
+        config.put("username", "user");
+        config.put("password", "pass");
+        config.put("bucket", "test_bucket");
+        config.put("collection", "_default");
+        return config;
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
 
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
index 7d5ce27ef5..df73d5ac12 100644
--- 
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
+++ 
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterConstructorLeakTest.java
@@ -23,16 +23,22 @@ import 
org.apache.seatunnel.api.table.catalog.TableIdentifier;
 import org.apache.seatunnel.api.table.catalog.TableSchema;
 
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.MethodSource;
 import org.mockito.MockedStatic;
 import org.mockito.Mockito;
 
+import com.couchbase.client.core.error.AuthenticationFailureException;
+import com.couchbase.client.core.error.UnambiguousTimeoutException;
 import com.couchbase.client.java.Bucket;
 import com.couchbase.client.java.Cluster;
 import com.couchbase.client.java.Collection;
 import com.couchbase.client.java.Scope;
 
 import java.time.Duration;
+import java.util.stream.Stream;
 
+import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyString;
@@ -59,6 +65,39 @@ import static org.mockito.Mockito.when;
  */
 class CouchbaseWriterConstructorLeakTest {
 
+    @ParameterizedTest
+    @MethodSource("readinessFailures")
+    void testReadinessFailureIsNotRetriedOrMasked(RuntimeException failure) {
+        Bucket bucket = mock(Bucket.class);
+        doThrow(failure).when(bucket).waitUntilReady(any(Duration.class));
+        Cluster cluster = mockCluster(bucket);
+        RuntimeException disconnectFailure = new RuntimeException("disconnect 
failed");
+        doThrow(disconnectFailure).when(cluster).disconnect();
+
+        try (MockedStatic<Cluster> staticCluster = 
Mockito.mockStatic(Cluster.class)) {
+            staticCluster
+                    .when(() -> Cluster.connect(anyString(), anyString(), 
anyString()))
+                    .thenReturn(cluster);
+            RuntimeException thrown =
+                    assertThrows(
+                            RuntimeException.class,
+                            () ->
+                                    new CouchbaseWriter(
+                                            minimalOptions(), 
minimalCatalogTable(), null));
+            assertSame(failure, thrown);
+            assertSame(disconnectFailure, thrown.getSuppressed()[0]);
+            verify(bucket).waitUntilReady(Duration.ofSeconds(30));
+            verify(cluster).disconnect();
+        }
+    }
+
+    private static Stream<RuntimeException> readinessFailures() {
+        return Stream.of(
+                new UnambiguousTimeoutException("readiness timed out", null),
+                new AuthenticationFailureException("authentication failed", 
null, null),
+                new RuntimeException(new InterruptedException("readiness 
interrupted")));
+    }
+
     /** Minimal {@link CouchbaseWriterOptions} pointing at a fake cluster. */
     private static CouchbaseWriterOptions minimalOptions() {
         return CouchbaseWriterOptions.builder()
diff --git 
a/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptionsTest.java
 
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptionsTest.java
new file mode 100644
index 0000000000..da2d7098a2
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-couchbase/src/test/java/org/apache/seatunnel/connectors/seatunnel/couchbase/sink/CouchbaseWriterOptionsTest.java
@@ -0,0 +1,77 @@
+/*
+ * 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.seatunnel.connectors.seatunnel.couchbase.sink;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
+import java.io.ObjectInputStream;
+import java.io.ObjectOutputStream;
+import java.util.Base64;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class CouchbaseWriterOptionsTest {
+
+    // Serialized default options from before readyTimeout was added 
(serialVersionUID = 1).
+    private static final String LEGACY_OPTIONS =
+            
"rO0ABXNyAE9vcmcuYXBhY2hlLnNlYXR1bm5lbC5jb25uZWN0b3JzLnNlYXR1bm5lbC5jb3VjaGJhc2Uuc2luay5Db3VjaGJhc2VXcml0ZXJPcHRpb25zAAAAAAAAAAECAAtJAAlmbHVzaFNpemVKAA1yZXRyeUludGVydmFsSQAIcmV0cnlNYXhaAAx1cHNlcnRFbmFibGVMAAZidWNrZXR0ABJMamF2YS9sYW5nL1N0cmluZztMAApjb2xsZWN0aW9ucQB+AAFMABBjb25uZWN0aW9uU3RyaW5ncQB+AAFMAAhwYXNzd29yZHEAfgABWwAKcHJpbWFyeUtleXQAE1tMamF2YS9sYW5nL1N0cmluZztMAAVzY29wZXEAfgABTAAIdXNlcm5hbWVxAH4AAXhwAAAD6AAAAAAAAAPoAAAAAwBwcHBwdXIAE1tMamF2YS5sYW5nLlN0cmluZzut0lbn6R17RwI
 [...]
+
+    @Test
+    void testDefaultReadinessTimeout() {
+        assertEquals(30, 
CouchbaseWriterOptions.builder().build().getReadyTimeout());
+    }
+
+    @ParameterizedTest
+    @ValueSource(ints = {0, -1})
+    void testBuilderRejectsNonPositiveReadinessTimeout(int timeout) {
+        IllegalArgumentException error =
+                assertThrows(
+                        IllegalArgumentException.class,
+                        () -> 
CouchbaseWriterOptions.builder().withReadyTimeout(timeout));
+        assertTrue(error.getMessage().contains("ready.timeout"));
+    }
+
+    @Test
+    void testConfiguredReadinessTimeoutSurvivesSerialization() throws 
Exception {
+        ByteArrayOutputStream bytes = new ByteArrayOutputStream();
+        try (ObjectOutputStream output = new ObjectOutputStream(bytes)) {
+            
output.writeObject(CouchbaseWriterOptions.builder().withReadyTimeout(60).build());
+        }
+        assertEquals(60, deserialize(bytes.toByteArray()).getReadyTimeout());
+    }
+
+    @Test
+    void testLegacySerializedOptionsRetainThirtySecondTimeout() throws 
Exception {
+        CouchbaseWriterOptions options = 
deserialize(Base64.getDecoder().decode(LEGACY_OPTIONS));
+        assertEquals(30, options.getReadyTimeout());
+        assertEquals(3, options.getRetryMax());
+        assertEquals("_default", options.getScope());
+    }
+
+    private CouchbaseWriterOptions deserialize(byte[] bytes) throws Exception {
+        try (ObjectInputStream input = new ObjectInputStream(new 
ByteArrayInputStream(bytes))) {
+            return (CouchbaseWriterOptions) input.readObject();
+        }
+    }
+}
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseReadinessIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseReadinessIT.java
new file mode 100644
index 0000000000..372a15855a
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-couchbase-e2e/src/test/java/org/apache/seatunnel/e2e/connector/couchbase/CouchbaseReadinessIT.java
@@ -0,0 +1,167 @@
+/*
+ * 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.seatunnel.e2e.connector.couchbase;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.PhysicalColumn;
+import org.apache.seatunnel.api.table.catalog.TableIdentifier;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import org.apache.seatunnel.api.table.type.BasicType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.connectors.seatunnel.couchbase.sink.CouchbaseSink;
+import 
org.apache.seatunnel.connectors.seatunnel.couchbase.sink.CouchbaseSinkFactory;
+import 
org.apache.seatunnel.connectors.seatunnel.couchbase.sink.CouchbaseWriter;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.api.TestInstance;
+import org.junit.jupiter.api.Timeout;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.couchbase.BucketDefinition;
+import org.testcontainers.couchbase.CouchbaseContainer;
+import org.testcontainers.couchbase.CouchbaseService;
+import org.testcontainers.utility.DockerImageName;
+
+import com.couchbase.client.core.error.UnambiguousTimeoutException;
+import com.couchbase.client.java.Cluster;
+import lombok.extern.slf4j.Slf4j;
+
+import java.time.Duration;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.concurrent.TimeUnit;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+
+/** Factory-level readiness tests; no engine containers are needed for this 
client contract. */
+@TestInstance(TestInstance.Lifecycle.PER_CLASS)
+@Timeout(value = 2, unit = TimeUnit.MINUTES, threadMode = 
Timeout.ThreadMode.SAME_THREAD)
+@Slf4j
+class CouchbaseReadinessIT {
+
+    private final CouchbaseContainer server =
+            new CouchbaseContainer(
+                            DockerImageName.parse(
+                                    System.getProperty(
+                                            "couchbase.test.image",
+                                            
"couchbase/server:community-7.1.1")))
+                    .withNetworkMode("bridge")
+                    .withCredentials("Administrator", "password")
+                    .withEnabledServices(CouchbaseService.KV)
+                    .withBucket(
+                            new BucketDefinition("readiness")
+                                    .withQuota(128)
+                                    .withReplicas(0)
+                                    .withPrimaryIndex(false))
+                    .withStartupTimeout(Duration.ofMinutes(3))
+                    .withStartupAttempts(3)
+                    .withLogConsumer(new 
Slf4jLogConsumer(log).withPrefix("couchbase-readiness"));
+
+    private Cluster verification;
+
+    @BeforeAll
+    void startServer() {
+        server.start();
+        verification =
+                Cluster.connect(
+                        server.getConnectionString(), server.getUsername(), 
server.getPassword());
+        
verification.bucket("readiness").waitUntilReady(Duration.ofSeconds(60));
+    }
+
+    @AfterAll
+    void closeServer() {
+        try {
+            if (verification != null) {
+                verification.disconnect();
+            }
+        } finally {
+            server.stop();
+        }
+    }
+
+    @Test
+    void testReadinessTimeoutAndRecovery() throws Exception {
+        CouchbaseSink unavailableSink = createSink(5, server.getPassword());
+        
server.getDockerClient().pauseContainerCmd(server.getContainerId()).exec();
+        try {
+            assertThrows(
+                    UnambiguousTimeoutException.class, () -> 
unavailableSink.createWriter(null));
+        } finally {
+            
server.getDockerClient().unpauseContainerCmd(server.getContainerId()).exec();
+        }
+
+        CouchbaseWriter writer = createSink(60, 
server.getPassword()).createWriter(null);
+        try {
+            writer.write(new SeaTunnelRow(new Object[] {"recovered"}));
+            writer.prepareCommit();
+            assertEquals(
+                    "recovered",
+                    verification
+                            .bucket("readiness")
+                            .defaultCollection()
+                            .get("9:recovered")
+                            .contentAsObject()
+                            .getString("id"));
+        } finally {
+            writer.close();
+        }
+    }
+
+    @Test
+    void testInvalidCredentialsStillFailReadiness() {
+        CouchbaseSink sink = createSink(5, "incorrect-password");
+        assertThrows(UnambiguousTimeoutException.class, () -> 
sink.createWriter(null));
+    }
+
+    private CouchbaseSink createSink(int timeout, String password) {
+        Map<String, Object> config = new HashMap<>();
+        config.put("connection.string", server.getConnectionString());
+        config.put("username", server.getUsername());
+        config.put("password", password);
+        config.put("bucket", "readiness");
+        config.put("collection", "_default");
+        config.put("primary-key", Collections.singletonList("id"));
+        config.put("upsert-enable", true);
+        config.put("ready.timeout", timeout);
+        CatalogTable table =
+                CatalogTable.of(
+                        TableIdentifier.of("catalog", "database", "table"),
+                        TableSchema.builder()
+                                .column(
+                                        PhysicalColumn.of(
+                                                "id", BasicType.STRING_TYPE, 
64L, false, null, ""))
+                                .build(),
+                        Collections.emptyMap(),
+                        Collections.emptyList(),
+                        "");
+        CouchbaseSinkFactory factory = new CouchbaseSinkFactory();
+        ReadonlyConfig options = ReadonlyConfig.fromMap(config);
+        ConfigValidator.of(options).validate(factory.optionRule());
+        return (CouchbaseSink)
+                factory.createSink(
+                                new TableSinkFactoryContext(
+                                        table, options, 
getClass().getClassLoader()))
+                        .createSink();
+    }
+}

Reply via email to