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-12402-398fdfac5c6f48e894ce94c7b0ba7eeb74c96f7c in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit c12764c05bdfb1f8b2d908d0e34ab2d27bb6999e Author: Goutam Adwant <[email protected]> AuthorDate: Tue Sep 22 02:57:56 2026 +0000 [Fix][Connector-V2] Honor Redis named-user authentication (#12402) Signed-off-by: Goutam Adwant <[email protected]> --- docs/en/connectors/sink/Redis.md | 13 ++ docs/en/connectors/source/Redis.md | 12 +- .../introduction/concepts/incompatible-changes.md | 12 + docs/zh/connectors/sink/Redis.md | 12 + docs/zh/connectors/source/Redis.md | 10 +- .../introduction/concepts/incompatible-changes.md | 11 + .../seatunnel/redis/config/RedisBaseOptions.java | 11 +- .../seatunnel/redis/config/RedisParameters.java | 65 ++++-- .../connectors/seatunnel/redis/Redis7Test.java | 257 +++++++++++++++++++++ .../seatunnel/redis/RedisFactoryTest.java | 74 ++++++ .../seatunnel/e2e/connector/redis/Redis7IT.java | 46 ++++ .../connector/redis/RedisTestCaseTemplateIT.java | 2 +- .../src/test/resources/redis-named-user.conf | 50 ++++ 13 files changed, 555 insertions(+), 20 deletions(-) diff --git a/docs/en/connectors/sink/Redis.md b/docs/en/connectors/sink/Redis.md index eaa81845b0..760176e433 100644 --- a/docs/en/connectors/sink/Redis.md +++ b/docs/en/connectors/sink/Redis.md @@ -58,6 +58,19 @@ downloaded from Maven Central. | multi_table_sink_replica | int | No | 1 | Writer replica count for multi-table writes. | | common-options | config | No | - | Sink plugin common parameters. See [Sink Common Options](../common-options/sink-common-options.md). | +### Authentication + +In both `SINGLE` and `CLUSTER` mode, a nonblank `user` selects Redis ACL authentication +(`AUTH user auth`, Redis 6 or later). The connector does not create or modify ACL users. +Create the user and grant its required command and key permissions before starting the job, including +`INFO` for connector initialization, `SELECT` in `SINGLE` mode, and `CLUSTER SLOTS` for topology +discovery in `CLUSTER` mode. +The password is passed unchanged, including whitespace; omitted or empty `auth` is sent as an empty +password and only works if the ACL user accepts it (for example, a `nopass` user). + +If `user` is omitted, empty, or whitespace-only, nonblank `auth` uses password-only authentication +as the default user. If both options are omitted or blank, no authentication command is sent. + ## Write Rules ### key diff --git a/docs/en/connectors/source/Redis.md b/docs/en/connectors/source/Redis.md index 77a203d473..093b4e5cdb 100644 --- a/docs/en/connectors/source/Redis.md +++ b/docs/en/connectors/source/Redis.md @@ -232,11 +232,19 @@ redis data types, support `key` `string` `hash` `list` `set` `zset` ### user [string] -redis authentication user, you need it when you connect to an encrypted cluster +Redis ACL username (Redis 6 or later), supported in both `SINGLE` and `CLUSTER` mode. +When nonblank, the connector authenticates with `AUTH user auth`; it does not create or modify ACL users. +Create the user and grant the required command and key permissions before starting the job, including +`INFO` for connector initialization, `SELECT` in `SINGLE` mode, and `CLUSTER SLOTS` for topology +discovery in `CLUSTER` mode. +If `user` is omitted, empty, or whitespace-only, a nonblank `auth` uses password-only authentication +as the default user; otherwise no authentication command is sent. ### auth [string] -redis authentication password, you need it when you connect to an encrypted cluster +Redis authentication password. With a nonblank `user`, the password is passed unchanged, including +whitespace. An omitted or empty password is sent as an empty string and works only if that ACL user +accepts it (for example, a user configured with `nopass`). ### db_num [int] diff --git a/docs/en/introduction/concepts/incompatible-changes.md b/docs/en/introduction/concepts/incompatible-changes.md index 99181067e1..b90ad91efd 100644 --- a/docs/en/introduction/concepts/incompatible-changes.md +++ b/docs/en/introduction/concepts/incompatible-changes.md @@ -5,6 +5,18 @@ You need to check this document before you upgrade to related version. ## dev +### Redis Authentication + +- Redis sources and sinks now authenticate as the configured nonblank `user` in both `SINGLE` and + `CLUSTER` mode. Previously, `SINGLE` used password-only authentication followed by `ACL SETUSER`, + and `CLUSTER` ignored `user`. Connection setup no longer creates or modifies ACL users. +- Before upgrading, create the intended ACL user and grant its required command and key permissions, + including `INFO` for connector initialization, `SELECT` in `SINGLE` mode, and `CLUSTER SLOTS` for + topology discovery in `CLUSTER` mode. Set `auth` to that user's password. An omitted or empty password + is sent as an empty string when `user` is nonblank. +- To keep using the default user, remove `user` and retain `auth` when a password is required. + Named users require Redis 6 or later. Legacy configurations without a username remain unchanged. + ### RabbitMQ Connector - **Breaking Change: `amqps://` connections now verify broker certificates** diff --git a/docs/zh/connectors/sink/Redis.md b/docs/zh/connectors/sink/Redis.md index c430f69bb6..3c59c1c1ab 100644 --- a/docs/zh/connectors/sink/Redis.md +++ b/docs/zh/connectors/sink/Redis.md @@ -57,6 +57,18 @@ Redis 接收器连接器可以在批处理或流处理作业中把上游数据 | multi_table_sink_replica | int | 否 | 1 | 多表写入时的写入器副本数。 | | common-options | config | 否 | - | 接收器插件通用参数,详情请参考[接收器通用选项](../common-options/sink-common-options.md)。 | +### 认证 + +在 `SINGLE` 和 `CLUSTER` 模式下,非空白的 `user` 使用 Redis ACL 认证 +(`AUTH user auth`,需要 Redis 6 或更新版本)。连接器不会创建或修改 ACL 用户。 +启动作业前,请创建用户并授予所需的命令和键权限,包括初始化连接器所需的 `INFO`, +`SINGLE` 模式所需的 `SELECT`,以及 `CLUSTER` 模式下拓扑发现所需的 `CLUSTER SLOTS`。 +密码将原样传递,包括空白字符;省略 `auth` 或使用空字符串时将发送空密码, +仅当该 ACL 用户允许时才能成功认证(例如配置了 `nopass` 的用户)。 + +若省略 `user`,或其值为空字符串、仅包含空白字符,则非空白的 `auth` 将用于默认用户的密码认证。 +如果两个选项都省略或为空白,则不发送认证命令。 + ## 写入规则 ### key diff --git a/docs/zh/connectors/source/Redis.md b/docs/zh/connectors/source/Redis.md index b930d1c824..77f84c018d 100644 --- a/docs/zh/connectors/source/Redis.md +++ b/docs/zh/connectors/source/Redis.md @@ -224,11 +224,17 @@ redis 数据类型, 支持 `key` `string` `hash` `list` `set` `zset`。 ### user [string] -Redis 认证身份用户,当连接到加密集群时需要使用 +Redis ACL 用户名(需要 Redis 6 或更新版本),支持 `SINGLE` 和 `CLUSTER` 模式。 +当用户名非空白时,连接器通过 `AUTH user auth` 认证,不会创建或修改 ACL 用户。 +启动作业前,请创建用户并授予所需的命令和键权限,包括初始化连接器所需的 `INFO`, +`SINGLE` 模式所需的 `SELECT`,以及 `CLUSTER` 模式下拓扑发现所需的 `CLUSTER SLOTS`。 +若省略 `user`,或其值为空字符串、仅包含空白字符,则非空白的 `auth` 将用于默认用户的密码认证; +否则不发送认证命令。 ### auth [string] -Redis 认证密钥,当连接到加密集群时需要使用 +Redis 认证密码。当 `user` 非空白时,密码将原样传递,包括空白字符。 +省略密码或使用空字符串时,将发送空密码,仅当该 ACL 用户允许时才能成功认证(例如配置了 `nopass` 的用户)。 ### db_num [int] diff --git a/docs/zh/introduction/concepts/incompatible-changes.md b/docs/zh/introduction/concepts/incompatible-changes.md index 8a3b55c330..fb0a541081 100644 --- a/docs/zh/introduction/concepts/incompatible-changes.md +++ b/docs/zh/introduction/concepts/incompatible-changes.md @@ -4,6 +4,17 @@ ## dev +### Redis 认证 + +- Redis Source 和 Sink 现在会在 `SINGLE` 和 `CLUSTER` 模式下以非空白的 `user` 指定的用户认证。 + 此前,`SINGLE` 模式先使用仅密码认证,再执行 `ACL SETUSER`;`CLUSTER` 模式忽略 `user`。 + 连接初始化不再创建或修改 ACL 用户。 +- 升级前,请创建目标 ACL 用户并授予所需的命令和键权限,包括初始化连接器所需的 `INFO`, + `SINGLE` 模式所需的 `SELECT`,以及 `CLUSTER` 模式下拓扑发现所需的 `CLUSTER SLOTS`。 + 将 `auth` 设置为该用户的密码。当 `user` 非空白时,省略密码或使用空字符串将发送空密码。 +- 如需继续使用默认用户,请移除 `user`,并在需要密码时保留 `auth`。 + 命名用户需要 Redis 6 或更新版本;未配置用户名的旧配置行为保持不变。 + ### RabbitMQ Connector - **破坏性变更:`amqps://` 连接现在会校验 Broker 证书** diff --git a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java index 690f69fd92..8ee37aae4a 100644 --- a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java +++ b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisBaseOptions.java @@ -50,7 +50,11 @@ public class RedisBaseOptions extends ConnectorCommonOptions { .stringType() .noDefaultValue() .withDescription( - "redis authentication password, you need it when you connect to an encrypted cluster"); + "Redis authentication password for SINGLE and CLUSTER modes. With a nonblank user," + + " the password is passed unchanged; omitted or empty auth sends an empty" + + " password, which the ACL user must accept (for example, via nopass)." + + " With a blank user, nonblank auth authenticates as the default user;" + + " otherwise no authentication command is sent."); public static final Option<Integer> DB_NUM = Options.key("db_num") @@ -64,7 +68,10 @@ public class RedisBaseOptions extends ConnectorCommonOptions { .stringType() .noDefaultValue() .withDescription( - "redis authentication user, you need it when you connect to an encrypted cluster"); + "Redis ACL username (Redis 6 or later) for SINGLE and CLUSTER modes." + + " When nonblank, the connector authenticates with AUTH user auth" + + " without creating or modifying ACL users. Omitted, empty, or" + + " whitespace-only user retains default-user authentication behavior."); public static final Option<String> KEY_PATTERN = Options.key("keys") diff --git a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java index cac06d2b47..d8458a282b 100644 --- a/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java +++ b/seatunnel-connectors-v2/connector-redis/src/main/java/org/apache/seatunnel/connectors/seatunnel/redis/config/RedisParameters.java @@ -29,6 +29,7 @@ import org.apache.seatunnel.connectors.seatunnel.redis.exception.RedisConnectorE import lombok.Data; import lombok.extern.slf4j.Slf4j; import redis.clients.jedis.ConnectionPoolConfig; +import redis.clients.jedis.DefaultJedisClientConfig; import redis.clients.jedis.HostAndPort; import redis.clients.jedis.Jedis; import redis.clients.jedis.JedisCluster; @@ -144,11 +145,16 @@ public class RedisParameters implements Serializable { public RedisClient buildRedisClient() { Jedis jedis = this.buildJedis(); - this.redisVersion = extractRedisVersion(jedis); - if (mode.equals(RedisBaseOptions.RedisMode.SINGLE)) { - return new RedisSingleClient(this, jedis, redisVersion); - } else { - return new RedisClusterClient(this, jedis, redisVersion); + try { + this.redisVersion = extractRedisVersion(jedis); + if (mode.equals(RedisBaseOptions.RedisMode.SINGLE)) { + return new RedisSingleClient(this, jedis, redisVersion); + } else { + return new RedisClusterClient(this, jedis, redisVersion); + } + } catch (RuntimeException | Error failure) { + closeAfterFailure(jedis, failure); + throw failure; } } @@ -180,18 +186,27 @@ public class RedisParameters implements Serializable { "Did not get the expected redis_version from the jedis.info() method"); } + /** + * Uses named-user AUTH when user is nonblank; otherwise preserves password-only or no-auth + * behavior. Authentication selects an existing identity, never administers users via ACL + * SETUSER. + */ public Jedis buildJedis() { switch (mode) { case SINGLE: Jedis jedis = new Jedis(host, port); - if (StringUtils.isNotBlank(auth)) { - jedis.auth(auth); + try { + if (StringUtils.isNotBlank(user)) { + jedis.auth(user, StringUtils.defaultString(auth)); + } else if (StringUtils.isNotBlank(auth)) { + jedis.auth(auth); + } + jedis.select(dbNum); + return jedis; + } catch (RuntimeException | Error failure) { + closeAfterFailure(jedis, failure); + throw failure; } - if (StringUtils.isNotBlank(user)) { - jedis.aclSetUser(user); - } - jedis.select(dbNum); - return jedis; case CLUSTER: HashSet<HostAndPort> nodes = new HashSet<>(); for (String redisNode : redisNodes) { @@ -202,7 +217,19 @@ public class RedisParameters implements Serializable { } ConnectionPoolConfig connectionPoolConfig = new ConnectionPoolConfig(); JedisCluster jedisCluster; - if (StringUtils.isNotBlank(auth)) { + if (StringUtils.isNotBlank(user)) { + jedisCluster = + new JedisCluster( + nodes, + DefaultJedisClientConfig.builder() + .user(user) + .password(StringUtils.defaultString(auth)) + .connectionTimeoutMillis(JedisCluster.DEFAULT_TIMEOUT) + .socketTimeoutMillis(JedisCluster.DEFAULT_TIMEOUT) + .build(), + JedisCluster.DEFAULT_MAX_ATTEMPTS, + connectionPoolConfig); + } else if (StringUtils.isNotBlank(auth)) { jedisCluster = new JedisCluster( nodes, @@ -222,4 +249,16 @@ public class RedisParameters implements Serializable { CommonErrorCode.OPERATION_NOT_SUPPORTED, "Not support this redis mode"); } } + + /** + * Releases a connection after failed initialization, including an Error, before the caller + * rethrows the original failure. Cleanup failures are suppressed to preserve that failure. + */ + private void closeAfterFailure(Jedis jedis, Throwable failure) { + try { + jedis.close(); + } catch (RuntimeException | Error closeFailure) { + failure.addSuppressed(closeFailure); + } + } } diff --git a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java index a2f008e799..87b7e3c3a0 100644 --- a/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java +++ b/seatunnel-connectors-v2/connector-redis/src/test/java/org/apache/seatunnel/connectors/seatunnel/redis/Redis7Test.java @@ -16,16 +16,273 @@ */ package org.apache.seatunnel.connectors.seatunnel.redis; +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.redis.client.RedisClient; +import org.apache.seatunnel.connectors.seatunnel.redis.config.JedisWrapper; +import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisBaseOptions; import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisContainerInfo; +import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; import org.junit.jupiter.api.condition.DisabledOnOs; import org.junit.jupiter.api.condition.OS; +import org.testcontainers.containers.GenericContainer; + +import redis.clients.jedis.Jedis; +import redis.clients.jedis.exceptions.JedisDataException; + +import java.net.InetAddress; +import java.net.UnknownHostException; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.BooleanSupplier; +import java.util.stream.IntStream; @DisabledOnOs( value = OS.WINDOWS, disabledReason = "There is no docker environment on the windows test system") public class Redis7Test extends RedisTemplateTest { + @Test + public void namedUserAuthentication() { + String user = "seatunnel_named"; + // Sharing the default user's password must not silently select the default identity. + jedis.aclSetUser( + user, + "reset", + "on", + ">" + password, + "~auth:*", + "+select", + "+info", + "+get", + "+set", + "+acl|whoami"); + List<String> aclBefore = jedis.aclList(); + try { + RedisParameters parameters = connectionParameters(user, password); + parameters.setDbNum(2); + try (Jedis connection = parameters.buildJedis()) { + Assertions.assertEquals(user, connection.aclWhoAmI()); + connection.set("auth:selected", "value"); + Assertions.assertThrows( + JedisDataException.class, () -> connection.set("outside:auth", "value")); + Assertions.assertThrows( + JedisDataException.class, () -> connection.aclSetUser("unexpected")); + } + jedis.select(2); + Assertions.assertEquals("value", jedis.get("auth:selected")); + jedis.del("auth:selected"); + jedis.select(0); + RedisClient client = parameters.buildRedisClient(); + try { + Assertions.assertEquals(7, parameters.getRedisVersion()); + } finally { + client.close(); + } + Assertions.assertEquals(aclBefore, jedis.aclList()); + jedis.aclSetUser(user, "resetpass", ">named-password"); + try (Jedis connection = connectionParameters(user, "named-password").buildJedis()) { + Assertions.assertEquals(user, connection.aclWhoAmI()); + } + } finally { + jedis.select(0); + jedis.aclDelUser(user); + } + } + + @Test + public void namedUserMissingBlankAndWrongPassword() { + String user = "seatunnel_nopass"; + jedis.aclSetUser(user, "reset", "on", "nopass", "+select", "+acl|whoami"); + try { + for (String auth : new String[] {null, "", " "}) { + try (Jedis connection = connectionParameters(user, auth).buildJedis()) { + Assertions.assertEquals(user, connection.aclWhoAmI()); + } + } + jedis.aclSetUser(user, "resetpass", ">named-password"); + for (String auth : new String[] {null, "", " ", "wrong-password"}) { + Assertions.assertThrows( + JedisDataException.class, + () -> connectionParameters(user, auth).buildJedis()); + } + jedis.aclSetUser(user, "resetpass", "> "); + try (Jedis connection = connectionParameters(user, " ").buildJedis()) { + Assertions.assertEquals(user, connection.aclWhoAmI()); + } + } finally { + jedis.aclDelUser(user); + } + } + + @Test + public void legacyDefaultUserAuthentication() { + for (String user : new String[] {null, "", " "}) { + try (Jedis connection = connectionParameters(user, password).buildJedis()) { + Assertions.assertEquals("default", connection.aclWhoAmI()); + } + } + jedis.aclSetUser("default", "nopass"); + try { + for (String user : new String[] {null, "", " "}) { + for (String auth : new String[] {null, "", " "}) { + try (Jedis connection = connectionParameters(user, auth).buildJedis()) { + Assertions.assertEquals("default", connection.aclWhoAmI()); + } + } + } + List<String> aclBefore = jedis.aclList(); + Assertions.assertThrows( + JedisDataException.class, + () -> connectionParameters("nonexistent", null).buildJedis()); + Assertions.assertEquals(aclBefore, jedis.aclList()); + } finally { + jedis.aclSetUser("default", "resetpass", ">" + password); + } + } + + @Test + public void failedInitializationClosesConnections() throws InterruptedException { + String user = "seatunnel_denied"; + jedis.aclSetUser(user, "reset", "on", ">named-password", "+select"); + try { + assertConnectionCountUnchanged(() -> connectionParameters(user, "wrong").buildJedis()); + RedisParameters parameters = connectionParameters(user, "named-password"); + parameters.setDbNum(-1); + assertConnectionCountUnchanged(parameters::buildJedis); + parameters.setDbNum(0); + assertConnectionCountUnchanged(parameters::buildRedisClient); + jedis.aclSetUser(user, "-select"); + assertConnectionCountUnchanged(parameters::buildJedis); + } finally { + jedis.aclDelUser(user); + } + } + + private void assertConnectionCountUnchanged(Runnable connect) throws InterruptedException { + int connections = jedis.clientList().split("\n").length; + Assertions.assertThrows(JedisDataException.class, connect::run); + awaitCondition(() -> connections == jedis.clientList().split("\n").length); + } + + @Test + public void namedClusterAuthentication() throws InterruptedException, UnknownHostException { + // A real cluster-enabled server owns all slots; no external cluster is required. + try (GenericContainer<?> cluster = + new GenericContainer<>("redis:7") + .withExposedPorts(6379) + .withCommand( + "redis-server", + "--cluster-enabled", + "yes", + "--save", + "", + "--appendonly", + "no")) { + cluster.start(); + try (Jedis admin = new Jedis(cluster.getHost(), cluster.getFirstMappedPort())) { + admin.configSet( + "cluster-announce-ip", + InetAddress.getByName(cluster.getHost()).getHostAddress()); + admin.configSet("cluster-announce-port", cluster.getFirstMappedPort().toString()); + admin.clusterAddSlots(IntStream.range(0, 16384).toArray()); + awaitCondition(() -> admin.clusterInfo().contains("cluster_state:ok")); + admin.aclSetUser( + "seatunnel_cluster", + "reset", + "on", + ">named-password", + "~auth:*", + "+cluster|slots", + "+info", + "+get", + "+set", + "+acl|whoami"); + RedisParameters parameters = + connectionParameters("seatunnel_cluster", "named-password"); + parameters.setMode(RedisBaseOptions.RedisMode.CLUSTER); + parameters.setRedisNodes( + Collections.singletonList( + cluster.getHost() + ":" + cluster.getFirstMappedPort())); + List<String> aclBefore = admin.aclList(); + try (JedisWrapper connection = (JedisWrapper) parameters.buildJedis()) { + for (String node : connection.getClusterNodes().keySet()) { + Assertions.assertEquals( + "seatunnel_cluster", connection.getJedis(node).aclWhoAmI()); + } + connection.set("auth:cluster", "value"); + Assertions.assertEquals("value", connection.get("auth:cluster")); + } + RedisClient client = parameters.buildRedisClient(); + client.close(); + Assertions.assertEquals(aclBefore, admin.aclList()); + + parameters.setUser(""); + parameters.setAuth(""); + assertClusterIdentity(parameters, "default"); + admin.aclSetUser("default", "resetpass", ">named-password"); + parameters.setAuth("named-password"); + assertClusterIdentity(parameters, "default"); + parameters.setUser("seatunnel_cluster"); + assertClusterIdentity(parameters, "seatunnel_cluster"); + + parameters.setAuth("wrong-password"); + int beforeAuthenticationFailure = admin.clientList().split("\n").length; + Assertions.assertThrows(JedisDataException.class, parameters::buildJedis); + awaitCondition( + () -> beforeAuthenticationFailure == admin.clientList().split("\n").length); + parameters.setAuth("named-password"); + admin.aclSetUser("seatunnel_cluster", "-info"); + int connections = admin.clientList().split("\n").length; + Assertions.assertThrows(RuntimeException.class, parameters::buildRedisClient); + awaitCondition(() -> connections == admin.clientList().split("\n").length); + + admin.aclSetUser("seatunnel_cluster", "nopass"); + for (String auth : new String[] {null, "", " "}) { + parameters.setAuth(auth); + assertClusterIdentity(parameters, "seatunnel_cluster"); + } + } + } + } + + private void assertClusterIdentity(RedisParameters parameters, String user) { + try (JedisWrapper connection = (JedisWrapper) parameters.buildJedis()) { + for (String node : connection.getClusterNodes().keySet()) { + Assertions.assertEquals(user, connection.getJedis(node).aclWhoAmI()); + } + } + } + + private void awaitCondition(BooleanSupplier condition) throws InterruptedException { + for (int attempt = 0; attempt < 100; attempt++) { + if (condition.getAsBoolean()) { + return; + } + Thread.sleep(50); + } + Assertions.assertTrue(condition.getAsBoolean()); + } + + private RedisParameters connectionParameters(String user, String auth) { + Map<String, Object> config = new HashMap<>(); + config.put("host", redisContainer.getHost()); + config.put("port", redisContainer.getFirstMappedPort()); + if (user != null) { + config.put("user", user); + } + if (auth != null) { + config.put("auth", auth); + } + RedisParameters parameters = new RedisParameters(); + parameters.buildConnectionConfig(ReadonlyConfig.fromMap(config)); + return parameters; + } + @Override public RedisContainerInfo getRedisContainerInfo() { return new RedisContainerInfo("redis-e2e", 6379, "SeaTunnel", "redis:7"); 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 67fb4a8989..f1b6d3850e 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 @@ -29,6 +29,8 @@ import org.apache.seatunnel.api.table.catalog.TableSchema; 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.config.RedisBaseOptions; +import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisParameters; import org.apache.seatunnel.connectors.seatunnel.redis.sink.RedisSink; import org.apache.seatunnel.connectors.seatunnel.redis.sink.RedisSinkFactory; import org.apache.seatunnel.connectors.seatunnel.redis.source.RedisSourceFactory; @@ -38,6 +40,11 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.Arguments; import org.junit.jupiter.params.provider.MethodSource; +import org.mockito.MockedConstruction; +import org.mockito.Mockito; + +import redis.clients.jedis.Jedis; +import redis.clients.jedis.exceptions.JedisDataException; import java.io.IOException; import java.util.ArrayList; @@ -53,6 +60,73 @@ class RedisFactoryTest { private static final OptionRule SOURCE_RULE = new RedisSourceFactory().optionRule(); private static final OptionRule SINK_RULE = new RedisSinkFactory().optionRule(); + @Test + void namedAuthenticationDoesNotAdministerUsers() { + try (MockedConstruction<Jedis> construction = Mockito.mockConstruction(Jedis.class)) { + RedisParameters parameters = singleConnectionParameters(); + parameters.setUser("named-user"); + parameters.setAuth("named-password"); + parameters.setDbNum(3); + Jedis connection = parameters.buildJedis(); + Assertions.assertSame(construction.constructed().get(0), connection); + Mockito.verify(connection).auth("named-user", "named-password"); + Mockito.verify(connection).select(3); + Mockito.verifyNoMoreInteractions(connection); + connection.close(); + } + } + + @Test + void authenticationFailureClosesConnectionAndPreservesCause() { + JedisDataException failure = new JedisDataException("Authentication failed"); + RuntimeException closeFailure = new RuntimeException("Close failed"); + try (MockedConstruction<Jedis> construction = + Mockito.mockConstruction( + Jedis.class, + (connection, context) -> { + Mockito.when(connection.auth("named-user", "wrong")).thenThrow(failure); + Mockito.doThrow(closeFailure).when(connection).close(); + })) { + RedisParameters parameters = singleConnectionParameters(); + parameters.setUser("named-user"); + parameters.setAuth("wrong"); + Assertions.assertSame( + failure, + Assertions.assertThrows(JedisDataException.class, parameters::buildJedis)); + Mockito.verify(construction.constructed().get(0)).close(); + Assertions.assertArrayEquals(new Throwable[] {closeFailure}, failure.getSuppressed()); + } + } + + @Test + void versionInitializationFailureClosesConnection() { + for (String info : Arrays.asList("redis_version:not-a-version", "no-version")) { + RedisParameters parameters = Mockito.spy(singleConnectionParameters()); + Jedis connection = Mockito.mock(Jedis.class); + Mockito.doReturn(connection).when(parameters).buildJedis(); + Mockito.when(connection.info()).thenReturn(info); + Assertions.assertThrows(RuntimeException.class, parameters::buildRedisClient); + Mockito.verify(connection).close(); + } + RedisParameters parameters = Mockito.spy(singleConnectionParameters()); + Jedis connection = Mockito.mock(Jedis.class); + Mockito.doReturn(connection).when(parameters).buildJedis(); + JedisDataException denied = new JedisDataException("INFO denied"); + Mockito.when(connection.info()).thenThrow(denied); + Assertions.assertSame( + denied, + Assertions.assertThrows(JedisDataException.class, parameters::buildRedisClient)); + Mockito.verify(connection).close(); + } + + private RedisParameters singleConnectionParameters() { + RedisParameters parameters = new RedisParameters(); + parameters.setHost("localhost"); + parameters.setPort(6379); + parameters.setMode(RedisBaseOptions.RedisMode.SINGLE); + return parameters; + } + @Test void optionRule() { Assertions.assertNotNull(SOURCE_RULE); diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java index 7e065837ff..33c6516539 100644 --- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/Redis7IT.java @@ -17,12 +17,58 @@ package org.apache.seatunnel.e2e.connector.redis; import org.apache.seatunnel.connectors.seatunnel.redis.config.RedisContainerInfo; +import org.apache.seatunnel.e2e.common.container.TestContainer; +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.TestTemplate; import org.junit.jupiter.api.parallel.ResourceLock; +import org.testcontainers.containers.Container; + +import java.io.IOException; +import java.util.Collections; +import java.util.List; @ResourceLock("redis-standalone-e2e") public class Redis7IT extends RedisTestCaseTemplateIT { + @TestTemplate + public void testNamedUserSourceAndSink(TestContainer container) + throws IOException, InterruptedException { + jedis.aclSetUser( + "seatunnel_reader", + "reset", + "on", + ">reader-password", + "~acl:source:*", + "+select", + "+info", + "+scan", + "+type", + "+get", + "+mget"); + jedis.aclSetUser( + "seatunnel_writer", + "reset", + "on", + ">writer-password", + "~acl:result", + "+select", + "+info", + "+lpush"); + List<String> aclBefore = jedis.aclList(); + jedis.set("acl:source:1", "{\"value\":\"named-user\"}"); + try { + Container.ExecResult result = container.executeJob("/redis-named-user.conf"); + Assertions.assertEquals(0, result.getExitCode()); + Assertions.assertEquals( + Collections.singletonList("named-user"), jedis.lrange("acl:result", 0, -1)); + Assertions.assertEquals(aclBefore, jedis.aclList()); + } finally { + jedis.del("acl:source:1", "acl:result"); + jedis.aclDelUser("seatunnel_reader", "seatunnel_writer"); + } + } + @Override public RedisContainerInfo getRedisContainerInfo() { return new RedisContainerInfo("redis-e2e", 6379, "SeaTunnel", "redis:7"); diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java index bb1650f93d..e6aabdae18 100644 --- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/java/org/apache/seatunnel/e2e/connector/redis/RedisTestCaseTemplateIT.java @@ -89,7 +89,7 @@ public abstract class RedisTestCaseTemplateIT extends TestSuiteBase implements T private GenericContainer<?> redisContainer; - private Jedis jedis; + protected Jedis jedis; @BeforeAll @Override diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/resources/redis-named-user.conf b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/resources/redis-named-user.conf new file mode 100644 index 0000000000..2f88da84e2 --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-redis-e2e/src/test/resources/redis-named-user.conf @@ -0,0 +1,50 @@ +# +# 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. +# + +env { + parallelism = 1 + job.mode = "BATCH" +} + +source { + Redis { + host = "redis-e2e" + port = 6379 + user = "seatunnel_reader" + auth = "reader-password" + keys = "acl:source:*" + data_type = string + format = json + schema = { + fields { + value = string + } + } + } +} + +sink { + Redis { + host = "redis-e2e" + port = 6379 + user = "seatunnel_writer" + auth = "writer-password" + key = "acl:result" + data_type = list + value_field = "value" + } +}
