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-12242-5e918f399afad1876366c14ff2bc6b84ed09c5e5 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit 3339132550594fe3128534cf37c86f0cbbe435cc Author: Goutam Adwant <[email protected]> AuthorDate: Sun Sep 13 10:31:49 2026 +0000 [Feature][Connector-V2] Support S3 source connectivity dry-run (#12242) Signed-off-by: Goutam Adwant <[email protected]> --- docs/en/connectors/source/S3File.md | 40 ++ docs/en/engines/zeta/user-command.md | 1 + docs/zh/connectors/source/S3File.md | 32 ++ docs/zh/engines/zeta/user-command.md | 1 + .../file/s3/source/S3FileSourceFactory.java | 19 +- .../file/s3/source/S3SourceDryRunValidator.java | 284 +++++++++++++ .../seatunnel/file/s3/S3FileSourceFactoryTest.java | 8 + .../s3/source/S3SourceDryRunValidatorTest.java | 463 +++++++++++++++++++++ .../connector/file/s3/S3FileConnectDryRunIT.java | 190 +++++++++ .../file/s3/S3FileConnectDryRunWireIT.java | 249 +++++++++++ 10 files changed, 1286 insertions(+), 1 deletion(-) diff --git a/docs/en/connectors/source/S3File.md b/docs/en/connectors/source/S3File.md index 37ec943dd6..5672df90fb 100644 --- a/docs/en/connectors/source/S3File.md +++ b/docs/en/connectors/source/S3File.md @@ -41,6 +41,46 @@ import ChangeLog from '../changelog/connector-file-s3.md'; Read data from aws s3 file system. +### Connectivity dry-run + +`--dry-run connect` can check a single S3A source before submitting a job. The +source must use `bucket = "s3a://your-bucket"`, an absolute `path`, an explicit +inline `schema.fields` or `schema.columns`, and `file_format_type` of `text`, +`csv`, `json`, or `xml`. Set `parse_partition_from_path = false` and omit +`read_columns`; file-derived schemas, projection, and partition inference are +not validated by this metadata-only check. Unsupported configurations fail the +connect dry-run with an explanation; normal job execution is unchanged. + +The check reuses Hadoop S3A endpoint, credential-chain, proxy and path-style +configuration. It checks object metadata with HEAD, or makes one prefix listing +with `maxKeys=1` and delimiter `/`. It does not open file contents, recursively +list files, create readers, upload, delete, or initialize a shared filesystem. +An exact object does not require listing permission. A prefix requires listing +permission; a successful empty listing is accepted for `discovery_mode = +"continuous"`, since files may arrive later. For a batch source, an empty +virtual prefix without a directory marker fails; an accessible empty bucket +root is accepted. Missing buckets and denied requests fail in both modes. + +Validation-only connection establishment and socket timeouts are capped at +5 seconds, preserving smaller positive values. SDK request retries are disabled +for validation, including bucket-specific overrides. These are network timeout +settings, not a total deadline for DNS, credential-provider initialization, or +SDK setup. Runtime timeouts and retries are unchanged. + +This initial check does not support `tables_configs`, legacy `s3n` buckets, +SSE-C customer-provided encryption keys, `fs.s3a.security.credential.provider.path`, +S3Guard, multipart purge, or custom S3 client factories. It does not prove object +content readability, file-format correctness, schema compatibility with the +stored data, worker-side credentials, or target/update/post-sync permissions. + +Save your job configuration meeting the requirements above as +`config/s3-to-console.conf` (this is a user-created file, not a bundled template), +then run from the SeaTunnel installation directory: + +```bash +bin/seatunnel.sh --config config/s3-to-console.conf --dry-run connect -e local +``` + ## Supported DataSource Info | Datasource | Supported versions | diff --git a/docs/en/engines/zeta/user-command.md b/docs/en/engines/zeta/user-command.md index 433f5517af..cbe38044e6 100644 --- a/docs/en/engines/zeta/user-command.md +++ b/docs/en/engines/zeta/user-command.md @@ -83,6 +83,7 @@ The `--dry-run connect` option runs the static checks first, then uses connector | Jdbc | Yes (connectivity + schema inference) | Yes (connectivity + table existence + field compatibility) | | Kafka | Yes ([topic metadata + runtime output schema](../../connectors/source/Kafka.md#connectivity-dry-run), not consumer/group permissions) | No | | FakeSource | Yes (schema inference only, no external system) | - | +| S3File | Yes (metadata connectivity + inline schema, single-table text/csv/json/xml; see [supported scope](../../connectors/source/S3File.md#connectivity-dry-run)) | - | Every plugin in the job is reported in a validation summary with one of two statuses: diff --git a/docs/zh/connectors/source/S3File.md b/docs/zh/connectors/source/S3File.md index 479ffe6418..4e9e0a4011 100644 --- a/docs/zh/connectors/source/S3File.md +++ b/docs/zh/connectors/source/S3File.md @@ -41,6 +41,38 @@ import ChangeLog from '../changelog/connector-file-s3.md'; 从aws s3文件系统读取数据。 +### 连接预检查 + +`--dry-run connect` 可以在提交作业前检查单表 S3A 数据源。需要配置 +`bucket = "s3a://your-bucket"`、绝对 `path`、显式内联 `schema.fields` 或 +`schema.columns`,并将 `file_format_type` 设置为 `text`、`csv`、`json` 或 +`xml`。同时设置 `parse_partition_from_path = false`,且不配置 `read_columns`。 +此元数据检查不推断文件结构、投影或分区字段。不支持的配置会使连接预检查失败并 +返回说明,不改变正常作业的行为。 + +检查复用 Hadoop S3A 的端点、凭证链、代理和路径访问配置,使用 HEAD 获取对象元数据, +或执行一次 `maxKeys=1`、分隔符为 `/` 的前缀列表请求。它不会打开文件内容、递归列举 +文件、创建读取器、上传或删除对象,也不会初始化共享文件系统。检查具体对象不需要 +列举权限,检查前缀需要列举权限。对于 `discovery_mode = "continuous"`,成功的空列表 +可以通过检查,因为文件可能稍后到达。批处理模式下,没有目录标记的空虚拟前缀会失败, +但可访问的空桶根路径可以通过。桶不存在或请求被拒绝时,两种模式均失败。 + +预检查将建立连接和套接字超时限制为最多 5 秒,并保留更小的正值,禁用 SDK 请求重试, +这些限制也适用于桶级覆盖配置。它们不是 DNS、凭证提供器初始化或 SDK 初始化的总耗时 +限制,不改变正常作业的超时和重试设置。 + +首个版本不支持 `tables_configs`、旧 `s3n` 桶、SSE-C 客户提供的加密密钥、 +`fs.s3a.security.credential.provider.path`、S3Guard、分段上传清理以及自定义 S3 客户端工厂。 +通过检查不代表文件内容可读、格式或数据结构正确,也不验证工作节点凭证、目标端、更新 +或同步后操作的权限。 + +将满足上述条件的作业配置保存为 `config/s3-to-console.conf`(这是用户自行创建的文件, +并非内置模板),然后在 SeaTunnel 安装目录中执行: + +```bash +bin/seatunnel.sh --config config/s3-to-console.conf --dry-run connect -e local +``` + ## 支持的数据源信息 | 数据源 | 支持的版本 | diff --git a/docs/zh/engines/zeta/user-command.md b/docs/zh/engines/zeta/user-command.md index cded907c75..d708124799 100644 --- a/docs/zh/engines/zeta/user-command.md +++ b/docs/zh/engines/zeta/user-command.md @@ -99,6 +99,7 @@ bin/seatunnel.sh --config $SEATUNNEL_HOME/config/v2.batch.config.template --dry- | Jdbc | 支持(连通性 + schema 推断) | 支持(连通性 + 表存在性 + 字段兼容性) | | Kafka | 支持([主题元数据 + 运行时输出 schema](../../connectors/source/Kafka.md#连通性-dry-run),不含消费或消费组权限) | 不支持 | | FakeSource | 支持(仅 schema 推断,无外部系统) | - | +| S3File | 支持(元数据连通性 + 内联 schema,仅单表 text/csv/json/xml;参见[支持范围](../../connectors/source/S3File.md#连接预检查)) | - | 作业中的每个插件都会在校验汇总中报告以下两种状态之一: diff --git a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java index 083d6ee211..02a85e7cd4 100644 --- a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java +++ b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java @@ -21,8 +21,10 @@ import org.apache.seatunnel.api.configuration.util.OptionRule; import org.apache.seatunnel.api.options.ConnectorCommonOptions; import org.apache.seatunnel.api.source.SeaTunnelSource; import org.apache.seatunnel.api.source.SourceSplit; +import org.apache.seatunnel.api.table.catalog.CatalogTable; import org.apache.seatunnel.api.table.connector.TableSource; import org.apache.seatunnel.api.table.factory.Factory; +import org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation; import org.apache.seatunnel.api.table.factory.TableSourceFactory; import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext; import org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions; @@ -36,9 +38,24 @@ import com.google.auto.service.AutoService; import java.io.Serializable; import java.util.Arrays; +import java.util.List; @AutoService(Factory.class) -public class S3FileSourceFactory implements TableSourceFactory { +public class S3FileSourceFactory implements TableSourceFactory, SupportSourceDryRunValidation { + /** Returns the configured schema without discovering or opening source files. */ + @Override + public List<CatalogTable> inferSchemaForDryRun(TableSourceFactoryContext context) { + S3SourceDryRunValidator.validateSchemaOptions(context.getOptions()); + return discoverTableSchemas(context); + } + + /** Checks S3A path metadata using an independently owned client. */ + @Override + public void validateConnectionForDryRun( + TableSourceFactoryContext context, List<CatalogTable> catalogTables) throws Exception { + S3SourceDryRunValidator.validate(context.getOptions()); + } + @Override public String factoryIdentifier() { return FileSystemType.S3.getFileSystemPluginName(); diff --git a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3SourceDryRunValidator.java b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3SourceDryRunValidator.java new file mode 100644 index 0000000000..ec1da6dd5b --- /dev/null +++ b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3SourceDryRunValidator.java @@ -0,0 +1,284 @@ +/* + * 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.file.s3.source; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.options.ConnectorCommonOptions; +import org.apache.seatunnel.api.options.table.ColumnOptions; +import org.apache.seatunnel.api.options.table.FieldOptions; +import org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions; +import org.apache.seatunnel.connectors.seatunnel.file.config.FileDiscoveryMode; +import org.apache.seatunnel.connectors.seatunnel.file.config.FileFormat; +import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf; +import org.apache.seatunnel.connectors.seatunnel.file.s3.config.S3FileSourceOptions; +import org.apache.seatunnel.connectors.seatunnel.file.s3.config.S3HadoopConf; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.fs.s3a.AWSCredentialProviderList; +import org.apache.hadoop.fs.s3a.Constants; +import org.apache.hadoop.fs.s3a.DefaultS3ClientFactory; +import org.apache.hadoop.fs.s3a.S3AUtils; + +import com.amazonaws.ClientConfiguration; +import com.amazonaws.auth.AWSCredentialsProvider; +import com.amazonaws.services.s3.AmazonS3; +import com.amazonaws.services.s3.model.AmazonS3Exception; +import com.amazonaws.services.s3.model.ListObjectsRequest; +import com.amazonaws.services.s3.model.ListObjectsV2Request; +import com.amazonaws.services.s3.model.ListObjectsV2Result; +import com.amazonaws.services.s3.model.ObjectListing; + +import java.io.FileNotFoundException; +import java.io.IOException; +import java.net.URI; +import java.util.Arrays; +import java.util.Map; + +final class S3SourceDryRunValidator { + private static final int NETWORK_TIMEOUT_MILLIS = 5000; + + private S3SourceDryRunValidator() {} + + /** Rejects schemas that require opening files or altering the configured field layout. */ + static void validateSchemaOptions(ReadonlyConfig options) { + if (options.getOptional(ConnectorCommonOptions.TABLE_CONFIGS).isPresent() + || !options.getOptional(ConnectorCommonOptions.SCHEMA).isPresent()) { + throw new IllegalArgumentException( + "S3File connect dry-run requires a single table with an explicit inline schema."); + } + Map<String, Object> schema = options.get(ConnectorCommonOptions.SCHEMA); + if (schema.containsKey(ColumnOptions.METADATA_TABLE_ID.key()) + || (!options.getOptional(FieldOptions.FIELDS).isPresent() + && !schema.containsKey(ColumnOptions.COLUMNS.key()))) { + throw new IllegalArgumentException( + "S3File connect dry-run requires schema.fields or schema.columns, not metadata schema lookup."); + } + if (!Arrays.asList(FileFormat.TEXT, FileFormat.CSV, FileFormat.JSON, FileFormat.XML) + .contains(options.get(FileBaseSourceOptions.FILE_FORMAT_TYPE))) { + throw new IllegalArgumentException( + "S3File connect dry-run supports text, csv, json and xml without reading file contents."); + } + if (options.get(FileBaseSourceOptions.PARSE_PARTITION_FROM_PATH) + || options.getOptional(FileBaseSourceOptions.READ_COLUMNS).isPresent()) { + throw new IllegalArgumentException( + "S3File connect dry-run requires parse_partition_from_path=false and no read_columns to preserve the runtime schema."); + } + } + + /** Checks metadata with a private client, closing credentials even when client setup fails. */ + static void validate(ReadonlyConfig options) throws IOException { + // The connection hook can be called directly, without the preceding schema hook. + validateSchemaOptions(options); + URI bucket = bucketUri(options); + Path path = sourcePath(options, bucket); + Configuration configuration = validationConfiguration(options, bucket); + // Validate pure SDK settings before Hadoop creates any credential providers. + DefaultS3ClientFactory.createAwsConf(configuration); + try (DryRunClientFactory factory = new DryRunClientFactory()) { + factory.setConf(configuration); + AmazonS3 client = factory.createS3Client(bucket); + validatePath( + client, + bucket.getHost(), + path, + options.get(FileBaseSourceOptions.DISCOVERY_MODE), + configuration.getInt(Constants.LIST_VERSION, Constants.DEFAULT_LIST_VERSION)); + } + } + + /** Accepts only a bare S3A bucket URI, without echoing potentially sensitive input. */ + static URI bucketUri(ReadonlyConfig options) { + URI bucket; + try { + bucket = URI.create(options.get(S3FileSourceOptions.S3_BUCKET)); + } catch (IllegalArgumentException e) { + throw new IllegalArgumentException( + "S3File connect dry-run requires an s3a://bucket URI."); + } + if (!"s3a".equals(bucket.getScheme()) + || bucket.getHost() == null + || bucket.getUserInfo() != null + || bucket.getPort() != -1 + || bucket.getQuery() != null + || bucket.getFragment() != null + || (!bucket.getPath().isEmpty() && !"/".equals(bucket.getPath()))) { + throw new IllegalArgumentException( + "S3File connect dry-run requires an s3a://bucket URI."); + } + return bucket; + } + + /** Keeps metadata requests inside the configured bucket, using the source's absolute path. */ + static Path sourcePath(ReadonlyConfig options, URI bucket) { + Path path; + try { + path = new Path(options.get(FileBaseSourceOptions.FILE_PATH)); + } catch (IllegalArgumentException e) { + throw new IllegalArgumentException( + "S3File connect dry-run requires an absolute path in the configured bucket."); + } + URI uri = path.toUri(); + if (!path.isAbsolute() + || (uri.getScheme() != null && !"s3a".equals(uri.getScheme())) + || (uri.getAuthority() != null + && !bucket.getAuthority().equals(uri.getAuthority()))) { + throw new IllegalArgumentException( + "S3File connect dry-run requires an absolute path in the configured bucket."); + } + return path; + } + + /** + * Preserves S3A connection settings but rejects filesystem side effects, custom client + * lifecycle behavior, and encryption/credential-store paths not supported by this metadata-only + * client. Disables retries so a preflight check does not inherit the production retry budget. + */ + static Configuration validationConfiguration(ReadonlyConfig options, URI bucket) { + HadoopConf hadoopConf = S3HadoopConf.buildWithReadOnlyConfig(options); + Configuration configuration = hadoopConf.toConfiguration(); + hadoopConf.setExtraOptionsForConfiguration(configuration); + configuration = S3AUtils.propagateBucketOptions(configuration, bucket.getHost()); + if (configuration.getBoolean(Constants.PURGE_EXISTING_MULTIPART, false) + || !Constants.S3GUARD_METASTORE_NULL.equals( + configuration.get( + Constants.S3_METADATA_STORE_IMPL, Constants.S3GUARD_METASTORE_NULL)) + || !Constants.DEFAULT_S3_CLIENT_FACTORY_IMPL + .getName() + .equals( + configuration.get( + Constants.S3_CLIENT_FACTORY_IMPL, + Constants.DEFAULT_S3_CLIENT_FACTORY_IMPL.getName()))) { + throw new IllegalArgumentException( + "S3File connect dry-run does not support multipart purge, S3Guard or a custom S3 client factory."); + } + if (!configuration.getTrimmed(Constants.S3A_SECURITY_CREDENTIAL_PROVIDER_PATH, "").isEmpty() + || "SSE-C" + .equalsIgnoreCase( + configuration.getTrimmed( + Constants.SERVER_SIDE_ENCRYPTION_ALGORITHM, ""))) { + throw new IllegalArgumentException( + "S3File connect dry-run does not support SSE-C or fs.s3a.security.credential.provider.path."); + } + clampTimeout(configuration, bucket, Constants.ESTABLISH_TIMEOUT); + clampTimeout(configuration, bucket, Constants.SOCKET_TIMEOUT); + for (String key : + Arrays.asList( + Constants.MAX_ERROR_RETRIES, + Constants.RETRY_LIMIT, + Constants.RETRY_THROTTLE_LIMIT)) { + S3AUtils.clearBucketOption(configuration, bucket.getHost(), key); + configuration.setInt(key, 0); + } + return configuration; + } + + /** Caps each network wait at five seconds for preflight, retaining stricter user limits. */ + private static void clampTimeout(Configuration configuration, URI bucket, String key) { + int configured = configuration.getInt(key, NETWORK_TIMEOUT_MILLIS); + if (configured < 0) { + throw new IllegalArgumentException( + "S3File connect dry-run requires non-negative " + key); + } + S3AUtils.clearBucketOption(configuration, bucket.getHost(), key); + configuration.setInt( + key, + configured == 0 + ? NETWORK_TIMEOUT_MILLIS + : Math.min(configured, NETWORK_TIMEOUT_MILLIS)); + } + + /** + * Tries HEAD first so exact objects need no listing permission. Only a missing object falls + * back to one bounded prefix listing; permission and other failures must not be treated as + * absence. + */ + static void validatePath( + AmazonS3 client, + String bucket, + Path path, + FileDiscoveryMode discoveryMode, + int listVersion) + throws IOException { + String key = path.toUri().getPath().substring(1); + if (!key.isEmpty()) { + try { + client.getObjectMetadata(bucket, key); + return; + } catch (AmazonS3Exception e) { + if (e.getStatusCode() != 404) { + throw e; + } + } + } + String prefix = key.isEmpty() || key.endsWith("/") ? key : key + "/"; + boolean empty; + if (listVersion == 1) { + ObjectListing result = + client.listObjects( + new ListObjectsRequest() + .withBucketName(bucket) + .withPrefix(prefix) + .withDelimiter("/") + .withMaxKeys(1)); + empty = result.getObjectSummaries().isEmpty() && result.getCommonPrefixes().isEmpty(); + } else { + ListObjectsV2Result result = + client.listObjectsV2( + new ListObjectsV2Request() + .withBucketName(bucket) + .withPrefix(prefix) + .withDelimiter("/") + .withMaxKeys(1)); + empty = result.getObjectSummaries().isEmpty() && result.getCommonPrefixes().isEmpty(); + } + // Only a successful empty listing proves an accessible prefix awaiting its first file. + if (empty && !key.isEmpty() && discoveryMode != FileDiscoveryMode.CONTINUOUS) { + throw new FileNotFoundException( + "S3File connect dry-run could not find the configured batch source path."); + } + } + + /** Owns the client and credential chain without creating a shared S3A filesystem. */ + private static final class DryRunClientFactory extends DefaultS3ClientFactory + implements AutoCloseable { + private AmazonS3 client; + private AWSCredentialProviderList credentials; + + @Override + protected AmazonS3 newAmazonS3Client( + AWSCredentialsProvider credentials, ClientConfiguration configuration) { + this.credentials = (AWSCredentialProviderList) credentials; + this.client = super.newAmazonS3Client(credentials, configuration); + return client; + } + + @Override + public void close() { + try { + if (client != null) { + client.shutdown(); + } + } finally { + if (credentials != null) { + credentials.close(); + } + } + } + } +} diff --git a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/S3FileSourceFactoryTest.java b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/S3FileSourceFactoryTest.java index a60221fbc2..c222b7d19a 100644 --- a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/S3FileSourceFactoryTest.java +++ b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/S3FileSourceFactoryTest.java @@ -21,6 +21,7 @@ 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.api.table.factory.SupportSourceDryRunValidation; import org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions; import org.apache.seatunnel.connectors.seatunnel.file.s3.config.S3FileSourceOptions; import org.apache.seatunnel.connectors.seatunnel.file.s3.source.S3FileSourceFactory; @@ -36,6 +37,13 @@ import static org.junit.jupiter.api.Assertions.assertTrue; class S3FileSourceFactoryTest { + @Test + void shouldSupportConnectivityDryRun() { + assertTrue( + new S3FileSourceFactory() instanceof SupportSourceDryRunValidation, + "S3File source connectivity is skipped when the factory does not expose the dry-run SPI"); + } + @Test void shouldExposeContinuousDiscoveryOptions() { OptionRule optionRule = new S3FileSourceFactory().optionRule(); diff --git a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3SourceDryRunValidatorTest.java b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3SourceDryRunValidatorTest.java new file mode 100644 index 0000000000..5be35b9230 --- /dev/null +++ b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3SourceDryRunValidatorTest.java @@ -0,0 +1,463 @@ +/* + * 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.file.s3.source; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.table.catalog.CatalogTable; +import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext; +import org.apache.seatunnel.connectors.seatunnel.file.config.FileDiscoveryMode; +import org.apache.seatunnel.connectors.seatunnel.file.s3.config.S3FileSourceOptions; + +import org.apache.hadoop.conf.Configuration; +import org.apache.hadoop.fs.Path; +import org.apache.hadoop.fs.s3a.Constants; +import org.apache.hadoop.fs.s3a.S3AUtils; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.params.ParameterizedTest; +import org.junit.jupiter.params.provider.ValueSource; + +import com.amazonaws.services.s3.AmazonS3; +import com.amazonaws.services.s3.model.AmazonS3Exception; +import com.amazonaws.services.s3.model.ListObjectsRequest; +import com.amazonaws.services.s3.model.ListObjectsV2Request; +import com.amazonaws.services.s3.model.ListObjectsV2Result; +import com.amazonaws.services.s3.model.ObjectListing; +import com.amazonaws.services.s3.model.ObjectMetadata; +import com.amazonaws.services.s3.model.S3ObjectSummary; + +import java.io.FileNotFoundException; +import java.io.IOException; +import java.net.URI; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.verifyNoMoreInteractions; +import static org.mockito.Mockito.when; + +class S3SourceDryRunValidatorTest { + private static final URI BUCKET = URI.create("s3a://warehouse"); + + @ParameterizedTest + @ValueSource(strings = {"text", "csv", "json", "xml"}) + void shouldReturnConfiguredSchemaWithoutCreatingSource(String format) { + Map<String, Object> config = sourceConfig(); + config.put("file_format_type", format); + TableSourceFactoryContext context = context(config); + S3FileSourceFactory factory = new S3FileSourceFactory(); + + List<CatalogTable> tables = factory.inferSchemaForDryRun(context); + + assertEquals(1, tables.size()); + assertEquals( + factory.discoverTableSchemas(context).get(0).getSeaTunnelRowType(), + tables.get(0).getSeaTunnelRowType()); + assertEquals("event_id", tables.get(0).getSeaTunnelRowType().getFieldName(0)); + } + + @ParameterizedTest + @ValueSource(strings = {"binary", "markdown", "pdf", "parquet", "orc", "excel"}) + void shouldRefuseFormatsRequiringDifferentSchemaHandling(String format) { + Map<String, Object> config = sourceConfig(); + config.put("file_format_type", format); + assertSchemaFailure(config, "supports text, csv, json and xml"); + } + + @Test + void shouldRefuseMissingSchema() { + Map<String, Object> config = sourceConfig(); + config.remove("schema"); + assertSchemaFailure(config, "explicit inline schema"); + } + + @Test + void shouldRefuseMetadataSchemaLookup() { + Map<String, Object> config = sourceConfig(); + config.put("schema", Collections.singletonMap("metadata_table_id", "warehouse.events")); + assertSchemaFailure(config, "not metadata schema lookup"); + } + + @Test + void shouldRefuseTableConfigs() { + Map<String, Object> config = sourceConfig(); + config.put("tables_configs", Collections.singletonList(sourceConfig())); + assertSchemaFailure(config, "single table"); + } + + @Test + void shouldRefuseDefaultPartitionInference() { + Map<String, Object> config = sourceConfig(); + config.remove("parse_partition_from_path"); + assertSchemaFailure(config, "parse_partition_from_path=false"); + } + + @Test + void shouldRefuseProjection() { + Map<String, Object> config = sourceConfig(); + config.put("read_columns", Collections.singletonList("event_id")); + assertSchemaFailure(config, "no read_columns"); + } + + @ParameterizedTest + @ValueSource( + strings = { + "s3n://warehouse", + "s3a://warehouse/folder", + "s3a://secret@warehouse", + "s3a://warehouse?key=secret" + }) + void shouldRefuseUnsupportedBucketUriWithoutEchoingIt(String bucket) { + Map<String, Object> config = sourceConfig(); + config.put("bucket", bucket); + IllegalArgumentException failure = + assertThrows( + IllegalArgumentException.class, + () -> S3SourceDryRunValidator.bucketUri(ReadonlyConfig.fromMap(config))); + assertTrue(failure.getMessage().contains("s3a://bucket URI")); + assertFalse(failure.getMessage().contains("secret")); + } + + @ParameterizedTest + @ValueSource( + strings = { + "events", + "s3a://other/events", + "file:///events", + "s3a://secret@warehouse/events", + "s3a://[secret" + }) + void shouldRefusePathOutsideConfiguredBucket(String path) { + Map<String, Object> config = sourceConfig(); + config.put("path", path); + IllegalArgumentException failure = + assertThrows( + IllegalArgumentException.class, + () -> + S3SourceDryRunValidator.sourcePath( + ReadonlyConfig.fromMap(config), BUCKET)); + assertTrue(failure.getMessage().contains("absolute path in the configured bucket")); + assertFalse(failure.getMessage().contains("secret")); + } + + @Test + void shouldRetainEndpointAndCredentialsWithoutChangingInputProperties() { + Map<String, Object> config = sourceConfig(); + Map<String, String> properties = new HashMap<>(); + properties.put(Constants.PATH_STYLE_ACCESS, "true"); + properties.put(Constants.SOCKET_TIMEOUT, "1200"); + config.put("hadoop_s3_properties", properties); + + Configuration configuration = configuration(config); + + assertEquals("http://localhost:9000", configuration.get(Constants.ENDPOINT)); + assertEquals("access-key", configuration.get(Constants.ACCESS_KEY)); + assertEquals("secret-key", configuration.get(Constants.SECRET_KEY)); + assertEquals( + S3FileSourceOptions.SIMPLE_AWS_CREDENTIALS_PROVIDER, + configuration.get(Constants.AWS_CREDENTIALS_PROVIDER)); + assertEquals(1200, configuration.getInt(Constants.SOCKET_TIMEOUT, -1)); + assertTrue(configuration.getBoolean(Constants.PATH_STYLE_ACCESS, false)); + assertEquals(2, properties.size()); + assertEquals("1200", properties.get(Constants.SOCKET_TIMEOUT)); + } + + @Test + void shouldBoundBucketSpecificTimeoutsAndRetriesAfterS3AReappliesOverrides() { + Map<String, Object> config = sourceConfig(); + Map<String, String> properties = new HashMap<>(); + properties.put("fs.s3a.bucket.warehouse.connection.timeout", "0"); + properties.put("fs.s3a.bucket.warehouse.connection.establish.timeout", "2000"); + properties.put("fs.s3a.bucket.warehouse.attempts.maximum", "99"); + properties.put("fs.s3a.bucket.warehouse.retry.limit", "99"); + properties.put("fs.s3a.bucket.warehouse.retry.throttle.limit", "99"); + config.put("hadoop_s3_properties", properties); + + Configuration configuration = + S3AUtils.propagateBucketOptions(configuration(config), "warehouse"); + + assertEquals(5000, configuration.getInt(Constants.SOCKET_TIMEOUT, -1)); + assertEquals(2000, configuration.getInt(Constants.ESTABLISH_TIMEOUT, -1)); + assertEquals(0, configuration.getInt(Constants.MAX_ERROR_RETRIES, -1)); + assertEquals(0, configuration.getInt(Constants.RETRY_LIMIT, -1)); + assertEquals(0, configuration.getInt(Constants.RETRY_THROTTLE_LIMIT, -1)); + assertEquals("99", properties.get("fs.s3a.bucket.warehouse.retry.limit")); + } + + @ParameterizedTest + @ValueSource(strings = {"fs.s3a.multipart.purge", "fs.s3a.bucket.warehouse.multipart.purge"}) + void shouldRefuseMultipartPurgeBeforeOpeningFilesystem(String key) { + assertConfigurationFailure(key, "true", "multipart purge"); + } + + @ParameterizedTest + @ValueSource( + strings = {"fs.s3a.metadatastore.impl", "fs.s3a.bucket.warehouse.metadatastore.impl"}) + void shouldRefuseMetadataStoreInitialization(String key) { + assertConfigurationFailure(key, Constants.S3GUARD_METASTORE_DYNAMO, "S3Guard"); + } + + @ParameterizedTest + @ValueSource( + strings = { + "fs.s3a.s3.client.factory.impl", + "fs.s3a.bucket.warehouse.s3.client.factory.impl" + }) + void shouldRefuseCustomClientInitialization(String key) { + assertConfigurationFailure(key, "example.CustomClientFactory", "custom S3 client factory"); + } + + @Test + void shouldRefuseNegativeTimeout() { + assertConfigurationFailure(Constants.SOCKET_TIMEOUT, "-1", "non-negative"); + } + + @Test + void shouldCheckOnlyPathMetadata() throws IOException { + AmazonS3 client = mock(AmazonS3.class); + Path path = new Path("/events"); + when(client.getObjectMetadata("warehouse", "events")).thenReturn(new ObjectMetadata()); + S3SourceDryRunValidator.validatePath(client, "warehouse", path, FileDiscoveryMode.ONCE, 2); + verify(client).getObjectMetadata("warehouse", "events"); + verifyNoMoreInteractions(client); + } + + @Test + void shouldFailForMissingBatchPath() throws IOException { + AmazonS3 client = emptyPrefixClient(); + Path path = new Path("/events"); + FileNotFoundException failure = + assertThrows( + FileNotFoundException.class, + () -> + S3SourceDryRunValidator.validatePath( + client, "warehouse", path, FileDiscoveryMode.ONCE, 2)); + assertTrue(failure.getMessage().contains("batch source path")); + } + + @Test + void shouldAllowEmptyContinuousPrefix() throws IOException { + AmazonS3 client = emptyPrefixClient(); + Path path = new Path("/events"); + assertDoesNotThrow( + () -> + S3SourceDryRunValidator.validatePath( + client, "warehouse", path, FileDiscoveryMode.CONTINUOUS, 2)); + } + + @Test + void shouldPropagateDeniedAccessForContinuousDiscovery() throws IOException { + AmazonS3 client = emptyPrefixClient(); + Path path = new Path("/events"); + AmazonS3Exception denied = serviceFailure(403); + when(client.listObjectsV2(any(ListObjectsV2Request.class))).thenThrow(denied); + AmazonS3Exception failure = + assertThrows( + AmazonS3Exception.class, + () -> + S3SourceDryRunValidator.validatePath( + client, + "warehouse", + path, + FileDiscoveryMode.CONTINUOUS, + 2)); + assertSame(denied, failure); + assertTrue(failure.getMessage().contains("metadata failure")); + } + + @Test + void shouldFailForMissingBucketDuringContinuousDiscovery() { + AmazonS3 client = emptyPrefixClient(); + AmazonS3Exception missing = serviceFailure(404); + when(client.listObjectsV2(any(ListObjectsV2Request.class))).thenThrow(missing); + AmazonS3Exception failure = + assertThrows( + AmazonS3Exception.class, + () -> + S3SourceDryRunValidator.validatePath( + client, + "warehouse", + new Path("/events"), + FileDiscoveryMode.CONTINUOUS, + 2)); + assertSame(missing, failure); + assertEquals(404, failure.getStatusCode()); + } + + @Test + void shouldBoundPrefixListingWithoutFollowingContinuation() throws IOException { + AmazonS3 client = emptyPrefixClient(); + when(client.listObjectsV2(any(ListObjectsV2Request.class))) + .thenAnswer( + invocation -> { + ListObjectsV2Request request = invocation.getArgument(0); + assertEquals("warehouse", request.getBucketName()); + assertEquals("events/", request.getPrefix()); + assertEquals("/", request.getDelimiter()); + assertEquals(1, request.getMaxKeys()); + ListObjectsV2Result result = new ListObjectsV2Result(); + result.getObjectSummaries().add(new S3ObjectSummary()); + result.setTruncated(true); + result.setNextContinuationToken("more-events"); + return result; + }); + S3SourceDryRunValidator.validatePath( + client, "warehouse", new Path("/events"), FileDiscoveryMode.ONCE, 2); + verify(client).getObjectMetadata("warehouse", "events"); + verify(client).listObjectsV2(any(ListObjectsV2Request.class)); + verifyNoMoreInteractions(client); + } + + @Test + void shouldRefuseCustomerProvidedEncryptionKeys() { + assertConfigurationFailure(Constants.SERVER_SIDE_ENCRYPTION_ALGORITHM, "SSE-C", "SSE-C"); + } + + @Test + void shouldHonorLegacyListVersion() throws IOException { + AmazonS3 client = mock(AmazonS3.class); + when(client.getObjectMetadata("warehouse", "events")).thenThrow(serviceFailure(404)); + when(client.listObjects(any(ListObjectsRequest.class))) + .thenAnswer( + invocation -> { + ListObjectsRequest request = invocation.getArgument(0); + assertEquals(Integer.valueOf(1), request.getMaxKeys()); + assertEquals("events/", request.getPrefix()); + assertEquals("/", request.getDelimiter()); + return new ObjectListing(); + }); + S3SourceDryRunValidator.validatePath( + client, "warehouse", new Path("/events"), FileDiscoveryMode.CONTINUOUS, 1); + verify(client).getObjectMetadata("warehouse", "events"); + verify(client).listObjects(any(ListObjectsRequest.class)); + verifyNoMoreInteractions(client); + } + + @ParameterizedTest + @ValueSource(ints = {0, 3}) + void shouldMatchS3AFallbackForOtherListVersions(int listVersion) throws IOException { + AmazonS3 client = emptyPrefixClient(); + S3SourceDryRunValidator.validatePath( + client, + "warehouse", + new Path("/events"), + FileDiscoveryMode.CONTINUOUS, + listVersion); + verify(client).getObjectMetadata("warehouse", "events"); + verify(client).listObjectsV2(any(ListObjectsV2Request.class)); + verifyNoMoreInteractions(client); + } + + @Test + void shouldRefuseS3SpecificCredentialStoreBeforeClientInitialization() { + assertConfigurationFailure( + Constants.S3A_SECURITY_CREDENTIAL_PROVIDER_PATH, + "jceks://file/credentials", + "credential.provider.path"); + } + + private static AmazonS3 emptyPrefixClient() { + AmazonS3 client = mock(AmazonS3.class); + when(client.getObjectMetadata("warehouse", "events")).thenThrow(serviceFailure(404)); + when(client.listObjectsV2(any(ListObjectsV2Request.class))) + .thenReturn(new ListObjectsV2Result()); + return client; + } + + @Test + void shouldReturnExplicitColumnSchema() { + Map<String, Object> config = sourceConfig(); + Map<String, Object> column = new HashMap<>(); + column.put("name", "event_id"); + column.put("type", "bigint"); + config.put( + "schema", Collections.singletonMap("columns", Collections.singletonList(column))); + CatalogTable table = new S3FileSourceFactory().inferSchemaForDryRun(context(config)).get(0); + assertEquals("event_id", table.getSeaTunnelRowType().getFieldName(0)); + } + + @Test + void shouldRetainDefaultCredentialProvider() { + Map<String, Object> config = sourceConfig(); + config.remove("fs.s3a.aws.credentials.provider"); + Configuration configuration = configuration(config); + assertEquals( + S3FileSourceOptions.INSTANCE_PROFILE_CREDENTIALS_PROVIDER, + configuration.get(Constants.AWS_CREDENTIALS_PROVIDER)); + } + + private static AmazonS3Exception serviceFailure(int status) { + AmazonS3Exception failure = new AmazonS3Exception("metadata failure"); + failure.setStatusCode(status); + return failure; + } + + static Map<String, Object> sourceConfig() { + Map<String, Object> config = new HashMap<>(); + config.put("path", "/events"); + config.put("file_format_type", "json"); + config.put("bucket", BUCKET.toString()); + config.put("fs.s3a.endpoint", "http://localhost:9000"); + config.put( + "fs.s3a.aws.credentials.provider", + S3FileSourceOptions.SIMPLE_AWS_CREDENTIALS_PROVIDER); + config.put("access_key", "access-key"); + config.put("secret_key", "secret-key"); + config.put("parse_partition_from_path", false); + config.put( + "schema", + Collections.singletonMap("fields", Collections.singletonMap("event_id", "bigint"))); + return config; + } + + private static TableSourceFactoryContext context(Map<String, Object> config) { + return new TableSourceFactoryContext( + ReadonlyConfig.fromMap(config), S3SourceDryRunValidatorTest.class.getClassLoader()); + } + + private static Configuration configuration(Map<String, Object> config) { + return S3SourceDryRunValidator.validationConfiguration( + ReadonlyConfig.fromMap(config), BUCKET); + } + + private static void assertSchemaFailure(Map<String, Object> config, String message) { + IllegalArgumentException failure = + assertThrows( + IllegalArgumentException.class, + () -> new S3FileSourceFactory().inferSchemaForDryRun(context(config))); + assertTrue(failure.getMessage().contains(message), failure.getMessage()); + } + + private static void assertConfigurationFailure(String key, String value, String message) { + Map<String, Object> config = sourceConfig(); + config.put("hadoop_s3_properties", Collections.singletonMap(key, value)); + IllegalArgumentException failure = + assertThrows(IllegalArgumentException.class, () -> configuration(config)); + assertTrue(failure.getMessage().contains(message), failure.getMessage()); + } +} diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileConnectDryRunIT.java b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileConnectDryRunIT.java new file mode 100644 index 0000000000..954f64013c --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileConnectDryRunIT.java @@ -0,0 +1,190 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.e2e.connector.file.s3; + +import org.apache.seatunnel.api.configuration.ReadonlyConfig; +import org.apache.seatunnel.api.table.factory.FactoryUtil; +import org.apache.seatunnel.api.table.factory.SupportSourceDryRunValidation; +import org.apache.seatunnel.api.table.factory.TableSourceFactory; +import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext; +import org.apache.seatunnel.e2e.common.TestResource; +import org.apache.seatunnel.e2e.common.TestSuiteBase; + +import org.junit.jupiter.api.AfterAll; +import org.junit.jupiter.api.BeforeAll; +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; +import org.testcontainers.containers.GenericContainer; +import org.testcontainers.containers.output.Slf4jLogConsumer; +import org.testcontainers.containers.wait.strategy.Wait; +import org.testcontainers.utility.DockerImageName; +import org.testcontainers.utility.DockerLoggerFactory; + +import com.amazonaws.auth.AWSStaticCredentialsProvider; +import com.amazonaws.auth.BasicAWSCredentials; +import com.amazonaws.client.builder.AwsClientBuilder; +import com.amazonaws.services.s3.AmazonS3; +import com.amazonaws.services.s3.AmazonS3ClientBuilder; +import com.amazonaws.services.s3.model.AmazonS3Exception; + +import java.io.FileNotFoundException; +import java.time.Duration; +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** Exercises the source factory metadata contract against MinIO without submitting a job. */ +@Timeout(60) +public class S3FileConnectDryRunIT extends TestSuiteBase implements TestResource { + private static final String IMAGE = "minio/minio:RELEASE.2024-06-13T22-53-53Z"; + private static final String BUCKET = "dry-run-events"; + private static final String ACCESS_KEY = "minioadmin"; + private static final String SECRET_KEY = "minioadmin"; + private GenericContainer<?> minio; + private AmazonS3 admin; + private String endpoint; + + @BeforeAll + @Override + public void startUp() { + minio = + new GenericContainer<>(DockerImageName.parse(IMAGE)) + .withEnv("MINIO_ROOT_USER", ACCESS_KEY) + .withEnv("MINIO_ROOT_PASSWORD", SECRET_KEY) + .withCommand("server", "/data") + .withExposedPorts(9000) + .waitingFor(Wait.forHttp("/minio/health/ready").forPort(9000)) + .withStartupTimeout(Duration.ofMinutes(2)) + .withLogConsumer( + new Slf4jLogConsumer(DockerLoggerFactory.getLogger(IMAGE))); + minio.start(); + endpoint = "http://" + minio.getHost() + ":" + minio.getMappedPort(9000); + admin = + AmazonS3ClientBuilder.standard() + .withCredentials( + new AWSStaticCredentialsProvider( + new BasicAWSCredentials(ACCESS_KEY, SECRET_KEY))) + .withEndpointConfiguration( + new AwsClientBuilder.EndpointConfiguration(endpoint, "us-east-1")) + .withPathStyleAccessEnabled(true) + .build(); + admin.createBucket(BUCKET); + admin.putObject(BUCKET, "events/data.json", "{\"event_id\":1}"); + admin.createBucket("dry-run-empty"); + } + + @AfterAll + @Override + public void tearDown() { + try { + if (admin != null) { + admin.shutdown(); + } + } finally { + if (minio != null) { + minio.stop(); + } + } + } + + @Test + void shouldValidateObjectAndPrefixWithoutChangingStoredData() throws Exception { + String etag = admin.getObjectMetadata(BUCKET, "events/data.json").getETag(); + validate(sourceConfig(endpoint, BUCKET, "/events/data.json")); + validate(sourceConfig(endpoint, BUCKET, "/events")); + assertEquals(etag, admin.getObjectMetadata(BUCKET, "events/data.json").getETag()); + assertEquals("{\"event_id\":1}", admin.getObjectAsString(BUCKET, "events/data.json")); + assertEquals(1, admin.listObjectsV2(BUCKET).getKeyCount()); + } + + @Test + void shouldRejectMissingBatchPrefix() { + FileNotFoundException failure = + assertThrows( + FileNotFoundException.class, + () -> validate(sourceConfig(endpoint, BUCKET, "/future"))); + assertTrue(failure.getMessage().contains("batch source path")); + } + + @Test + void shouldAcceptEmptyContinuousPrefix() { + Map<String, Object> config = sourceConfig(endpoint, BUCKET, "/future"); + config.put("discovery_mode", "continuous"); + assertDoesNotThrow(() -> validate(config)); + } + + @Test + void shouldAcceptAccessibleEmptyBucketRoot() { + assertDoesNotThrow(() -> validate(sourceConfig(endpoint, "dry-run-empty", "/"))); + } + + @Test + void shouldRejectMissingBucketInContinuousMode() { + Map<String, Object> config = sourceConfig(endpoint, "dry-run-missing", "/future"); + config.put("discovery_mode", "continuous"); + AmazonS3Exception failure = assertThrows(AmazonS3Exception.class, () -> validate(config)); + assertEquals(404, failure.getStatusCode()); + assertEquals("NoSuchBucket", failure.getErrorCode()); + } + + @Test + void shouldRejectIncorrectCredentials() { + Map<String, Object> config = sourceConfig(endpoint, BUCKET, "/events/data.json"); + config.put("secret_key", "incorrect-test-secret"); + AmazonS3Exception failure = assertThrows(AmazonS3Exception.class, () -> validate(config)); + assertEquals(403, failure.getStatusCode()); + assertTrue(failure.getMessage().contains("403")); + } + + static Map<String, Object> sourceConfig(String endpoint, String bucket, String path) { + Map<String, Object> config = new HashMap<>(); + config.put("path", path); + config.put("file_format_type", "json"); + config.put("bucket", "s3a://" + bucket); + config.put("fs.s3a.endpoint", endpoint); + config.put( + "fs.s3a.aws.credentials.provider", + "org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"); + config.put("access_key", ACCESS_KEY); + config.put("secret_key", SECRET_KEY); + config.put("parse_partition_from_path", false); + config.put( + "hadoop_s3_properties", + Collections.singletonMap("fs.s3a.path.style.access", "true")); + config.put( + "schema", + Collections.singletonMap("fields", Collections.singletonMap("event_id", "bigint"))); + return config; + } + + static void validate(Map<String, Object> config) throws Exception { + ClassLoader classLoader = S3FileConnectDryRunIT.class.getClassLoader(); + TableSourceFactory factory = + FactoryUtil.discoverFactory(classLoader, TableSourceFactory.class, "S3File"); + assertTrue(factory instanceof SupportSourceDryRunValidation); + SupportSourceDryRunValidation validator = (SupportSourceDryRunValidation) factory; + TableSourceFactoryContext context = + new TableSourceFactoryContext(ReadonlyConfig.fromMap(config), classLoader); + validator.validateConnectionForDryRun(context, validator.inferSchemaForDryRun(context)); + } +} diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileConnectDryRunWireIT.java b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileConnectDryRunWireIT.java new file mode 100644 index 0000000000..2cedd84c2c --- /dev/null +++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileConnectDryRunWireIT.java @@ -0,0 +1,249 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.seatunnel.e2e.connector.file.s3; + +import org.apache.seatunnel.e2e.common.TestSuiteBase; + +import org.junit.jupiter.api.Test; +import org.junit.jupiter.api.Timeout; + +import com.amazonaws.AmazonClientException; +import com.amazonaws.auth.AWSCredentials; +import com.amazonaws.auth.AWSCredentialsProvider; +import com.amazonaws.auth.BasicAWSCredentials; +import com.amazonaws.services.s3.model.AmazonS3Exception; +import com.sun.net.httpserver.HttpExchange; +import com.sun.net.httpserver.HttpServer; + +import java.io.Closeable; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.nio.charset.StandardCharsets; +import java.util.ArrayList; +import java.util.Arrays; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.concurrent.ConcurrentLinkedQueue; +import java.util.concurrent.CountDownLatch; +import java.util.concurrent.ExecutorService; +import java.util.concurrent.Executors; +import java.util.concurrent.atomic.AtomicInteger; + +import static org.junit.jupiter.api.Assertions.assertDoesNotThrow; +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertThrows; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** Records actual SDK requests; no SeaTunnel jobs or object-content requests are used. */ +@Timeout(30) +public class S3FileConnectDryRunWireIT extends TestSuiteBase { + @Test + void shouldUseOnlyObjectHeadForExactPath() throws Exception { + try (MetadataServer server = new MetadataServer("object")) { + S3FileConnectDryRunIT.validate(server.config()); + assertEquals(Arrays.asList("HEAD /dry-run-events/events"), server.requests()); + } + } + + @Test + void shouldListOnlyOnePrefixEntryWithoutFetchingContentsOrNextPage() throws Exception { + try (MetadataServer server = new MetadataServer("prefix")) { + assertDoesNotThrow( + () -> S3FileConnectDryRunIT.validate(server.config()), + () -> server.requests().toString()); + List<String> requests = server.requests(); + assertEquals(2, requests.size()); + assertEquals("HEAD /dry-run-events/events", requests.get(0)); + assertTrue(requests.get(1).startsWith("GET /dry-run-events/?"), requests.toString()); + assertTrue(requests.get(1).contains("max-keys=1")); + assertTrue(requests.get(1).contains("prefix=events%2F")); + assertTrue(requests.get(1).contains("delimiter=%2F")); + assertFalse(requests.get(1).contains("continuation-token")); + } + } + + @Test + void shouldPropagateForbiddenMetadataWithoutRetrying() throws Exception { + try (MetadataServer server = new MetadataServer("forbidden")) { + AmazonS3Exception failure = + assertThrows( + AmazonS3Exception.class, + () -> S3FileConnectDryRunIT.validate(server.config())); + assertEquals(403, failure.getStatusCode()); + assertEquals(1, server.requests().size()); + } + } + + @Test + void shouldUseConfiguredSmallerSocketTimeoutAndCloseCredentialsOnFailure() throws Exception { + try (MetadataServer server = new MetadataServer("stalled")) { + Map<String, Object> config = server.config(); + Map<String, String> properties = new HashMap<>(); + properties.put("fs.s3a.path.style.access", "true"); + properties.put("fs.s3a.connection.timeout", "200"); + config.put("hadoop_s3_properties", properties); + config.put( + "fs.s3a.aws.credentials.provider", TrackingCredentialsProvider.class.getName()); + int closedBefore = TrackingCredentialsProvider.CLOSED.get(); + AmazonClientException failure = + assertThrows( + AmazonClientException.class, + () -> S3FileConnectDryRunIT.validate(config)); + assertTrue(failure.getMessage().contains("Unable to execute HTTP request")); + assertEquals(1, server.requests().size()); + assertEquals(closedBefore + 1, TrackingCredentialsProvider.CLOSED.get()); + } + } + + @Test + void shouldCloseCredentialsAfterSuccessfulMetadataRequest() throws Exception { + try (MetadataServer server = new MetadataServer("object")) { + Map<String, Object> config = server.config(); + config.put( + "fs.s3a.aws.credentials.provider", TrackingCredentialsProvider.class.getName()); + int closedBefore = TrackingCredentialsProvider.CLOSED.get(); + S3FileConnectDryRunIT.validate(config); + assertEquals(closedBefore + 1, TrackingCredentialsProvider.CLOSED.get()); + } + } + + public static class TrackingCredentialsProvider implements AWSCredentialsProvider, Closeable { + private static final AtomicInteger CLOSED = new AtomicInteger(); + private static final AtomicInteger CREATED = new AtomicInteger(); + + public TrackingCredentialsProvider() { + CREATED.incrementAndGet(); + } + + @Override + public AWSCredentials getCredentials() { + return new BasicAWSCredentials("minioadmin", "minioadmin"); + } + + @Override + public void refresh() {} + + @Override + public void close() { + CLOSED.incrementAndGet(); + } + } + + @Test + void shouldCloseCredentialsWhenEndpointInitializationFails() { + Map<String, Object> config = + S3FileConnectDryRunIT.sourceConfig("http://[", "dry-run-events", "/events"); + config.put("fs.s3a.aws.credentials.provider", TrackingCredentialsProvider.class.getName()); + int closedBefore = TrackingCredentialsProvider.CLOSED.get(); + IllegalArgumentException failure = + assertThrows( + IllegalArgumentException.class, + () -> S3FileConnectDryRunIT.validate(config)); + assertTrue(failure.getMessage().contains("endpoint")); + assertEquals(closedBefore + 1, TrackingCredentialsProvider.CLOSED.get()); + } + + @Test + void shouldRejectInvalidClientSettingsBeforeCreatingCredentials() { + Map<String, Object> config = + S3FileConnectDryRunIT.sourceConfig("http://localhost", "dry-run-events", "/events"); + config.put("fs.s3a.aws.credentials.provider", TrackingCredentialsProvider.class.getName()); + Map<String, String> properties = new HashMap<>(); + properties.put("fs.s3a.connection.maximum", "0"); + config.put("hadoop_s3_properties", properties); + int createdBefore = TrackingCredentialsProvider.CREATED.get(); + IllegalArgumentException failure = + assertThrows( + IllegalArgumentException.class, + () -> S3FileConnectDryRunIT.validate(config)); + assertTrue(failure.getMessage().contains("fs.s3a.connection.maximum")); + assertEquals(createdBefore, TrackingCredentialsProvider.CREATED.get()); + } + + private static final class MetadataServer implements AutoCloseable { + private final HttpServer server; + private final ExecutorService executor = Executors.newCachedThreadPool(); + private final ConcurrentLinkedQueue<String> requests = new ConcurrentLinkedQueue<>(); + private final CountDownLatch release = new CountDownLatch(1); + + private MetadataServer(String response) throws IOException { + server = HttpServer.create(new InetSocketAddress("127.0.0.1", 0), 0); + server.setExecutor(executor); + server.createContext("/", exchange -> reply(exchange, response)); + server.start(); + } + + private void reply(HttpExchange exchange, String response) throws IOException { + requests.add(exchange.getRequestMethod() + " " + exchange.getRequestURI()); + try { + exchange.getResponseHeaders().add("Connection", "close"); + if ("stalled".equals(response)) { + try { + release.await(); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + } + return; + } + if ("HEAD".equals(exchange.getRequestMethod())) { + int status = + "object".equals(response) + ? 200 + : "forbidden".equals(response) ? 403 : 404; + exchange.getResponseHeaders().add("ETag", "\"event-etag\""); + exchange.getResponseHeaders() + .add("Last-Modified", "Tue, 08 Sep 2026 00:00:00 GMT"); + exchange.sendResponseHeaders(status, -1); + return; + } + String xml = + "<ListBucketResult xmlns=\"http://s3.amazonaws.com/doc/2006-03-01/\">" + + "<Name>dry-run-events</Name><Prefix>events/</Prefix><MaxKeys>1</MaxKeys>" + + "<IsTruncated>true</IsTruncated><NextContinuationToken>remaining-events</NextContinuationToken>" + + "<Contents><Key>events/data.json</Key><LastModified>2026-09-08T00:00:00.000Z</LastModified>" + + "<ETag>event-etag</ETag><Size>14</Size><StorageClass>STANDARD</StorageClass></Contents></ListBucketResult>"; + byte[] body = xml.getBytes(StandardCharsets.UTF_8); + exchange.getResponseHeaders().add("Content-Type", "application/xml"); + exchange.sendResponseHeaders(200, body.length); + exchange.getResponseBody().write(body); + } finally { + exchange.close(); + } + } + + private Map<String, Object> config() { + return S3FileConnectDryRunIT.sourceConfig( + "http://127.0.0.1:" + server.getAddress().getPort(), + "dry-run-events", + "/events"); + } + + private List<String> requests() { + return new ArrayList<>(requests); + } + + @Override + public void close() { + release.countDown(); + server.stop(0); + executor.shutdownNow(); + } + } +}
