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-11576-0d9f9e2303da0b40b20e55fd341d564757e86dc0 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 38a37104e59ee3317cf0bc299c9bfee22149391c Author: Jast <[email protected]> AuthorDate: Tue Sep 15 04:16:29 2026 +0000 [Fix][Connector-V2] Support secure passive RabbitMQ connections (#11576) Co-authored-by: zhangshenghang <[email protected]> Co-authored-by: zhangshenghang <[email protected]> --- docs/en/connectors/sink/Rabbitmq.md | 20 ++++++ docs/en/connectors/source/Rabbitmq.md | 20 ++++++ .../introduction/concepts/incompatible-changes.md | 13 ++++ docs/zh/connectors/sink/Rabbitmq.md | 20 ++++++ docs/zh/connectors/source/Rabbitmq.md | 20 ++++++ .../introduction/concepts/incompatible-changes.md | 11 +++ .../seatunnel/rabbitmq/client/RabbitmqClient.java | 55 +++++++++++--- .../rabbitmq/config/RabbitmqBaseOptions.java | 20 ++++++ .../seatunnel/rabbitmq/config/RabbitmqConfig.java | 25 ++++++- .../exception/RabbitmqConnectorErrorCode.java | 3 +- .../rabbitmq/sink/RabbitmqSinkFactory.java | 4 ++ .../rabbitmq/source/RabbitmqSourceFactory.java | 4 ++ .../seatunnel/rabbitmq/RabbitmqFactoryTest.java | 42 +++++++++++ .../rabbitmq/client/RabbitmqClientTest.java | 83 ++++++++++++++++++++++ 14 files changed, 328 insertions(+), 12 deletions(-) diff --git a/docs/en/connectors/sink/Rabbitmq.md b/docs/en/connectors/sink/Rabbitmq.md index 70a507bb9a..bab1147b37 100644 --- a/docs/en/connectors/sink/Rabbitmq.md +++ b/docs/en/connectors/sink/Rabbitmq.md @@ -32,6 +32,8 @@ Used to write data to RabbitMQ queues. | protobuf_schema | string | no | - | | protobuf_message_name | string | no | - | | url | string | no | - | +| uri | string | no | - | +| ssl | boolean | no | false | | routing_key | string | no | - | | exchange | string | no | - | | network_recovery_interval | int | no | - | @@ -43,6 +45,7 @@ Used to write data to RabbitMQ queues. | durable | boolean | no | true | | exclusive | boolean | no | false | | auto_delete | boolean | no | false | +| passive | boolean | no | false | ### host [string] @@ -70,6 +73,16 @@ the password to use when connecting to the broker convenience method for setting the fields in an AMQP URI: host, port, username, password and virtual host +### uri [string] + +Legacy alias for `url`. Configure only one of `url` and `uri`. + +### ssl [boolean] + +Enables SSL/TLS for host-and-port configuration. Use `url` with an `amqps://` URI when the URI itself supplies the connection settings. + +When `url` uses an `amqps://` URI, the broker certificate is verified against the JVM trust store with hostname verification enabled. Connections that previously relied on the implicit trust-all behavior with self-signed or private-CA certificates must import the broker certificate into the trust store, or they will fail to connect. + ### queue_name [string] the queue to write the message to. If `routing_key` is not configured, the connector publishes messages to this queue through the default exchange. @@ -135,10 +148,17 @@ Sink plugin common parameters, please refer to [Sink Common Options](../common-o - true: The queue will be deleted automatically when the last consumer unsubscribes. - false: The queue will not be automatically deleted. +### passive + +- false: Declare the queue with the configured durable, exclusive, and auto-delete settings. +- true: Verify that the queue already exists without creating or modifying it. Use this for accounts that can publish but cannot declare queues. + ## Configuration Notes - If you configure `username`, you must also configure `password`, and vice versa. +- Configure only one of `url` and `uri`. `uri` is retained for existing configurations; use `url` in new configurations. +- Set `ssl = true` when connecting to an AMQPS endpoint with `host` and `port` settings. - `host`, `port`, `virtual_host`, and `queue_name` are required connector options. `url` can additionally provide the AMQP URI used by the RabbitMQ client. - `durable`, `exclusive`, and `auto_delete` are used when the connector declares the target queue. - When `format` is `protobuf`, configure both `protobuf_schema` and `protobuf_message_name`. diff --git a/docs/en/connectors/source/Rabbitmq.md b/docs/en/connectors/source/Rabbitmq.md index 3cd9d38bbe..1c50b4b0ba 100644 --- a/docs/en/connectors/source/Rabbitmq.md +++ b/docs/en/connectors/source/Rabbitmq.md @@ -46,6 +46,8 @@ The source must be non-parallel (parallelism set to 1) in order to achieve exact | protobuf_schema | string | no | - | | protobuf_message_name | string | no | - | | url | string | no | - | +| uri | string | no | - | +| ssl | boolean | no | false | | routing_key | string | no | - | | exchange | string | no | - | | network_recovery_interval | int | no | - | @@ -62,6 +64,7 @@ The source must be non-parallel (parallelism set to 1) in order to achieve exact | durable | boolean | no | true | | exclusive | boolean | no | false | | auto_delete | boolean | no | false | +| passive | boolean | no | false | ### host [string] @@ -89,6 +92,16 @@ the password to use when connecting to the broker convenience method for setting the fields in an AMQP URI: host, port, username, password and virtual host +### uri [string] + +Legacy alias for `url`. Configure only one of `url` and `uri`. + +### ssl [boolean] + +Enables SSL/TLS for host-and-port configuration. Use `url` with an `amqps://` URI when the URI itself supplies the connection settings. + +When `url` uses an `amqps://` URI, the broker certificate is verified against the JVM trust store with hostname verification enabled. Connections that previously relied on the implicit trust-all behavior with self-signed or private-CA certificates must import the broker certificate into the trust store, or they will fail to connect. + ### queue_name [string] the queue to consume messages from. *Note: Required if `tables_configs` is not configured.* @@ -186,6 +199,11 @@ Source plugin common parameters, please refer to [Source Common Options](../comm - true: The queue will be deleted automatically when the last consumer unsubscribes. - false: The queue will not be automatically deleted. +### passive + +- false: Declare the queue with the configured durable, exclusive, and auto-delete settings. +- true: Verify that the queue already exists without creating or modifying it. Use this for consumer accounts without queue-declaration permission. + ## Migration Guide & Configuration Rules If you are upgrading from a previous version that only supported single-table reads, your existing configuration will work without any changes. @@ -197,6 +215,8 @@ If you are upgrading from a previous version that only supported single-table re - In multi-table mode, put each queue's `schema` inside its own `tables_configs` item. - When `format` is `protobuf`, configure both `protobuf_schema` and `protobuf_message_name` at the same level as the queue configuration. - If you configure `username`, you must also configure `password`, and vice versa. +- Configure only one of `url` and `uri`. `uri` is retained for existing configurations; use `url` in new configurations. +- Set `ssl = true` when connecting to an AMQPS endpoint with `host` and `port` settings. - `host` and `port` are always required. `virtual_host` is optional unless your RabbitMQ deployment requires a non-default virtual host. ## Example diff --git a/docs/en/introduction/concepts/incompatible-changes.md b/docs/en/introduction/concepts/incompatible-changes.md index 0dc9cb8392..ebe8cebd76 100644 --- a/docs/en/introduction/concepts/incompatible-changes.md +++ b/docs/en/introduction/concepts/incompatible-changes.md @@ -5,6 +5,19 @@ You need to check this document before you upgrade to related version. ## dev +### RabbitMQ Connector + +- **Breaking Change: `amqps://` connections now verify broker certificates** + - **Affected component**: `seatunnel-connectors-v2/connector-rabbitmq` + - **Description**: Previously, connecting with an `amqps://` `url`/`uri` implicitly installed a + trust-all trust manager without hostname verification. Certificate verification is now + enforced for `amqps://` connections, consistent with the `ssl = true` host/port path. + - **Impact**: Jobs that connect with `amqps://` URLs to brokers using self-signed or private-CA + certificates will fail to connect after upgrading. + - **Migration Guide**: Import the broker certificate (or your private CA chain) into the JVM + trust store of the SeaTunnel runtime, or switch to the `host`/`port` + `ssl = true` + configuration with a properly configured trust store. + ### Zeta REST Pagination Parameter Validation - **Behavior change: `page` and `rows` are validated on paginated endpoints** diff --git a/docs/zh/connectors/sink/Rabbitmq.md b/docs/zh/connectors/sink/Rabbitmq.md index e8028f7586..12b9222797 100644 --- a/docs/zh/connectors/sink/Rabbitmq.md +++ b/docs/zh/connectors/sink/Rabbitmq.md @@ -33,6 +33,8 @@ import ChangeLog from '../changelog/connector-rabbitmq.md'; | protobuf_schema | string | 否 | - | | protobuf_message_name | string | 否 | - | | url | string | 否 | - | +| uri | string | 否 | - | +| ssl | boolean | 否 | false | | routing_key | string | 否 | - | | exchange | string | 否 | - | | network_recovery_interval | int | 否 | - | @@ -43,6 +45,7 @@ import ChangeLog from '../changelog/connector-rabbitmq.md'; | durable | boolean | 否 | true | | exclusive | boolean | 否 | false | | auto_delete | boolean | 否 | false | +| passive | boolean | 否 | false | | common-options | | 否 | - | ### host [string] @@ -71,6 +74,16 @@ virtual host,连接 broker 使用的 vhost 设置host、port、username、password和virtual host的简便方式。 +### uri [string] + +`url` 的兼容别名。`url` 和 `uri` 只能配置一个。 + +### ssl [boolean] + +使用 `host` 和 `port` 配置连接时启用 SSL/TLS。若 URI 本身提供连接信息,请使用 `amqps://` 开头的 `url`。 + +当 `url` 使用 `amqps://` 时,将按 JVM 信任库校验 Broker 证书并启用主机名校验。此前依赖隐式信任所有证书、使用自签名或私有 CA 证书的连接,需要将 Broker 证书导入信任库,否则将无法建立连接。 + ### queue_name [string] 数据写入的队列名。如果没有配置 `routing_key`,连接器会通过默认 exchange 将消息直接写入该队列。 @@ -110,6 +123,11 @@ virtual host,连接 broker 使用的 vhost - true:队列将在最后一个消费者取消订阅时自动删除。 - false:队列不会自动删除。 +### passive [boolean] + +- false:按已配置的 durable、exclusive 和 auto_delete 参数声明队列。 +- true:只校验队列已存在,不创建或修改队列。适用于可发布但没有队列声明权限的账号。 + ### network_recovery_interval [int] 自动恢复需等待多长时间才尝试重连,单位为毫秒。 @@ -139,6 +157,8 @@ Sink插件常用参数,请参考[Sink常用选项](../common-options/sink-comm ## 配置说明 - 如果配置了 `username`,也必须配置 `password`,反过来也一样。 +- `url` 和 `uri` 只能配置一个。`uri` 为兼容已有配置保留,新配置请使用 `url`。 +- 使用 `host` 和 `port` 连接 AMQPS 端点时,请设置 `ssl = true`。 - `host`、`port`、`virtual_host` 和 `queue_name` 是连接器必填项。`url` 可额外提供 RabbitMQ 客户端使用的 AMQP URI。 - `durable`、`exclusive` 和 `auto_delete` 用于连接器声明目标队列。 - 当 `format` 为 `protobuf` 时,需要同时配置 `protobuf_schema` 和 `protobuf_message_name`。 diff --git a/docs/zh/connectors/source/Rabbitmq.md b/docs/zh/connectors/source/Rabbitmq.md index 1be0f4ef75..5975930117 100644 --- a/docs/zh/connectors/source/Rabbitmq.md +++ b/docs/zh/connectors/source/Rabbitmq.md @@ -46,6 +46,8 @@ import ChangeLog from '../changelog/connector-rabbitmq.md'; | protobuf_schema | string | 否 | - | 当 format 为 protobuf 时生效,用于解析消息体的 Protobuf Schema | | protobuf_message_name | string | 否 | - | 当 format 为 protobuf 时生效,指定要解析的 Protobuf Message 名称 | | url | string | 否 | - | 便捷方法,用于设置 AMQP URI 中的字段:主机、端口、用户名、密码和虚拟主机 | +| uri | string | 否 | - | `url` 的兼容别名 | +| ssl | boolean | 否 | false | 使用 host 和 port 连接时是否启用 SSL/TLS | | routing_key | string | 否 | - | RabbitMQ 共享配置中的可选路由键 | | exchange | string | 否 | - | RabbitMQ 共享配置中的可选 exchange | | network_recovery_interval | int | 否 | - | 自动恢复在尝试重新连接之前等待多长时间(毫秒) | @@ -61,6 +63,7 @@ import ChangeLog from '../changelog/connector-rabbitmq.md'; | durable | boolean | 否 | true | 队列是否在服务器重启时保留 | | exclusive | boolean | 否 | false | 队列是否仅由当前连接使用 | | auto_delete | boolean | 否 | false | 队列是否在最后一个消费者取消订阅时自动删除 | +| passive | boolean | 否 | false | 是否只校验已有队列而不声明或创建队列 | | common-options | | 否 | - | 源插件通用参数 | ### host [string] @@ -89,6 +92,16 @@ import ChangeLog from '../changelog/connector-rabbitmq.md'; 便捷方法,用于设置 AMQP URI 中的字段:主机、端口、用户名、密码和虚拟主机 +### uri [string] + +`url` 的兼容别名。`url` 和 `uri` 只能配置一个。 + +### ssl [boolean] + +使用 `host` 和 `port` 配置连接时启用 SSL/TLS。若 URI 本身提供连接信息,请使用 `amqps://` 开头的 `url`。 + +当 `url` 使用 `amqps://` 时,将按 JVM 信任库校验 Broker 证书并启用主机名校验。此前依赖隐式信任所有证书、使用自签名或私有 CA 证书的连接,需要将 Broker 证书导入信任库,否则将无法建立连接。 + ### queue_name [string] 要消费消息的队列。*注意:如果未配置 `tables_configs`,则为必填项。* @@ -184,6 +197,11 @@ RabbitMQ 共享配置中的可选 exchange。普通队列消费不需要配置 - true:队列将在最后一个消费者取消订阅时自动删除。 - false:队列不会自动删除。 +### passive + +- false:按已配置的 durable、exclusive 和 auto_delete 参数声明队列。 +- true:只校验队列已存在,不创建或修改队列。适用于没有队列声明权限的消费者账号。 + ## 迁移指南与配置规则 如果您从仅支持单表读取的早期版本升级,您现有的配置无需任何更改即可正常工作。 @@ -195,6 +213,8 @@ RabbitMQ 共享配置中的可选 exchange。普通队列消费不需要配置 - 多表模式下,每个队列自己的 `schema` 应放在对应的 `tables_configs` 条目里。 - 当 `format` 为 `protobuf` 时,需要在队列配置所在层级同时配置 `protobuf_schema` 和 `protobuf_message_name`。 - 如果配置了 `username`,也必须配置 `password`,反过来也一样。 +- `url` 和 `uri` 只能配置一个。`uri` 为兼容已有配置保留,新配置请使用 `url`。 +- 使用 `host` 和 `port` 连接 AMQPS 端点时,请设置 `ssl = true`。 - `host` 和 `port` 总是必填。`virtual_host` 是可选项,除非您的 RabbitMQ 环境要求使用非默认虚拟主机。 ## 示例 diff --git a/docs/zh/introduction/concepts/incompatible-changes.md b/docs/zh/introduction/concepts/incompatible-changes.md index 21a5217a28..ea8fc9d80b 100644 --- a/docs/zh/introduction/concepts/incompatible-changes.md +++ b/docs/zh/introduction/concepts/incompatible-changes.md @@ -4,6 +4,17 @@ ## dev +### RabbitMQ Connector + +- **破坏性变更:`amqps://` 连接现在会校验 Broker 证书** + - **影响范围**:`seatunnel-connectors-v2/connector-rabbitmq` + - **变更说明**:此前使用 `amqps://` 的 `url`/`uri` 建立连接时,会隐式启用“信任所有证书”的 + TrustManager 且不校验主机名。现在 `amqps://` 连接会强制校验证书,与 `ssl = true` 的 + host/port 路径行为保持一致。 + - **影响**:使用自签名或私有 CA 证书的 Broker,升级后通过 `amqps://` 建立的连接将失败。 + - **迁移指南**:将 Broker 证书(或私有 CA 证书链)导入 SeaTunnel 运行时的 JVM 信任库,或改用 + `host`/`port` + `ssl = true` 配置并正确设置信任库。 + ### Zeta REST 分页参数校验 - **行为变更:分页接口开始校验 `page` 与 `rows`** diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java index 85c6700678..af7ca2ef62 100644 --- a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClient.java @@ -29,6 +29,8 @@ import com.rabbitmq.client.ConnectionFactory; import com.rabbitmq.client.DefaultConsumer; import lombok.extern.slf4j.Slf4j; +import javax.net.ssl.SSLContext; + import java.io.IOException; import java.net.URISyntaxException; import java.security.KeyManagementException; @@ -38,6 +40,7 @@ import java.util.concurrent.TimeoutException; import static org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.CLOSE_CONNECTION_FAILED; import static org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.CREATE_RABBITMQ_CLIENT_FAILED; +import static org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.ILLEGAL_CONFIG; import static org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.INIT_SSL_CONTEXT_FAILED; import static org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.PARSE_URI_FAILED; import static org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.SEND_MESSAGE_FAILED; @@ -105,7 +108,9 @@ public class RabbitmqClient implements AutoCloseable { try { factory.setUri(config.getUri()); } catch (URISyntaxException e) { - throw new RabbitmqConnectorException(PARSE_URI_FAILED, e); + throw new RabbitmqConnectorException( + PARSE_URI_FAILED, + "Failed to parse RabbitMQ uri. Check the URI syntax without exposing credentials."); } catch (KeyManagementException e) { // this should never happen throw new RabbitmqConnectorException(INIT_SSL_CONTEXT_FAILED, e); @@ -113,6 +118,9 @@ public class RabbitmqClient implements AutoCloseable { // this should never happen throw new RabbitmqConnectorException(SETUP_SSL_FACTORY_FAILED, e); } + if (factory.isSSL()) { + configureSsl(factory); + } } else { factory.setHost(config.getHost()); factory.setPort(config.getPort()); @@ -121,6 +129,9 @@ public class RabbitmqClient implements AutoCloseable { } factory.setUsername(config.getUsername()); factory.setPassword(config.getPassword()); + if (config.isSsl()) { + configureSsl(factory); + } } if (config.getAutomaticRecovery() != null) { @@ -147,6 +158,16 @@ public class RabbitmqClient implements AutoCloseable { return factory; } + /** Configures TLS with the JVM trust store and hostname verification enabled. */ + private void configureSsl(ConnectionFactory factory) { + try { + factory.useSslProtocol(SSLContext.getDefault()); + factory.enableHostnameVerification(); + } catch (NoSuchAlgorithmException e) { + throw new RabbitmqConnectorException(INIT_SSL_CONTEXT_FAILED, e); + } + } + /** * Write data to RabbitMQ. * @@ -221,6 +242,29 @@ public class RabbitmqClient implements AutoCloseable { /** Declare a specific queue */ public void setupQueue(String queueName) throws IOException { if (StringUtils.isNotEmpty(queueName)) { + declareQueue(channel, config, queueName); + } + } + + private void declareQueueDefaults(Channel channel, RabbitmqConfig config) throws IOException { + declareQueue(channel, config, config.getQueueName()); + } + + static void declareQueue(Channel channel, RabbitmqConfig config, String queueName) + throws IOException { + if (config.isPassive()) { + try { + channel.queueDeclarePassive(queueName); + } catch (IOException e) { + throw new RabbitmqConnectorException( + ILLEGAL_CONFIG, + String.format( + "Cannot passively declare RabbitMQ queue '%s'. The queue may not exist " + + "or the account may lack access; create it first or set passive=false.", + queueName), + e); + } + } else { channel.queueDeclare( queueName, config.getDurable(), @@ -229,13 +273,4 @@ public class RabbitmqClient implements AutoCloseable { null); } } - - private void declareQueueDefaults(Channel channel, RabbitmqConfig config) throws IOException { - channel.queueDeclare( - config.getQueueName(), - config.getDurable(), - config.getExclusive(), - config.getAutoDelete(), - null); - } } diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java index 68d3ee15d0..1b93f5ccec 100644 --- a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqBaseOptions.java @@ -66,6 +66,19 @@ public class RabbitmqBaseOptions extends ConnectorCommonOptions { .withDescription( "convenience method for setting the fields in an AMQP URI: host, port, username, password and virtual host"); + public static final Option<String> URI = + Options.key("uri") + .stringType() + .noDefaultValue() + .withDescription("legacy alias of url for an AMQP URI"); + + public static final Option<Boolean> SSL = + Options.key("ssl") + .booleanType() + .defaultValue(false) + .withDescription( + "whether to enable SSL/TLS when connecting with host and port"); + public static final Option<String> ROUTING_KEY = Options.key("routing_key") .stringType() @@ -133,6 +146,13 @@ public class RabbitmqBaseOptions extends ConnectorCommonOptions { "true: The queue will be deleted automatically when the last consumer unsubscribes." + "false: The queue will not be automatically deleted."); + public static final Option<Boolean> PASSIVE = + Options.key("passive") + .booleanType() + .defaultValue(false) + .withDescription( + "whether to verify an existing queue without declaring or creating it"); + public static final Option<RabbitmqMessageFormat> FORMAT = Options.key("format") .enumType(RabbitmqMessageFormat.class) diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java index 1bab65f778..eabc9d4018 100644 --- a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/config/RabbitmqConfig.java @@ -18,6 +18,7 @@ package org.apache.seatunnel.connectors.seatunnel.rabbitmq.config; import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorException; import lombok.AllArgsConstructor; import lombok.Getter; @@ -28,17 +29,27 @@ import java.io.Serializable; import java.util.HashMap; import java.util.Map; +import static org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorErrorCode.ILLEGAL_CONFIG; + @Setter @Getter @NoArgsConstructor @AllArgsConstructor public class RabbitmqConfig implements Serializable { + /** + * Pinned to the default computed UID of the class before ssl/passive were added, so execution + * plans and checkpoint state serialized by older versions still deserialize without throwing + * InvalidClassException. + */ + private static final long serialVersionUID = -6715216959598971323L; + private String host; private Integer port; private String virtualHost; private String username; private String password; private String uri; + private boolean ssl; private Integer networkRecoveryInterval; private Boolean automaticRecovery; private Boolean topologyRecovery; @@ -52,6 +63,7 @@ public class RabbitmqConfig implements Serializable { private Boolean durable; private Boolean exclusive; private Boolean autoDelete; + private boolean passive; private RabbitmqMessageFormat format; private String protobufSchema; private String protobufMessageName; @@ -127,6 +139,8 @@ public class RabbitmqConfig implements Serializable { this.durable = config.get(RabbitmqBaseOptions.DURABLE); this.exclusive = config.get(RabbitmqBaseOptions.EXCLUSIVE); this.autoDelete = config.get(RabbitmqBaseOptions.AUTO_DELETE); + this.ssl = config.get(RabbitmqBaseOptions.SSL); + this.passive = config.get(RabbitmqBaseOptions.PASSIVE); this.format = config.get(RabbitmqBaseOptions.FORMAT); if (config.getOptional(RabbitmqBaseOptions.PROTOBUF_SCHEMA).isPresent()) { this.protobufSchema = config.get(RabbitmqBaseOptions.PROTOBUF_SCHEMA); @@ -137,8 +151,17 @@ public class RabbitmqConfig implements Serializable { if (config.getOptional(RabbitmqSinkOptions.RABBITMQ_CONFIG).isPresent()) { this.sinkOptionProps = config.get(RabbitmqSinkOptions.RABBITMQ_CONFIG); } - if (config.getOptional(RabbitmqBaseOptions.URL).isPresent()) { + boolean hasUrl = config.getOptional(RabbitmqBaseOptions.URL).isPresent(); + boolean hasUri = config.getOptional(RabbitmqBaseOptions.URI).isPresent(); + if (hasUrl && hasUri) { + throw new RabbitmqConnectorException( + ILLEGAL_CONFIG, + "RabbitMQ connector options 'url' and 'uri' are mutually exclusive; please configure only one of them. 'uri' is a legacy alias kept for backward compatibility, prefer 'url' for new configurations."); + } + if (hasUrl) { this.uri = config.get(RabbitmqBaseOptions.URL); + } else if (hasUri) { + this.uri = config.get(RabbitmqBaseOptions.URI); } } } diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java index 35e8eedd43..d31472d6f7 100644 --- a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/exception/RabbitmqConnectorErrorCode.java @@ -29,7 +29,8 @@ public enum RabbitmqConnectorErrorCode implements SeaTunnelErrorCode { MESSAGE_ACK_REJECTED("RABBITMQ-06", "messages could not be acknowledged with basicReject"), PARSE_URI_FAILED("RABBITMQ-07", "parse uri failed"), INIT_SSL_CONTEXT_FAILED("RABBITMQ-08", "initialize ssl context failed"), - SETUP_SSL_FACTORY_FAILED("RABBITMQ-09", "setup ssl factory failed"); + SETUP_SSL_FACTORY_FAILED("RABBITMQ-09", "setup ssl factory failed"), + ILLEGAL_CONFIG("RABBITMQ-10", "illegal connector configuration"); private final String code; private final String description; diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java index 906b5425a6..c41a0ea930 100644 --- a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/sink/RabbitmqSinkFactory.java @@ -22,6 +22,7 @@ import org.apache.seatunnel.api.table.connector.TableSink; import org.apache.seatunnel.api.table.factory.Factory; import org.apache.seatunnel.api.table.factory.TableSinkFactory; import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext; +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqBaseOptions; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqConfig; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqMessageFormat; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqSinkOptions; @@ -53,6 +54,8 @@ public class RabbitmqSinkFactory implements TableSinkFactory { RabbitmqSinkOptions.PROTOBUF_MESSAGE_NAME) .optional( RabbitmqSinkOptions.URL, + RabbitmqBaseOptions.URI, + RabbitmqSinkOptions.SSL, RabbitmqSinkOptions.ROUTING_KEY, RabbitmqSinkOptions.EXCHANGE, RabbitmqSinkOptions.NETWORK_RECOVERY_INTERVAL, @@ -63,6 +66,7 @@ public class RabbitmqSinkFactory implements TableSinkFactory { RabbitmqSinkOptions.DURABLE, RabbitmqSinkOptions.EXCLUSIVE, RabbitmqSinkOptions.AUTO_DELETE, + RabbitmqSinkOptions.PASSIVE, RabbitmqSinkOptions.RABBITMQ_CONFIG) .build(); } diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java index 48a248fe74..969a3880b8 100644 --- a/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/source/RabbitmqSourceFactory.java @@ -25,6 +25,7 @@ import org.apache.seatunnel.api.table.connector.TableSource; import org.apache.seatunnel.api.table.factory.Factory; import org.apache.seatunnel.api.table.factory.TableSourceFactory; import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext; +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqBaseOptions; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqMessageFormat; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqSingleTableValidator; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqSinkOptions; @@ -69,6 +70,8 @@ public class RabbitmqSourceFactory implements TableSourceFactory { .optional( RabbitmqSourceOptions.VIRTUAL_HOST, RabbitmqSourceOptions.URL, + RabbitmqBaseOptions.URI, + RabbitmqSourceOptions.SSL, RabbitmqSourceOptions.ROUTING_KEY, RabbitmqSourceOptions.EXCHANGE, RabbitmqSourceOptions.NETWORK_RECOVERY_INTERVAL, @@ -79,6 +82,7 @@ public class RabbitmqSourceFactory implements TableSourceFactory { RabbitmqSinkOptions.DURABLE, RabbitmqSinkOptions.EXCLUSIVE, RabbitmqSinkOptions.AUTO_DELETE, + RabbitmqSourceOptions.PASSIVE, RabbitmqSourceOptions.REQUESTED_CHANNEL_MAX, RabbitmqSourceOptions.REQUESTED_FRAME_MAX, RabbitmqSourceOptions.REQUESTED_HEARTBEAT, diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java index 36641e585f..e5ba329aa8 100644 --- a/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/RabbitmqFactoryTest.java @@ -17,12 +17,19 @@ package org.apache.seatunnel.connectors.seatunnel.rabbitmq; +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqBaseOptions; +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqConfig; +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorException; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.sink.RabbitmqSinkFactory; import org.apache.seatunnel.connectors.seatunnel.rabbitmq.source.RabbitmqSourceFactory; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.util.HashMap; +import java.util.Map; + class RabbitmqFactoryTest { @Test @@ -30,4 +37,39 @@ class RabbitmqFactoryTest { Assertions.assertNotNull((new RabbitmqSourceFactory()).optionRule()); Assertions.assertNotNull((new RabbitmqSinkFactory()).optionRule()); } + + @Test + void readsLegacyUriOption() { + Map<String, Object> options = new HashMap<>(); + options.put(RabbitmqBaseOptions.URI.key(), "amqps://guest:guest@localhost:5671/%2F"); + + RabbitmqConfig config = new RabbitmqConfig(ReadonlyConfig.fromMap(options)); + + Assertions.assertEquals("amqps://guest:guest@localhost:5671/%2F", config.getUri()); + } + + @Test + void readsSecurePassiveOptions() { + Map<String, Object> options = new HashMap<>(); + options.put(RabbitmqBaseOptions.HOST.key(), "localhost"); + options.put(RabbitmqBaseOptions.PORT.key(), 5671); + options.put(RabbitmqBaseOptions.SSL.key(), true); + options.put(RabbitmqBaseOptions.PASSIVE.key(), true); + + RabbitmqConfig config = new RabbitmqConfig(ReadonlyConfig.fromMap(options)); + + Assertions.assertTrue(config.isSsl()); + Assertions.assertTrue(config.isPassive()); + } + + @Test + void rejectsUrlAndUriSetTogether() { + Map<String, Object> options = new HashMap<>(); + options.put(RabbitmqBaseOptions.URL.key(), "amqp://host-a:5672/%2F"); + options.put(RabbitmqBaseOptions.URI.key(), "amqps://host-b:5671/%2F"); + + Assertions.assertThrows( + RabbitmqConnectorException.class, + () -> new RabbitmqConfig(ReadonlyConfig.fromMap(options))); + } } diff --git a/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClientTest.java b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClientTest.java new file mode 100644 index 0000000000..a2ab086521 --- /dev/null +++ b/seatunnel-connectors-v2/connector-rabbitmq/src/test/java/org/apache/seatunnel/connectors/seatunnel/rabbitmq/client/RabbitmqClientTest.java @@ -0,0 +1,83 @@ +/* + * 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.rabbitmq.client; + +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.config.RabbitmqConfig; +import org.apache.seatunnel.connectors.seatunnel.rabbitmq.exception.RabbitmqConnectorException; + +import org.junit.jupiter.api.Test; + +import com.rabbitmq.client.Channel; + +import java.io.IOException; + +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.Mockito.doThrow; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +class RabbitmqClientTest { + + @Test + void declaresExistingQueuePassivelyWhenConfigured() throws Exception { + Channel channel = mock(Channel.class); + RabbitmqConfig config = mock(RabbitmqConfig.class); + when(config.isPassive()).thenReturn(true); + + RabbitmqClient.declareQueue(channel, config, "existing-queue"); + + verify(channel).queueDeclarePassive("existing-queue"); + verify(channel, never()).queueDeclare("existing-queue", true, false, false, null); + } + + @Test + void declaresQueueWithConfiguredPropertiesByDefault() throws Exception { + Channel channel = mock(Channel.class); + RabbitmqConfig config = mock(RabbitmqConfig.class); + when(config.isPassive()).thenReturn(false); + when(config.getDurable()).thenReturn(true); + when(config.getExclusive()).thenReturn(false); + when(config.getAutoDelete()).thenReturn(false); + + RabbitmqClient.declareQueue(channel, config, "new-queue"); + + verify(channel).queueDeclare("new-queue", true, false, false, null); + verify(channel, never()).queueDeclarePassive("new-queue"); + } + + @Test + void explainsPassiveQueueDeclarationFailure() throws Exception { + Channel channel = mock(Channel.class); + RabbitmqConfig config = mock(RabbitmqConfig.class); + when(config.isPassive()).thenReturn(true); + doThrow(new IOException("queue not found")) + .when(channel) + .queueDeclarePassive("missing-queue"); + + RabbitmqConnectorException exception = + assertThrows( + RabbitmqConnectorException.class, + () -> RabbitmqClient.declareQueue(channel, config, "missing-queue")); + + assertTrue(exception.getMessage().contains("missing-queue")); + assertTrue(exception.getMessage().contains("passive=false")); + } +}
