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 2a77caf1cf [Feature][Connector-V2] Support --dry-run connect for Redis
source and sink (#12431)
2a77caf1cf is described below
commit 2a77caf1cf4a597c6ebb2a00322ca9b6279c5146
Author: junvelop <[email protected]>
AuthorDate: Fri Oct 2 08:11:31 2026 +0000
[Feature][Connector-V2] Support --dry-run connect for Redis source and sink
(#12431)
Co-authored-by: Kimyeojun <[email protected]>
---
docs/en/connectors/sink/Redis.md | 16 +
docs/en/connectors/source/Redis.md | 17 +
docs/en/engines/zeta/user-command.md | 1 +
docs/zh/connectors/sink/Redis.md | 13 +
docs/zh/connectors/source/Redis.md | 14 +
docs/zh/engines/zeta/user-command.md | 1 +
.../redis/client/RedisDryRunValidator.java | 98 ++++++
.../seatunnel/redis/sink/RedisSinkFactory.java | 17 +-
.../seatunnel/redis/source/RedisSourceFactory.java | 22 +-
.../seatunnel/redis/RedisFactoryTest.java | 105 ++++++
.../redis/client/RedisDryRunValidatorTest.java | 359 +++++++++++++++++++++
.../e2e/connector/redis/RedisConnectDryRunIT.java | 263 +++++++++++++++
12 files changed, 924 insertions(+), 2 deletions(-)
diff --git a/docs/en/connectors/sink/Redis.md b/docs/en/connectors/sink/Redis.md
index 760176e433..ddebaf9d80 100644
--- a/docs/en/connectors/sink/Redis.md
+++ b/docs/en/connectors/sink/Redis.md
@@ -12,6 +12,22 @@ Redis Cluster, and can write to `key`/`string`, `hash`,
`list`, `set`, and `zset
The configured `key` can be either a literal Redis key or an upstream field
name. When `support_custom_key = true`,
the connector can build the Redis key from one or more upstream fields, for
example `user:${id}`.
+### Connectivity dry-run
+
+Zeta's `--dry-run connect` checks that Redis is reachable and accepts the
configured credentials.
+The client is created through the same connection setup as normal job
execution, so `user` and
+`auth` are verified exactly as at runtime: `AUTH user auth` when `user` is
set, `AUTH auth` when only
+`auth` is set. In `SINGLE` mode it then sends `SELECT db_num` and `PING`. In
`CLUSTER` mode it
+initializes the cluster slot cache from `nodes` (`CLUSTER SLOTS`) and reads
`INFO` from one node.
+The runtime connect and socket timeouts (2 seconds) apply, and every client is
closed on success and
+on failure. No key is read, scanned, written or expired, no key space is
created and no ACL entry is
+modified. Normal job execution is unchanged.
+
+Successful validation does **not** prove write permission on the target keys.
`key`, `value_field`,
+`hash_key_field` and `hash_value_field` are not checked against the upstream
schema, because a name
+that is not an upstream field is written as a literal value at runtime. In
`CLUSTER` mode,
+validation passes as long as one node answers, so partially unreachable
clusters are not detected.
+
## Support Those Engines
> Spark<br/>
diff --git a/docs/en/connectors/source/Redis.md
b/docs/en/connectors/source/Redis.md
index 093b4e5cdb..3b38a657b6 100644
--- a/docs/en/connectors/source/Redis.md
+++ b/docs/en/connectors/source/Redis.md
@@ -8,6 +8,23 @@ import ChangeLog from '../changelog/connector-redis.md';
Used to read data from Redis.
+### Connectivity dry-run
+
+Zeta's `--dry-run connect` checks that Redis is reachable and accepts the
configured credentials.
+The client is created through the same connection setup as normal job
execution, so `user` and
+`auth` are verified exactly as at runtime: `AUTH user auth` when `user` is
set, `AUTH auth` when only
+`auth` is set. In `SINGLE` mode it then sends `SELECT db_num` and `PING`. In
`CLUSTER` mode it
+initializes the cluster slot cache from `nodes` (`CLUSTER SLOTS`) and reads
`INFO` from one node.
+The runtime connect and socket timeouts (2 seconds) apply, and every client is
closed on success and
+on failure. No key is read, scanned, written or expired, no key space is
created and no ACL entry is
+modified. Normal job execution is unchanged.
+
+Output schemas come from the configured `schema` or `tables_configs`, through
the same path as normal
+execution; no Redis value is inspected. Successful validation does **not**
prove that matching keys
+exist, that stored values match `data_type` or `format`, or that the
credentials may read those keys.
+In `CLUSTER` mode, validation passes as long as one node answers, so partially
unreachable clusters
+are not detected.
+
## Support Those Engines
> Spark<br/>
diff --git a/docs/en/engines/zeta/user-command.md
b/docs/en/engines/zeta/user-command.md
index cbe38044e6..53e0ea4141 100644
--- a/docs/en/engines/zeta/user-command.md
+++ b/docs/en/engines/zeta/user-command.md
@@ -84,6 +84,7 @@ The `--dry-run connect` option runs the static checks first,
then uses connector
| Kafka | Yes ([topic metadata + runtime output
schema](../../connectors/source/Kafka.md#connectivity-dry-run), not
consumer/group permissions) | No |
| FakeSource | Yes (schema inference only, no external system) | - |
| S3File | Yes (metadata connectivity + inline schema, single-table
text/csv/json/xml; see [supported
scope](../../connectors/source/S3File.md#connectivity-dry-run)) | - |
+| Redis | Yes ([connectivity +
authentication](../../connectors/source/Redis.md#connectivity-dry-run) +
configured schema; no key access) | Yes ([connectivity +
authentication](../../connectors/sink/Redis.md#connectivity-dry-run); no field
compatibility or write permission check) |
Every plugin in the job is reported in a validation summary with one of two
statuses:
diff --git a/docs/zh/connectors/sink/Redis.md b/docs/zh/connectors/sink/Redis.md
index 3c59c1c1ab..e4a97f7668 100644
--- a/docs/zh/connectors/sink/Redis.md
+++ b/docs/zh/connectors/sink/Redis.md
@@ -12,6 +12,19 @@ Redis 接收器连接器可以在批处理或流处理作业中把上游数据
`key` 可以是固定的 Redis key,也可以是上游字段名。开启 `support_custom_key = true` 后,还可以用上游字段
拼出 Redis key,例如 `user:${id}`。
+### 连通性 dry-run
+
+Zeta 的 `--dry-run connect` 会校验 Redis 是否可达以及是否接受配置的凭据。客户端通过与正常作业
+运行相同的连接逻辑创建,因此 `user` 和 `auth` 会按运行时的方式进行验证:配置了 `user` 时发送
+`AUTH user auth`,仅配置了 `auth` 时发送 `AUTH auth`。`SINGLE` 模式下随后发送 `SELECT db_num` 和
+`PING`。`CLUSTER` 模式下会基于 `nodes` 初始化集群 slot 缓存(`CLUSTER SLOTS`),并从一个节点读取
+`INFO`。连接超时和 socket 超时沿用运行时的默认值(2 秒),无论成功或失败都会关闭所有客户端。校验不会
+读取、扫描、写入任何 key,不会设置过期时间或创建 key 空间,也不会修改任何 ACL 条目。正常作业运行保持不变。
+
+校验成功**不代表**具备目标 key 的写入权限。`key`、`value_field`、`hash_key_field` 和
+`hash_value_field` 不会与上游 schema 进行比对,因为运行时如果名称不是上游字段,会作为字面值写入。
+`CLUSTER` 模式下只要有一个节点响应即可通过校验,因此无法发现集群中部分节点不可达的情况。
+
## 支持引擎
> Spark<br/>
diff --git a/docs/zh/connectors/source/Redis.md
b/docs/zh/connectors/source/Redis.md
index 77f84c018d..042868ab08 100644
--- a/docs/zh/connectors/source/Redis.md
+++ b/docs/zh/connectors/source/Redis.md
@@ -8,6 +8,20 @@ import ChangeLog from '../changelog/connector-redis.md';
用于从 `Redis` 读取数据
+### 连通性 dry-run
+
+Zeta 的 `--dry-run connect` 会校验 Redis 是否可达以及是否接受配置的凭据。客户端通过与正常作业
+运行相同的连接逻辑创建,因此 `user` 和 `auth` 会按运行时的方式进行验证:配置了 `user` 时发送
+`AUTH user auth`,仅配置了 `auth` 时发送 `AUTH auth`。`SINGLE` 模式下随后发送 `SELECT db_num` 和
+`PING`。`CLUSTER` 模式下会基于 `nodes` 初始化集群 slot 缓存(`CLUSTER SLOTS`),并从一个节点读取
+`INFO`。连接超时和 socket 超时沿用运行时的默认值(2 秒),无论成功或失败都会关闭所有客户端。校验不会
+读取、扫描、写入任何 key,不会设置过期时间或创建 key 空间,也不会修改任何 ACL 条目。正常作业运行保持不变。
+
+输出 schema 通过与正常运行相同的路径,从配置的 `schema` 或 `tables_configs` 中获取,不会读取任何
+Redis 值。校验成功**不代表**匹配的 key 存在、已存储的值与 `data_type` 或 `format` 相符,也不代表
+凭据具备读取这些 key 的权限。`CLUSTER` 模式下只要有一个节点响应即可通过校验,因此无法发现集群中
+部分节点不可达的情况。
+
## 支持引擎
> Spark<br/>
diff --git a/docs/zh/engines/zeta/user-command.md
b/docs/zh/engines/zeta/user-command.md
index d708124799..c491ac9cdd 100644
--- a/docs/zh/engines/zeta/user-command.md
+++ b/docs/zh/engines/zeta/user-command.md
@@ -100,6 +100,7 @@ bin/seatunnel.sh --config
$SEATUNNEL_HOME/config/v2.batch.config.template --dry-
| Kafka | 支持([主题元数据 + 运行时输出
schema](../../connectors/source/Kafka.md#连通性-dry-run),不含消费或消费组权限) | 不支持 |
| FakeSource | 支持(仅 schema 推断,无外部系统) | - |
| S3File | 支持(元数据连通性 + 内联 schema,仅单表
text/csv/json/xml;参见[支持范围](../../connectors/source/S3File.md#连接预检查)) | - |
+| Redis | 支持([连通性 + 认证](../../connectors/source/Redis.md#连通性-dry-run) + 配置的
schema,不访问 key) | 支持([连通性 +
认证](../../connectors/sink/Redis.md#连通性-dry-run),不校验字段兼容性和写入权限) |
作业中的每个插件都会在校验汇总中报告以下两种状态之一:
diff --git
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/client/RedisDryRunValidator.java
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/client/RedisDryRunValidator.java
new file mode 100644
index 0000000000..56bebbdccc
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/client/RedisDryRunValidator.java
@@ -0,0 +1,98 @@
+/*
+ * 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.redis.client;
+
+import org.apache.seatunnel.common.exception.CommonErrorCode;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters;
+import
org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisConnectorException;
+import
org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisErrorCode;
+
+import redis.clients.jedis.Jedis;
+
+/**
+ * Connectivity and authentication check for {@code --dry-run connect}, shared
by the Redis source
+ * and sink.
+ *
+ * <p>The client is created through {@link RedisParameters#buildJedis()}, the
same connection setup
+ * the runtime uses, so the {@code user} and {@code auth} options are verified
exactly as they are
+ * during job execution. Afterwards only {@code PING} (single node) or {@code
INFO} (cluster) is
+ * issued: no key is read, scanned, written or expired, no key space is
created and no ACL is
+ * modified. The client is closed on success and on failure.
+ */
+public final class RedisDryRunValidator {
+
+ private RedisDryRunValidator() {}
+
+ /**
+ * Opens one short-lived connection with the runtime connection setup and
closes it again.
+ *
+ * @param parameters connection parameters built from the source or sink
options
+ * @throws RedisConnectorException with {@link
RedisErrorCode#REDIS_CONNECTION_ERROR} naming the
+ * validated target when Redis is unreachable or rejects the
credentials
+ */
+ public static void validate(RedisParameters parameters) {
+ String target = target(parameters);
+ Jedis jedis;
+ try {
+ // Sends AUTH (and SELECT in SINGLE mode, CLUSTER SLOTS in CLUSTER
mode) and already
+ // closes the connection itself when that fails.
+ jedis = parameters.buildJedis();
+ } catch (RuntimeException e) {
+ throw connectionError(target, e);
+ }
+ try (Jedis client = jedis) {
+ switch (parameters.getMode()) {
+ case SINGLE:
+ client.ping();
+ break;
+ case CLUSTER:
+ // JedisWrapper has no connection of its own, so INFO is
read from a node.
+ client.info();
+ break;
+ default:
+ throw unsupportedMode();
+ }
+ } catch (RuntimeException e) {
+ throw connectionError(target, e);
+ }
+ }
+
+ private static String target(RedisParameters parameters) {
+ switch (parameters.getMode()) {
+ case SINGLE:
+ return parameters.getHost() + ":" + parameters.getPort();
+ case CLUSTER:
+ return String.valueOf(parameters.getRedisNodes());
+ default:
+ throw unsupportedMode();
+ }
+ }
+
+ private static RedisConnectorException unsupportedMode() {
+ return new RedisConnectorException(
+ CommonErrorCode.OPERATION_NOT_SUPPORTED, "Not support this
redis mode");
+ }
+
+ private static RedisConnectorException connectionError(String target,
Throwable cause) {
+ return new RedisConnectorException(
+ RedisErrorCode.REDIS_CONNECTION_ERROR,
+ String.format(
+ "Redis connect dry-run failed for %s: %s", target,
cause.getMessage()),
+ cause);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/sink/RedisSinkFactory.java
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/sink/RedisSinkFactory.java
index 7730fec61a..0984a6a900 100644
---
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/sink/RedisSinkFactory.java
+++
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/sink/RedisSinkFactory.java
@@ -23,16 +23,19 @@ import
org.apache.seatunnel.api.options.SinkConnectorCommonOptions;
import org.apache.seatunnel.api.table.catalog.CatalogTable;
import org.apache.seatunnel.api.table.connector.TableSink;
import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.SupportSinkDryRunValidation;
import org.apache.seatunnel.api.table.factory.TableSinkFactory;
import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import
org.apache.seatunnel.connectors.seatunnel.redis.client.RedisDryRunValidator;
import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisBaseOptions;
import
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisNodesValidator;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters;
import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisSinkOptions;
import com.google.auto.service.AutoService;
@AutoService(Factory.class)
-public class RedisSinkFactory implements TableSinkFactory {
+public class RedisSinkFactory implements TableSinkFactory,
SupportSinkDryRunValidation {
@Override
public String factoryIdentifier() {
return "Redis";
@@ -86,4 +89,16 @@ public class RedisSinkFactory implements TableSinkFactory {
Conditions.extension(RedisBaseOptions.NODES, new
RedisNodesValidator()))
.build();
}
+
+ /**
+ * Checks connectivity and authentication only; no key is written or
expired. Sink field options
+ * are not checked against the upstream schema because the runtime falls
back to the configured
+ * name as a literal value.
+ */
+ @Override
+ public void validateConnectionForDryRun(TableSinkFactoryContext context) {
+ RedisParameters redisParameters = new RedisParameters();
+ redisParameters.buildConnectionConfig(context.getOptions());
+ RedisDryRunValidator.validate(redisParameters);
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/source/RedisSourceFactory.java
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/source/RedisSourceFactory.java
index c44700e56e..93b750fb48 100644
---
a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/source/RedisSourceFactory.java
+++
b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/source/RedisSourceFactory.java
@@ -21,12 +21,16 @@ import
org.apache.seatunnel.api.configuration.util.Conditions;
import org.apache.seatunnel.api.configuration.util.OptionRule;
import org.apache.seatunnel.api.source.SeaTunnelSource;
import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
import org.apache.seatunnel.api.table.connector.TableSource;
import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation;
import org.apache.seatunnel.api.table.factory.TableSourceFactory;
import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import
org.apache.seatunnel.connectors.seatunnel.redis.client.RedisDryRunValidator;
import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisBaseOptions;
import
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisNodesValidator;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters;
import
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisSingleTableDataTypeValidator;
import
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisSourceOptions;
import
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisTableConfigsValidator;
@@ -34,9 +38,10 @@ import
org.apache.seatunnel.connectors.seatunnel.redis.config.RedisTableConfigsV
import com.google.auto.service.AutoService;
import java.io.Serializable;
+import java.util.List;
@AutoService(Factory.class)
-public class RedisSourceFactory implements TableSourceFactory {
+public class RedisSourceFactory implements TableSourceFactory,
SupportSourceDryRunValidation {
@Override
public String factoryIdentifier() {
return "Redis";
@@ -110,4 +115,19 @@ public class RedisSourceFactory implements
TableSourceFactory {
public Class<? extends SeaTunnelSource> getSourceClass() {
return RedisSource.class;
}
+
+ /** Uses the configured schema only; no Redis value is inspected and no
reader is created. */
+ @Override
+ public List<CatalogTable> inferSchemaForDryRun(TableSourceFactoryContext
context) {
+ return new
RedisSource(context.getOptions()).getProducedCatalogTables();
+ }
+
+ /** Checks connectivity and authentication only; no key is read or
scanned. */
+ @Override
+ public void validateConnectionForDryRun(
+ TableSourceFactoryContext context, List<CatalogTable>
catalogTables) {
+ RedisParameters redisParameters = new RedisParameters();
+ redisParameters.buildConnectionConfig(context.getOptions());
+ RedisDryRunValidator.validate(redisParameters);
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
index f1b6d3850e..5949155490 100644
---
a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
+++
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/RedisFactoryTest.java
@@ -51,10 +51,16 @@ import java.util.ArrayList;
import java.util.Arrays;
import java.util.Collections;
import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.List;
import java.util.Map;
+import java.util.stream.Collectors;
import java.util.stream.Stream;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+
class RedisFactoryTest {
private static final OptionRule SOURCE_RULE = new
RedisSourceFactory().optionRule();
@@ -235,6 +241,95 @@ class RedisFactoryTest {
tableSchema,
serializer.deserialize(serializer.serialize(tableSchema)));
}
+ // connect dry-run
+
+ @Test
+ void sourceDryRunInfersConfiguredSchemaForSingleTable() {
+ Map<String, Object> config = singleSourceConfig();
+ config.put("format", "JSON");
+ config.put("schema", schema("db.users"));
+ TableSourceFactoryContext context =
+ new TableSourceFactoryContext(
+ ReadonlyConfig.fromMap(config),
+ Thread.currentThread().getContextClassLoader());
+
+ // No Redis instance is running, so passing here proves no value is
read.
+ List<CatalogTable> tables = new
RedisSourceFactory().inferSchemaForDryRun(context);
+
+ Assertions.assertEquals(1, tables.size());
+ Assertions.assertEquals(
+ Arrays.asList("id", "name"),
+ Arrays.asList(tables.get(0).getTableSchema().getFieldNames()));
+ }
+
+ @Test
+ void sourceDryRunInfersOneTablePerTableConfig() {
+ Map<String, Object> first = tableEntry("key_a*");
+ first.put("schema", schema("db.a"));
+ Map<String, Object> second = tableEntry("key_b*");
+ second.put("schema", schema("db.b"));
+ Map<String, Object> config = multiTableSourceConfig();
+ config.put("tables_configs", Arrays.asList(first, second));
+ TableSourceFactoryContext context =
+ new TableSourceFactoryContext(
+ ReadonlyConfig.fromMap(config),
+ Thread.currentThread().getContextClassLoader());
+
+ List<CatalogTable> tables = new
RedisSourceFactory().inferSchemaForDryRun(context);
+
+ Assertions.assertEquals(
+ Arrays.asList("db.a", "db.b"),
+ tables.stream()
+ .map(table -> table.getTablePath().toString())
+ .collect(Collectors.toList()));
+ for (CatalogTable table : tables) {
+ Assertions.assertEquals(
+ Arrays.asList("id", "name"),
+ Arrays.asList(table.getTableSchema().getFieldNames()));
+ }
+ }
+
+ @Test
+ void sourceDryRunValidatesConnectionWithoutReadingKeys() throws Exception {
+ TableSourceFactoryContext context =
+ new TableSourceFactoryContext(
+ ReadonlyConfig.fromMap(singleSourceConfig()),
+ Thread.currentThread().getContextClassLoader());
+ try (MockedConstruction<Jedis> clients =
mockConstruction(Jedis.class)) {
+ new RedisSourceFactory().validateConnectionForDryRun(context,
Collections.emptyList());
+
+ Jedis jedis = clients.constructed().get(0);
+ verify(jedis).select(0);
+ verify(jedis).ping();
+ verify(jedis).close();
+ verifyNoMoreInteractions(jedis);
+ }
+ }
+
+ @Test
+ void sinkDryRunIgnoresFieldOptionsMissingFromUpstreamSchema() {
+ // Missing field names are literal values at runtime, so they must not
fail the dry run.
+ Map<String, Object> config = singleSinkConfig();
+ config.put("data_type", "HASH");
+ config.put("key", "missing_key_field");
+ config.put("hash_key_field", "missing_hash_key");
+ config.put("hash_value_field", "missing_hash_value");
+ TableSinkFactoryContext context =
+ new TableSinkFactoryContext(
+ catalogTable(),
+ ReadonlyConfig.fromMap(config),
+ Thread.currentThread().getContextClassLoader());
+ try (MockedConstruction<Jedis> clients =
mockConstruction(Jedis.class)) {
+ new RedisSinkFactory().validateConnectionForDryRun(context);
+
+ Jedis jedis = clients.constructed().get(0);
+ verify(jedis).select(0);
+ verify(jedis).ping();
+ verify(jedis).close();
+ verifyNoMoreInteractions(jedis);
+ }
+ }
+
// parameterized-case providers
static Stream<Arguments> invalidSingleConnections() {
@@ -332,6 +427,16 @@ class RedisFactoryTest {
"catalog");
}
+ private static Map<String, Object> schema(String table) {
+ Map<String, Object> fields = new LinkedHashMap<>();
+ fields.put("id", "bigint");
+ fields.put("name", "string");
+ Map<String, Object> schema = new HashMap<>();
+ schema.put("table", table);
+ schema.put("fields", fields);
+ return schema;
+ }
+
private static Map<String, Object> singleSourceConfig() {
Map<String, Object> config = new HashMap<>();
config.put("mode", "SINGLE");
diff --git
a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/client/RedisDryRunValidatorTest.java
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/client/RedisDryRunValidatorTest.java
new file mode 100644
index 0000000000..062ea88959
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/client/RedisDryRunValidatorTest.java
@@ -0,0 +1,359 @@
+/*
+ * 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.redis.client;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.JedisWrapper;
+import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters;
+import
org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisConnectorException;
+import
org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisErrorCode;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
+import org.mockito.MockedConstruction;
+
+import redis.clients.jedis.DefaultJedisClientConfig;
+import redis.clients.jedis.HostAndPort;
+import redis.clients.jedis.Jedis;
+import redis.clients.jedis.JedisCluster;
+import redis.clients.jedis.exceptions.JedisClusterOperationException;
+import redis.clients.jedis.exceptions.JedisConnectionException;
+import redis.clients.jedis.exceptions.JedisDataException;
+
+import java.util.Arrays;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.Mockito.inOrder;
+import static org.mockito.Mockito.mockConstruction;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
+import static org.mockito.Mockito.when;
+
+class RedisDryRunValidatorTest {
+
+ private static final String USER = "synthetic-user";
+ private static final String PASSWORD = "synthetic-secret";
+ private static final String SINGLE_TARGET = "localhost:6379";
+ private static final String CLUSTER_TARGET = "[127.0.0.1:7000,
127.0.0.1:7001]";
+
+ @Test
+ void singleNodeAuthenticatesNamedUserSelectsPingsAndCloses() {
+ try (MockedConstruction<Jedis> clients =
mockConstruction(Jedis.class)) {
+ RedisDryRunValidator.validate(singleParameters(USER, PASSWORD));
+
+ assertEquals(1, clients.constructed().size());
+ Jedis jedis = clients.constructed().get(0);
+ InOrder order = inOrder(jedis);
+ order.verify(jedis).auth(USER, PASSWORD);
+ order.verify(jedis).select(2);
+ order.verify(jedis).ping();
+ order.verify(jedis).close();
+ // Proves that no key or ACL command is issued.
+ verifyNoMoreInteractions(jedis);
+ }
+ }
+
+ @Test
+ void singleNodeAuthenticatesPasswordOnly() {
+ try (MockedConstruction<Jedis> clients =
mockConstruction(Jedis.class)) {
+ RedisDryRunValidator.validate(singleParameters(null, PASSWORD));
+
+ Jedis jedis = clients.constructed().get(0);
+ InOrder order = inOrder(jedis);
+ order.verify(jedis).auth(PASSWORD);
+ order.verify(jedis).select(2);
+ order.verify(jedis).ping();
+ order.verify(jedis).close();
+ verifyNoMoreInteractions(jedis);
+ }
+ }
+
+ @Test
+ void singleNodeUsesRuntimeConnectionSetup() {
+ try (MockedConstruction<Jedis> clients =
+ mockConstruction(
+ Jedis.class,
+ (jedis, context) ->
+ assertEquals(
+ Arrays.asList("localhost", 6379),
context.arguments()))) {
+ RedisDryRunValidator.validate(singleParameters(null, null));
+
+ assertEquals(1, clients.constructed().size());
+ }
+ }
+
+ @Test
+ void singleNodeWithoutAuthDoesNotAuthenticate() {
+ try (MockedConstruction<Jedis> clients =
mockConstruction(Jedis.class)) {
+ RedisDryRunValidator.validate(singleParameters(null, null));
+
+ Jedis jedis = clients.constructed().get(0);
+ verify(jedis, never()).auth(anyString());
+ verify(jedis, never()).auth(anyString(), anyString());
+ verify(jedis).select(2);
+ verify(jedis).ping();
+ verify(jedis).close();
+ verifyNoMoreInteractions(jedis);
+ }
+ }
+
+ @Test
+ void singleNodeAuthenticationFailureIsReportedAndClientClosed() {
+ JedisDataException cause =
+ new JedisDataException("WRONGPASS invalid username-password
pair");
+ try (MockedConstruction<Jedis> clients =
+ mockConstruction(
+ Jedis.class,
+ (jedis, context) -> when(jedis.auth(USER,
PASSWORD)).thenThrow(cause))) {
+ RedisConnectorException exception =
+ assertThrows(
+ RedisConnectorException.class,
+ () ->
RedisDryRunValidator.validate(singleParameters(USER, PASSWORD)));
+
+ assertConnectionError(exception, SINGLE_TARGET, cause);
+ Jedis jedis = clients.constructed().get(0);
+ verify(jedis, never()).select(anyInt());
+ verify(jedis, never()).ping();
+ verify(jedis).close();
+ }
+ }
+
+ @Test
+ void singleNodeConnectionFailureIsReportedAndClientClosed() {
+ JedisConnectionException cause =
+ new JedisConnectionException("java.net.ConnectException:
Connection refused");
+ try (MockedConstruction<Jedis> clients =
+ mockConstruction(
+ Jedis.class, (jedis, context) ->
when(jedis.select(2)).thenThrow(cause))) {
+ RedisConnectorException exception =
+ assertThrows(
+ RedisConnectorException.class,
+ () ->
RedisDryRunValidator.validate(singleParameters(null, null)));
+
+ assertConnectionError(exception, SINGLE_TARGET, cause);
+ verify(clients.constructed().get(0)).close();
+ }
+ }
+
+ @Test
+ void singleNodePingFailureIsReportedAndClientClosed() {
+ JedisConnectionException cause = new
JedisConnectionException("Unexpected end of stream.");
+ try (MockedConstruction<Jedis> clients =
+ mockConstruction(
+ Jedis.class, (jedis, context) ->
when(jedis.ping()).thenThrow(cause))) {
+ RedisConnectorException exception =
+ assertThrows(
+ RedisConnectorException.class,
+ () ->
RedisDryRunValidator.validate(singleParameters(null, PASSWORD)));
+
+ assertConnectionError(exception, SINGLE_TARGET, cause);
+ verify(clients.constructed().get(0)).close();
+ }
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void clusterAuthenticatesNamedUserReadsInfoAndCloses() {
+ try (MockedConstruction<JedisCluster> clusters =
+ mockConstruction(
+ JedisCluster.class,
+ (cluster, context) -> {
+ List<?> arguments = context.arguments();
+ assertClusterNodes((Set<HostAndPort>)
arguments.get(0));
+ DefaultJedisClientConfig config =
+ (DefaultJedisClientConfig)
arguments.get(1);
+ assertEquals(USER, config.getUser());
+ assertEquals(PASSWORD,
config.getPassword());
+ });
+ MockedConstruction<JedisWrapper> wrappers =
+ mockConstruction(
+ JedisWrapper.class,
+ (wrapper, context) ->
+
when(wrapper.info()).thenReturn("redis_version:7.0.0"))) {
+ RedisDryRunValidator.validate(clusterParameters(USER, PASSWORD));
+
+ assertEquals(1, clusters.constructed().size());
+ JedisWrapper wrapper = wrappers.constructed().get(0);
+ InOrder order = inOrder(wrapper);
+ order.verify(wrapper).info();
+ order.verify(wrapper).close();
+ verifyNoMoreInteractions(wrapper);
+ }
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void clusterAuthenticatesPasswordOnly() {
+ try (MockedConstruction<JedisCluster> clusters =
+ mockConstruction(
+ JedisCluster.class,
+ (cluster, context) -> {
+ List<?> arguments = context.arguments();
+ assertClusterNodes((Set<HostAndPort>)
arguments.get(0));
+ assertEquals(PASSWORD, arguments.get(4));
+ });
+ MockedConstruction<JedisWrapper> wrappers =
mockConstruction(JedisWrapper.class)) {
+ RedisDryRunValidator.validate(clusterParameters(null, PASSWORD));
+
+ assertEquals(1, clusters.constructed().size());
+ JedisWrapper wrapper = wrappers.constructed().get(0);
+ verify(wrapper).info();
+ verify(wrapper).close();
+ }
+ }
+
+ @Test
+ @SuppressWarnings("unchecked")
+ void clusterWithoutAuthPassesNoCredentials() {
+ try (MockedConstruction<JedisCluster> clusters =
+ mockConstruction(
+ JedisCluster.class,
+ (cluster, context) -> {
+ List<?> arguments = context.arguments();
+ assertEquals(1, arguments.size());
+ assertClusterNodes((Set<HostAndPort>)
arguments.get(0));
+ });
+ MockedConstruction<JedisWrapper> wrappers =
mockConstruction(JedisWrapper.class)) {
+ RedisDryRunValidator.validate(clusterParameters(null, null));
+
+ assertEquals(1, clusters.constructed().size());
+ verify(wrappers.constructed().get(0)).close();
+ }
+ }
+
+ @Test
+ void clusterWithNoReachableNodeIsReported() {
+ JedisClusterOperationException cause =
+ new JedisClusterOperationException("Could not initialize
cluster slots cache.");
+ try (MockedConstruction<JedisCluster> clusters =
+ mockConstruction(
+ JedisCluster.class,
+ (cluster, context) -> {
+ throw cause;
+ });
+ MockedConstruction<JedisWrapper> wrappers =
mockConstruction(JedisWrapper.class)) {
+ RedisConnectorException exception =
+ assertThrows(
+ RedisConnectorException.class,
+ () ->
RedisDryRunValidator.validate(clusterParameters(USER, PASSWORD)));
+
+ assertConnectionError(exception, CLUSTER_TARGET, cause);
+ assertTrue(wrappers.constructed().isEmpty());
+ }
+ }
+
+ @Test
+ void clusterInfoFailureIsReportedAndClientClosed() {
+ RedisConnectorException cause =
+ new RedisConnectorException(
+ RedisErrorCode.GET_REDIS_INFO_ERROR,
+ "Failed to get redis info from all node in cluster");
+ try (MockedConstruction<JedisCluster> clusters =
mockConstruction(JedisCluster.class);
+ MockedConstruction<JedisWrapper> wrappers =
+ mockConstruction(
+ JedisWrapper.class,
+ (wrapper, context) ->
when(wrapper.info()).thenThrow(cause))) {
+ RedisConnectorException exception =
+ assertThrows(
+ RedisConnectorException.class,
+ () ->
RedisDryRunValidator.validate(clusterParameters(null, PASSWORD)));
+
+ assertConnectionError(exception, CLUSTER_TARGET, cause);
+ verify(wrappers.constructed().get(0)).close();
+ }
+ }
+
+ @Test
+ void clusterMalformedNodeIsReportedAsConnectionError() {
+ Map<String, Object> config = new HashMap<>();
+ config.put("mode", "CLUSTER");
+ config.put("nodes", Arrays.asList("127.0.0.1:7000",
"host-without-port"));
+ RedisParameters parameters = parameters(config);
+
+ RedisConnectorException exception =
+ assertThrows(
+ RedisConnectorException.class,
+ () -> RedisDryRunValidator.validate(parameters));
+
+ assertEquals(RedisErrorCode.REDIS_CONNECTION_ERROR,
exception.getSeaTunnelErrorCode());
+ assertTrue(exception.getMessage().contains("host-without-port"),
exception.getMessage());
+ }
+
+ private static void assertClusterNodes(Set<HostAndPort> nodes) {
+ assertEquals(2, nodes.size());
+ assertTrue(nodes.contains(new HostAndPort("127.0.0.1", 7000)));
+ assertTrue(nodes.contains(new HostAndPort("127.0.0.1", 7001)));
+ }
+
+ private static void assertConnectionError(
+ RedisConnectorException exception, String target, Throwable cause)
{
+ assertEquals(RedisErrorCode.REDIS_CONNECTION_ERROR,
exception.getSeaTunnelErrorCode());
+ assertTrue(exception.getMessage().contains(target),
exception.getMessage());
+ assertFalse(exception.getMessage().contains(PASSWORD),
exception.getMessage());
+ // Mockito wraps exceptions thrown from a mocked constructor, so
search the cause chain.
+ Throwable current = exception.getCause();
+ while (current != null && current != cause) {
+ current = current.getCause();
+ }
+ assertSame(cause, current);
+ }
+
+ private static RedisParameters singleParameters(String user, String auth) {
+ Map<String, Object> config = new HashMap<>();
+ config.put("mode", "SINGLE");
+ config.put("host", "localhost");
+ config.put("port", 6379);
+ config.put("db_num", 2);
+ return parameters(withCredentials(config, user, auth));
+ }
+
+ private static RedisParameters clusterParameters(String user, String auth)
{
+ Map<String, Object> config = new HashMap<>();
+ config.put("mode", "CLUSTER");
+ config.put("nodes", Arrays.asList("127.0.0.1:7000", "127.0.0.1:7001"));
+ return parameters(withCredentials(config, user, auth));
+ }
+
+ private static Map<String, Object> withCredentials(
+ Map<String, Object> config, String user, String auth) {
+ if (user != null) {
+ config.put("user", user);
+ }
+ if (auth != null) {
+ config.put("auth", auth);
+ }
+ return config;
+ }
+
+ private static RedisParameters parameters(Map<String, Object> config) {
+ RedisParameters parameters = new RedisParameters();
+ parameters.buildConnectionConfig(ReadonlyConfig.fromMap(config));
+ return parameters;
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisConnectDryRunIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisConnectDryRunIT.java
new file mode 100644
index 0000000000..4146c21bd7
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisConnectDryRunIT.java
@@ -0,0 +1,263 @@
+/*
+ * 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.redis;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+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.SupportSinkDryRunValidation;
+import org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import org.apache.seatunnel.api.table.type.BasicType;
+import
org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisConnectorException;
+import
org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisErrorCode;
+import org.apache.seatunnel.connectors.seatunnel.redis.sink.RedisSinkFactory;
+import
org.apache.seatunnel.connectors.seatunnel.redis.source.RedisSourceFactory;
+import org.apache.seatunnel.e2e.common.TestResource;
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+
+import org.junit.jupiter.api.AfterAll;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.testcontainers.containers.GenericContainer;
+import org.testcontainers.containers.output.Slf4jLogConsumer;
+import org.testcontainers.containers.wait.strategy.HostPortWaitStrategy;
+import org.testcontainers.utility.DockerImageName;
+import org.testcontainers.utility.DockerLoggerFactory;
+
+import lombok.extern.slf4j.Slf4j;
+import redis.clients.jedis.Jedis;
+
+import java.io.IOException;
+import java.net.InetAddress;
+import java.net.ServerSocket;
+import java.net.Socket;
+import java.time.Duration;
+import java.util.ArrayList;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+/** Exercises the factory-level connect dry-run contract against Redis without
submitting a job. */
+@Slf4j
+public class RedisConnectDryRunIT extends TestSuiteBase implements
TestResource {
+
+ private static final String IMAGE = "redis:7";
+ private static final int REDIS_PORT = 6379;
+ private static final String PASSWORD = "test-only-password";
+ private static final String USER = "dry-run-user";
+ private static final String USER_PASSWORD = "test-only-user-password";
+ private static final int DB_NUM = 3;
+
+ private GenericContainer<?> redis;
+ private Jedis admin;
+
+ @BeforeAll
+ @Override
+ public void startUp() {
+ redis =
+ new GenericContainer<>(DockerImageName.parse(IMAGE))
+ .withExposedPorts(REDIS_PORT)
+ .withCommand("redis-server --requirepass " + PASSWORD)
+ .withLogConsumer(new
Slf4jLogConsumer(DockerLoggerFactory.getLogger(IMAGE)))
+ .waitingFor(
+ new HostPortWaitStrategy()
+
.withStartupTimeout(Duration.ofMinutes(2)));
+ redis.start();
+ log.info("Password-protected Redis dry-run fixture started");
+ admin = new Jedis(redis.getHost(), redis.getFirstMappedPort());
+ admin.auth(PASSWORD);
+ // A named ACL user, so that the user option is exercised against a
real server.
+ admin.aclSetUser(USER, "on", ">" + USER_PASSWORD, "+@all", "~*");
+ }
+
+ @AfterAll
+ @Override
+ public void tearDown() {
+ try {
+ if (admin != null) {
+ admin.close();
+ }
+ } finally {
+ if (redis != null) {
+ redis.stop();
+ }
+ }
+ }
+
+ @Test
+ public void testSourceAndSinkValidationHaveNoKeyOrAclSideEffects() throws
Exception {
+ List<String> aclBefore = admin.aclList();
+
+ List<CatalogTable> tables = validateSource(null, PASSWORD,
redis.getFirstMappedPort());
+ validateSink(null, PASSWORD, redis.getFirstMappedPort());
+ validateSource(USER, USER_PASSWORD, redis.getFirstMappedPort());
+ validateSink(USER, USER_PASSWORD, redis.getFirstMappedPort());
+
+ assertEquals(1, tables.size());
+ // No key space is created, neither in the default db nor in the
configured db_num.
+ admin.select(0);
+ assertEquals(0L, admin.dbSize());
+ admin.select(DB_NUM);
+ assertEquals(0L, admin.dbSize());
+ // Validating with a named user must not create or alter any ACL entry.
+ assertEquals(aclBefore, admin.aclList());
+ }
+
+ @Test
+ public void testIncorrectPasswordFailsValidation() {
+ assertConnectionError(() -> validateSource(null, "incorrect",
redis.getFirstMappedPort()));
+ assertConnectionError(() -> validateSink(null, "incorrect",
redis.getFirstMappedPort()));
+ }
+
+ @Test
+ public void testMissingPasswordFailsValidation() {
+ assertConnectionError(() -> validateSource(null, null,
redis.getFirstMappedPort()));
+ assertConnectionError(() -> validateSink(null, null,
redis.getFirstMappedPort()));
+ }
+
+ @Test
+ public void testNamedUserCredentialsAreVerified() {
+ int port = redis.getFirstMappedPort();
+ // The named user's password is checked, not the server-wide
requirepass.
+ assertConnectionError(() -> validateSource(USER, PASSWORD, port));
+ assertConnectionError(() -> validateSink(USER, PASSWORD, port));
+ assertConnectionError(() -> validateSource("unknown-user",
USER_PASSWORD, port));
+ assertConnectionError(() -> validateSink("unknown-user",
USER_PASSWORD, port));
+ }
+
+ @Test
+ public void testUnreachableServerFailsValidation() throws Exception {
+ // A listener that drops every connection stands in for an unreachable
Redis; unlike a
+ // released ephemeral port it cannot be taken over by another process
mid-test.
+ try (ServerSocket server = droppingServer()) {
+ String host = server.getInetAddress().getHostAddress();
+ int port = server.getLocalPort();
+ assertConnectionError(() -> validateSource(null, PASSWORD, host,
port));
+ assertConnectionError(() -> validateSink(null, PASSWORD, host,
port));
+ }
+ }
+
+ private List<CatalogTable> validateSource(String user, String password,
int port)
+ throws Exception {
+ return validateSource(user, password, redis.getHost(), port);
+ }
+
+ private List<CatalogTable> validateSource(String user, String password,
String host, int port)
+ throws Exception {
+ TableSourceFactoryContext context =
+ new TableSourceFactoryContext(
+ ReadonlyConfig.fromMap(options(user, password, host,
port)),
+ getClass().getClassLoader());
+ SupportSourceDryRunValidation validation = new RedisSourceFactory();
+ List<CatalogTable> tables = validation.inferSchemaForDryRun(context);
+ validation.validateConnectionForDryRun(context, tables);
+ return tables;
+ }
+
+ private void validateSink(String user, String password, int port) throws
Exception {
+ validateSink(user, password, redis.getHost(), port);
+ }
+
+ private void validateSink(String user, String password, String host, int
port)
+ throws Exception {
+ Map<String, Object> options = options(user, password, host, port);
+ options.remove("keys");
+ options.put("key", "id");
+ TableSinkFactoryContext context =
+ new TableSinkFactoryContext(
+ catalogTable(),
+ ReadonlyConfig.fromMap(options),
+ getClass().getClassLoader());
+ SupportSinkDryRunValidation validation = new RedisSinkFactory();
+ validation.validateConnectionForDryRun(context);
+ }
+
+ private static Map<String, Object> options(
+ String user, String password, String host, int port) {
+ Map<String, Object> options = new HashMap<>();
+ options.put("host", host);
+ options.put("port", port);
+ options.put("db_num", DB_NUM);
+ options.put("keys", "dry-run-*");
+ options.put("data_type", "KEY");
+ if (user != null) {
+ options.put("user", user);
+ }
+ if (password != null) {
+ options.put("auth", password);
+ }
+ return options;
+ }
+
+ private static void assertConnectionError(Validation validation) {
+ RedisConnectorException exception =
+ assertThrows(RedisConnectorException.class, validation::run);
+ assertEquals(RedisErrorCode.REDIS_CONNECTION_ERROR,
exception.getSeaTunnelErrorCode());
+ assertFalse(exception.getMessage().contains(PASSWORD),
exception.getMessage());
+ assertFalse(exception.getMessage().contains(USER_PASSWORD),
exception.getMessage());
+ assertTrue(exception.getMessage().contains("dry-run"),
exception.getMessage());
+ }
+
+ /** Listens on a loopback port and closes every accepted connection
without answering. */
+ private static ServerSocket droppingServer() throws IOException {
+ ServerSocket server = new ServerSocket(0, 1,
InetAddress.getLoopbackAddress());
+ Thread acceptor =
+ new Thread(
+ () -> {
+ while (!server.isClosed()) {
+ try (Socket ignored = server.accept()) {
+ // Close immediately so the client sees a
dropped connection.
+ } catch (IOException e) {
+ // The server socket was closed by the
test; stop accepting.
+ }
+ }
+ },
+ "redis-dry-run-dropping-server");
+ acceptor.setDaemon(true);
+ acceptor.start();
+ return server;
+ }
+
+ private static CatalogTable catalogTable() {
+ TableSchema schema =
+ TableSchema.builder()
+ .column(PhysicalColumn.of("id", BasicType.LONG_TYPE,
22, false, null, "id"))
+ .build();
+ return CatalogTable.of(
+ TableIdentifier.of("catalog", "default", null, "dry_run"),
+ schema,
+ new HashMap<>(),
+ new ArrayList<>(),
+ null,
+ "catalog");
+ }
+
+ @FunctionalInterface
+ private interface Validation {
+ void run() throws Exception;
+ }
+}