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-12065-83454b590dbf9442faecb28974b815fbb6ae2b82 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit ce4cde54c8abc7a6a9e3831110dbab029512ef7b Author: Yohanes <[email protected]> AuthorDate: Fri Sep 4 21:55:00 2026 +0000 [Improve][Core] Add reusable config validation result (#12065) Co-authored-by: Yohanes <[email protected]> --- .../starter/validation/ConfigValidationError.java | 74 +++++++++++++ .../starter/validation/ConfigValidationResult.java | 116 +++++++++++++++++++++ .../validation/ConfigValidationResultTest.java | 66 ++++++++++++ .../command/SeaTunnelConfValidateCommand.java | 90 ++++++++++++++++ .../command/SeaTunnelConfValidateCommandTest.java | 111 +++++++++++++++++++- 5 files changed, 456 insertions(+), 1 deletion(-) diff --git a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/validation/ConfigValidationError.java b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/validation/ConfigValidationError.java new file mode 100644 index 0000000000..fb57bfffbe --- /dev/null +++ b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/validation/ConfigValidationError.java @@ -0,0 +1,74 @@ +/* + * 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.core.starter.validation; + +import java.io.Serializable; + +/** A machine-readable validation failure at config level. */ +public final class ConfigValidationError implements Serializable { + + private static final long serialVersionUID = 1L; + + /** Plugin location, such as {@code source[0](FakeSource)}, when the phase can identify it. */ + private final String location; + + /** Plugin factory identifier parsed from {@link #location}, when available. */ + private final String plugin; + + /** Option or option group reported by the underlying option validator, when available. */ + private final String optionPath; + + /** Stable, closed category describing the validation rule that failed. */ + private final String ruleCategory; + + /** Sanitized diagnostic message suitable for programmatic consumers. */ + private final String message; + + public ConfigValidationError( + String location, + String plugin, + String optionPath, + String ruleCategory, + String message) { + this.location = location; + this.plugin = plugin; + this.optionPath = optionPath; + this.ruleCategory = ruleCategory; + this.message = message; + } + + public String getLocation() { + return location; + } + + public String getPlugin() { + return plugin; + } + + public String getOptionPath() { + return optionPath; + } + + public String getRuleCategory() { + return ruleCategory; + } + + public String getMessage() { + return message; + } +} diff --git a/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/validation/ConfigValidationResult.java b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/validation/ConfigValidationResult.java new file mode 100644 index 0000000000..49ea630661 --- /dev/null +++ b/seatunnel-core/seatunnel-core-starter/src/main/java/org/apache/seatunnel/core/starter/validation/ConfigValidationResult.java @@ -0,0 +1,116 @@ +/* + * 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.core.starter.validation; + +import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.node.ArrayNode; +import org.apache.seatunnel.shade.com.fasterxml.jackson.databind.node.ObjectNode; + +import org.apache.seatunnel.common.utils.JsonUtils; + +import java.io.Serializable; +import java.util.ArrayList; +import java.util.Collections; +import java.util.List; + +/** + * Versioned, config-level validation result shared by CLI and future adapters. + * + * <p>This model intentionally describes static/config validation only. It does not imply that a + * connector can reach its external system or that a generated job is runtime-equivalent. + */ +public final class ConfigValidationResult implements Serializable { + + private static final long serialVersionUID = 1L; + public static final String SCHEMA_VERSION = "1.0"; + + private final boolean valid; + private final String phase; + private final List<ConfigValidationError> errors; + + private ConfigValidationResult( + boolean valid, String phase, List<ConfigValidationError> errors) { + this.valid = valid; + this.phase = phase; + this.errors = Collections.unmodifiableList(new ArrayList<>(errors)); + } + + public static ConfigValidationResult success(String phase) { + return new ConfigValidationResult(true, phase, Collections.emptyList()); + } + + public static ConfigValidationResult failure(String phase, ConfigValidationError error) { + return new ConfigValidationResult(false, phase, Collections.singletonList(error)); + } + + public boolean isValid() { + return valid; + } + + public String getPhase() { + return phase; + } + + public List<ConfigValidationError> getErrors() { + return errors; + } + + /** Serialize fields in a fixed order so adapters can rely on a stable shape. */ + public String toJson() { + ObjectNode root = JsonUtils.createObjectNode(); + root.put("schemaVersion", SCHEMA_VERSION); + root.put("valid", valid); + root.put("phase", phase); + ArrayNode errorNodes = root.putArray("errors"); + for (ConfigValidationError error : errors) { + ObjectNode errorNode = errorNodes.addObject(); + putNullable(errorNode, "location", error.getLocation()); + putNullable(errorNode, "plugin", error.getPlugin()); + putNullable(errorNode, "optionPath", error.getOptionPath()); + putNullable(errorNode, "ruleCategory", error.getRuleCategory()); + putNullable(errorNode, "message", error.getMessage()); + } + return root.toString(); + } + + /** Preserve the existing --check message for the CLI adapter. */ + public String toHumanReadable() { + if (valid) { + return "VALID"; + } + String message = errors.isEmpty() ? "Validation failed" : errors.get(0).getMessage(); + return humanPhase(phase) + " failed: " + message; + } + + private static String humanPhase(String phase) { + if ("connectivity".equals(phase)) { + return "Connectivity check"; + } + if ("static".equals(phase)) { + return "Static analysis"; + } + return phase; + } + + private static void putNullable(ObjectNode node, String name, String value) { + if (value == null) { + node.putNull(name); + } else { + node.put(name, value); + } + } +} diff --git a/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/validation/ConfigValidationResultTest.java b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/validation/ConfigValidationResultTest.java new file mode 100644 index 0000000000..9d15fde2ad --- /dev/null +++ b/seatunnel-core/seatunnel-core-starter/src/test/java/org/apache/seatunnel/core/starter/validation/ConfigValidationResultTest.java @@ -0,0 +1,66 @@ +/* + * 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.core.starter.validation; + +import org.apache.seatunnel.common.utils.JsonUtils; + +import org.junit.jupiter.api.Assertions; +import org.junit.jupiter.api.Test; + +public class ConfigValidationResultTest { + + @Test + void serializesSuccessWithStableSchema() throws Exception { + ConfigValidationResult result = ConfigValidationResult.success("static"); + + Assertions.assertEquals( + "{\"schemaVersion\":\"1.0\",\"valid\":true," + + "\"phase\":\"static\",\"errors\":[]}", + result.toJson()); + Assertions.assertTrue( + JsonUtils.readTree( + result.toJson().getBytes(java.nio.charset.StandardCharsets.UTF_8)) + .get("valid") + .asBoolean()); + Assertions.assertEquals("VALID", result.toHumanReadable()); + } + + @Test + void serializesFailureAndPreservesNullableFields() throws Exception { + ConfigValidationError error = + new ConfigValidationError( + "source[0](Kafka)", "Kafka", null, "option", "Required option is missing"); + ConfigValidationResult result = ConfigValidationResult.failure("static", error); + + Assertions.assertEquals( + "{\"schemaVersion\":\"1.0\",\"valid\":false," + + "\"phase\":\"static\",\"errors\":[{" + + "\"location\":\"source[0](Kafka)\",\"plugin\":\"Kafka\"," + + "\"optionPath\":null,\"ruleCategory\":\"option\"," + + "\"message\":\"Required option is missing\"}]}", + result.toJson()); + Assertions.assertEquals( + "Static analysis failed: Required option is missing", result.toHumanReadable()); + Assertions.assertEquals( + "option", + JsonUtils.readTree( + result.toJson().getBytes(java.nio.charset.StandardCharsets.UTF_8)) + .at("/errors/0/ruleCategory") + .asText()); + } +} diff --git a/seatunnel-core/seatunnel-starter/src/main/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommand.java b/seatunnel-core/seatunnel-starter/src/main/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommand.java index b9bc0c5333..a1b77575e1 100644 --- a/seatunnel-core/seatunnel-starter/src/main/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommand.java +++ b/seatunnel-core/seatunnel-starter/src/main/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommand.java @@ -39,6 +39,8 @@ import org.apache.seatunnel.core.starter.exception.ConfigCheckException; import org.apache.seatunnel.core.starter.seatunnel.args.ClientCommandArgs; import org.apache.seatunnel.core.starter.utils.ConfigBuilder; import org.apache.seatunnel.core.starter.utils.FileUtils; +import org.apache.seatunnel.core.starter.validation.ConfigValidationError; +import org.apache.seatunnel.core.starter.validation.ConfigValidationResult; import org.apache.seatunnel.engine.core.parse.ConfigParserUtil; import org.apache.seatunnel.engine.core.parse.JobPluginClasspathHelper; @@ -50,6 +52,9 @@ import java.net.URL; import java.nio.file.Path; import java.util.Collections; import java.util.List; +import java.util.Locale; +import java.util.regex.Matcher; +import java.util.regex.Pattern; import static org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAME; @@ -83,6 +88,11 @@ import static org.apache.seatunnel.api.options.ConnectorCommonOptions.PLUGIN_NAM @Slf4j public class SeaTunnelConfValidateCommand implements Command<ClientCommandArgs> { + private static final Pattern PLUGIN_LOCATION_PATTERN = + Pattern.compile("((?:source|transform|sink)\\[\\d+\\]\\([^)]*\\))"); + private static final Pattern OPTION_PATH_PATTERN = + Pattern.compile("(?m)^\\s*options?:\\s*([^\\r\\n]+)"); + private final ClientCommandArgs clientCommandArgs; public SeaTunnelConfValidateCommand(ClientCommandArgs clientCommandArgs) { @@ -187,6 +197,86 @@ public class SeaTunnelConfValidateCommand implements Command<ClientCommandArgs> } } + /** + * Validate the configuration and return a reusable result for non-CLI integrations. + * + * <p>The result is deliberately config-level and does not claim runtime-equivalent validation. + */ + public ConfigValidationResult validateResult() { + try { + execute(); + return ConfigValidationResult.success(validationPhase()); + } catch (ConfigCheckException e) { + String message = e.getMessage(); + String prefix = validationMode() + " failed: "; + if (message != null && message.startsWith(prefix)) { + message = message.substring(prefix.length()); + } + // The result is intended for programmatic consumers, so never expose credentials + // even when the underlying validation phase is static. + message = DryRunConnectFailureMessageSanitizer.sanitize(message); + return ConfigValidationResult.failure( + validationPhase(), + toValidationError(message == null ? "Validation failed" : message)); + } + } + + private String validationPhase() { + return clientCommandArgs.getDryRun() == DryRun.CONNECT ? "connectivity" : "static"; + } + + private String validationMode() { + return clientCommandArgs.getDryRun() == DryRun.CONNECT + ? "Connectivity check" + : "Static analysis"; + } + + private ConfigValidationError toValidationError(String message) { + String location = null; + String plugin = null; + Matcher locationMatcher = PLUGIN_LOCATION_PATTERN.matcher(message); + if (locationMatcher.find()) { + location = locationMatcher.group(1); + int open = location.lastIndexOf('('); + plugin = location.substring(open + 1, location.length() - 1); + } + + Matcher optionPathMatcher = OPTION_PATH_PATTERN.matcher(message); + String optionPath = optionPathMatcher.find() ? optionPathMatcher.group(1).trim() : null; + + String lower = message.toLowerCase(Locale.ROOT); + ValidationRuleCategory ruleCategory; + if (lower.contains("parse") || lower.contains("syntax") || lower.contains("hocon")) { + ruleCategory = ValidationRuleCategory.PARSE; + } else if (lower.contains("option") + || lower.contains("required") + || lower.contains("unknown key") + || lower.contains("type")) { + ruleCategory = ValidationRuleCategory.OPTION; + } else if (lower.contains("plugin") + || lower.contains("factory") + || lower.contains("classloader")) { + ruleCategory = ValidationRuleCategory.PLUGIN; + } else { + ruleCategory = ValidationRuleCategory.VALIDATION; + } + return new ConfigValidationError(location, plugin, optionPath, ruleCategory.value, message); + } + + /** Closed categories exposed by the current structured validation result schema. */ + private enum ValidationRuleCategory { + PARSE("parse"), + OPTION("option"), + PLUGIN("plugin"), + VALIDATION("validation"); + + private final String value; + + ValidationRuleCategory(String value) { + this.value = value; + } + } + private void validateOptionTypes(ReadonlyConfig config, OptionRule rule) { if (rule == null) { return; diff --git a/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommandTest.java b/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommandTest.java index 1724b0eb00..d44673bf1a 100644 --- a/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommandTest.java +++ b/seatunnel-core/seatunnel-starter/src/test/java/org/apache/seatunnel/core/starter/seatunnel/command/SeaTunnelConfValidateCommandTest.java @@ -25,6 +25,7 @@ import org.apache.seatunnel.core.starter.exception.ConfigCheckException; import org.apache.seatunnel.core.starter.seatunnel.args.ClientCommandArgs; import org.apache.seatunnel.core.starter.utils.CommandLineUtils; import org.apache.seatunnel.core.starter.utils.ConfigBuilder; +import org.apache.seatunnel.core.starter.validation.ConfigValidationResult; import org.junit.jupiter.api.Assertions; import org.junit.jupiter.api.Test; @@ -49,6 +50,110 @@ public class SeaTunnelConfValidateCommandTest { Assertions.assertDoesNotThrow(command::execute); } + @Test + public void testValidationResultForValidConfig() { + SeaTunnelConfValidateCommand command = + new SeaTunnelConfValidateCommand(buildArgs("config/valid_static_dryrun.json")); + + ConfigValidationResult result = command.validateResult(); + + Assertions.assertTrue(result.isValid()); + Assertions.assertEquals("static", result.getPhase()); + Assertions.assertTrue(result.getErrors().isEmpty()); + } + + @Test + public void testValidationResultClassifiesMissingPluginAsOptionError() throws Exception { + Path configFile = Files.createTempFile("seatunnel-validation-result", ".conf"); + Files.write( + configFile, + ("source { FakeSource { plugin_output = output } }\n" + "sink { InMemory {} }") + .getBytes(StandardCharsets.UTF_8)); + configFile.toFile().deleteOnExit(); + + SeaTunnelConfValidateCommand command = + new SeaTunnelConfValidateCommand(buildArgsFromPath(configFile.toString())); + ConfigValidationResult result = command.validateResult(); + + Assertions.assertFalse(result.isValid()); + Assertions.assertEquals(1, result.getErrors().size()); + Assertions.assertEquals("option", result.getErrors().get(0).getRuleCategory()); + Assertions.assertTrue(result.toJson().contains("\"schemaVersion\":\"1.0\"")); + } + + @Test + public void testValidationResultClassifiesParseFailure() { + SeaTunnelConfValidateCommand command = + new SeaTunnelConfValidateCommand(buildArgs("config/invalid_hocon_syntax.conf")); + + ConfigValidationResult result = command.validateResult(); + + Assertions.assertFalse(result.isValid()); + Assertions.assertEquals("parse", result.getErrors().get(0).getRuleCategory()); + } + + @Test + public void testValidationResultClassifiesOptionFailure() { + SeaTunnelConfValidateCommand command = + new SeaTunnelConfValidateCommand(buildArgs("config/invalid_option_type.json")); + + ConfigValidationResult result = command.validateResult(); + + Assertions.assertFalse(result.isValid()); + Assertions.assertEquals("option", result.getErrors().get(0).getRuleCategory()); + } + + @Test + public void testValidationResultIncludesConnectFailureLocationAndPlugin() throws Exception { + Path configFile = Files.createTempFile("seatunnel-validation-connect-failure", ".conf"); + Files.write( + configFile, + ("source { DryRunTestSource { fail_connection = true } }\n" + + "sink { InMemory {} }") + .getBytes(StandardCharsets.UTF_8)); + configFile.toFile().deleteOnExit(); + + SeaTunnelConfValidateCommand command = + new SeaTunnelConfValidateCommand(buildConnectArgsFromPath(configFile.toString())); + ConfigValidationResult result = command.validateResult(); + + Assertions.assertFalse(result.isValid()); + Assertions.assertEquals("connectivity", result.getPhase()); + Assertions.assertEquals( + "source[0](DryRunTestSource)", result.getErrors().get(0).getLocation()); + Assertions.assertEquals("DryRunTestSource", result.getErrors().get(0).getPlugin()); + } + + @Test + public void testValidationResultSanitizesConnectFailure() throws Exception { + Path configFile = Files.createTempFile("seatunnel-validation-sensitive-failure", ".conf"); + Files.write( + configFile, + ("source { DryRunTestSource { sensitive_connection_failure = true } }\n" + + "sink { InMemory {} }") + .getBytes(StandardCharsets.UTF_8)); + configFile.toFile().deleteOnExit(); + + SeaTunnelConfValidateCommand command = + new SeaTunnelConfValidateCommand(buildConnectArgsFromPath(configFile.toString())); + String message = command.validateResult().getErrors().get(0).getMessage(); + + Assertions.assertFalse(message.contains("secret-password"), message); + Assertions.assertFalse(message.contains("secret-token"), message); + } + + @Test + public void testValidationResultClassifiesPluginLoadFailure() { + SeaTunnelConfValidateCommand command = + new SeaTunnelConfValidateCommand( + buildArgs("config/invalid_plugin_loadability.json")); + + ConfigValidationResult result = command.validateResult(); + + Assertions.assertFalse(result.isValid()); + Assertions.assertEquals("plugin", result.getErrors().get(0).getRuleCategory()); + } + @Test public void testValidConnectDryRun() { ClientCommandArgs args = buildConnectArgs("config/valid_static_dryrun.json"); @@ -497,7 +602,11 @@ public class SeaTunnelConfValidateCommandTest { } private ClientCommandArgs buildArgs(String configFile) { - String[] args = {"-c", resolveConfigPath(configFile), "--dry-run", "static"}; + return buildArgsFromPath(resolveConfigPath(configFile)); + } + + private ClientCommandArgs buildArgsFromPath(String configPath) { + String[] args = {"-c", configPath, "--dry-run", "static"}; return CommandLineUtils.parse(args, new ClientCommandArgs(), "seatunnel.sh", true); }
