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-12004-0071e8fa608a6b25fb0dd4bab044bc41a4e263b8 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit ffdffb98290f20f8a1b5725e72d3c31bb7a9c253 Author: Nikhil kumar <[email protected]> AuthorDate: Sun Aug 30 04:20:50 2026 +0000 [Improve][Connector-V2] Migrate ActiveMQ validation to OptionRule (#12004) --- .../activemq/sink/ActivemqSinkFactory.java | 8 ++- .../seatunnel/activemq/ActivemqFactoryTest.java | 72 +++++++++++++++++++++- 2 files changed, 77 insertions(+), 3 deletions(-) diff --git a/seatunnel-connectors-v2/connector-activemq/src/main/java/org/apache/seatunnel/connectors/seatunnel/activemq/sink/ActivemqSinkFactory.java b/seatunnel-connectors-v2/connector-activemq/src/main/java/org/apache/seatunnel/connectors/seatunnel/activemq/sink/ActivemqSinkFactory.java index a34ba7a703..55d14771a1 100644 --- a/seatunnel-connectors-v2/connector-activemq/src/main/java/org/apache/seatunnel/connectors/seatunnel/activemq/sink/ActivemqSinkFactory.java +++ b/seatunnel-connectors-v2/connector-activemq/src/main/java/org/apache/seatunnel/connectors/seatunnel/activemq/sink/ActivemqSinkFactory.java @@ -25,11 +25,13 @@ import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext; import com.google.auto.service.AutoService; +import static org.apache.seatunnel.api.configuration.util.Conditions.notBlank; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.ALWAYS_SESSION_ASYNC; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.ALWAYS_SYNC_SEND; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.CHECK_FOR_DUPLICATE; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.CLIENT_ID; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.CLOSE_TIMEOUT; +import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.CONSUMER_EXPIRY_CHECK_ENABLED; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.DISPATCH_ASYNC; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.NESTED_MAP_AND_LIST_ENABLED; import static org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions.PASSWORD; @@ -49,14 +51,16 @@ public class ActivemqSinkFactory implements TableSinkFactory { @Override public OptionRule optionRule() { return OptionRule.builder() - .required(QUEUE_NAME, URI) + .required(QUEUE_NAME, notBlank(QUEUE_NAME)) + .required(URI, notBlank(URI)) .bundled(USERNAME, PASSWORD) + .optional(CLIENT_ID, notBlank(CLIENT_ID)) .optional( - CLIENT_ID, CHECK_FOR_DUPLICATE, ALWAYS_SESSION_ASYNC, ALWAYS_SYNC_SEND, CLOSE_TIMEOUT, + CONSUMER_EXPIRY_CHECK_ENABLED, DISPATCH_ASYNC, NESTED_MAP_AND_LIST_ENABLED, WARN_ABOUT_UNSTARTED_CONNECTION_TIMEOUT) diff --git a/seatunnel-connectors-v2/connector-activemq/src/test/java/org/apache/seatunnel/connectors/seatunnel/activemq/ActivemqFactoryTest.java b/seatunnel-connectors-v2/connector-activemq/src/test/java/org/apache/seatunnel/connectors/seatunnel/activemq/ActivemqFactoryTest.java index 90732d8a0e..56defbb500 100644 --- a/seatunnel-connectors-v2/connector-activemq/src/test/java/org/apache/seatunnel/connectors/seatunnel/activemq/ActivemqFactoryTest.java +++ b/seatunnel-connectors-v2/connector-activemq/src/test/java/org/apache/seatunnel/connectors/seatunnel/activemq/ActivemqFactoryTest.java @@ -17,15 +17,85 @@ package org.apache.seatunnel.connectors.seatunnel.activemq; +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.configuration.util.ConfigValidator; +import org.apache.seatunnel.api.configuration.util.OptionRule; +import org.apache.seatunnel.api.configuration.util.OptionValidationException; +import org.apache.seatunnel.connectors.seatunnel.activemq.config.ActivemqSinkOptions; import org.apache.seatunnel.connectors.seatunnel.activemq.sink.ActivemqSinkFactory; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; +import java.util.HashMap; +import java.util.Map; + class ActivemqFactoryTest { + private final OptionRule optionRule = new ActivemqSinkFactory().optionRule(); + @Test void optionRule() { - Assertions.assertNotNull((new ActivemqSinkFactory()).optionRule()); + Assertions.assertNotNull(optionRule); + } + + @Test + void testValidRequiredOptions() { + Assertions.assertDoesNotThrow(() -> validate(requiredConfig())); + } + + @Test + void testBlankRequiredOptionsRejected() { + Map<String, Object> blankUriConfig = requiredConfig(); + blankUriConfig.put(ActivemqSinkOptions.URI.key(), " "); + Assertions.assertThrows(OptionValidationException.class, () -> validate(blankUriConfig)); + + Map<String, Object> blankQueueNameConfig = requiredConfig(); + blankQueueNameConfig.put(ActivemqSinkOptions.QUEUE_NAME.key(), "\t"); + Assertions.assertThrows( + OptionValidationException.class, () -> validate(blankQueueNameConfig)); + } + + @Test + void testBlankClientIdRejected() { + Map<String, Object> config = requiredConfig(); + config.put(ActivemqSinkOptions.CLIENT_ID.key(), ""); + Assertions.assertThrows(OptionValidationException.class, () -> validate(config)); + } + + @Test + void testCredentialsMustBeBundled() { + Map<String, Object> usernameOnlyConfig = requiredConfig(); + usernameOnlyConfig.put(ActivemqSinkOptions.USERNAME.key(), "user"); + Assertions.assertThrows( + OptionValidationException.class, () -> validate(usernameOnlyConfig)); + + Map<String, Object> passwordOnlyConfig = requiredConfig(); + passwordOnlyConfig.put(ActivemqSinkOptions.PASSWORD.key(), "password"); + Assertions.assertThrows( + OptionValidationException.class, () -> validate(passwordOnlyConfig)); + } + + @Test + void testSupportedOptionalOptions() { + Map<String, Object> config = requiredConfig(); + config.put(ActivemqSinkOptions.CLIENT_ID.key(), "client-id"); + config.put(ActivemqSinkOptions.CLOSE_TIMEOUT.key(), 1000); + config.put(ActivemqSinkOptions.CONSUMER_EXPIRY_CHECK_ENABLED.key(), true); + config.put(ActivemqSinkOptions.WARN_ABOUT_UNSTARTED_CONNECTION_TIMEOUT.key(), -1); + Assertions.assertDoesNotThrow(() -> validate(config)); + } + + private void validate(Map<String, Object> config) { + ReadonlyConfig readonlyConfig = ReadonlyConfig.fromMap(config); + ConfigValidator.validateUnknownKeys(readonlyConfig, optionRule, "ActiveMQSink"); + ConfigValidator.of(readonlyConfig).validate(optionRule); + } + + private Map<String, Object> requiredConfig() { + Map<String, Object> config = new HashMap<>(); + config.put(ActivemqSinkOptions.URI.key(), "tcp://localhost:61616"); + config.put(ActivemqSinkOptions.QUEUE_NAME.key(), "test-queue"); + return config; } }
