This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12431-483aa86c096a2d72eba28ac85eb8d3a5c71e9295 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
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; + } +}
