This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-25706-21a3215b66a8620cb932899de419e6b43d1848e6 in repository https://gitbox.apache.org/repos/asf/datafusion.git
commit e2ca7f38051744b2010cf09b80db8e9dfa4b5d38 Author: Andrew Duffy <[email protected]> AuthorDate: Thu Sep 24 14:54:42 2026 +0000 fix datafusion-cli CI: switch from minio to rustfs (#25706) ## Which issue does this PR close? Closes https://github.com/apache/datafusion/issues/25705 ## Rationale for this change This morning it seems like Minio pulled its images from quay.io, which has caused recent CI runs to start failing, e.g. https://github.com/apache/datafusion/actions/runs/36006680489/job/107657362588 RustFS is a pretty comparable (and more permissively licensed) alternative that provides a suitable S3 interface for integration tests. This was substantially authored by Codex, but reviewed by me and tested locally. ## What is the testing strategy for this PR? Existing integration tests ## Are there any user-facing changes? No Signed-off-by: Andrew Duffy <[email protected]> --- .github/workflows/rust.yml | 17 ++- datafusion-cli/CONTRIBUTING.md | 6 +- datafusion-cli/Cargo.toml | 2 +- datafusion-cli/tests/cli_integration.rs | 246 ++++++++++++++++++-------------- 4 files changed, 149 insertions(+), 122 deletions(-) diff --git a/.github/workflows/rust.yml b/.github/workflows/rust.yml index b9b8ed22cc..b6c9114942 100644 --- a/.github/workflows/rust.yml +++ b/.github/workflows/rust.yml @@ -480,22 +480,21 @@ jobs: with: save-if: false # set in linux-test shared-key: "amd-ci" - # The storage integration tests start MinIO containers. Pulling the image + # The storage integration tests start RustFS containers. Pulling the image # once up front keeps the pull off the critical path of the tests, which # would otherwise pull it several times concurrently and occasionally fail # with transient Docker transport errors. The tests retry the pull # themselves, so a failure here is only a warning. # - # MINIO_IMAGE must match what the tests start: the tag comes from the - # `minio` module of `testcontainers-modules`, the registry from - # MINIO_IMAGE_NAME in `datafusion-cli/tests/cli_integration.rs`. The - # `minio_image_matches_ci_prepull` test fails if either drifts. - - name: Pre-pull MinIO image + # RUSTFS_IMAGE must match RUSTFS_IMAGE_NAME and RUSTFS_IMAGE_TAG in + # `datafusion-cli/tests/cli_integration.rs`. The + # `rustfs_image_matches_ci_prepull` test fails if they differ. + - name: Pre-pull RustFS image env: - MINIO_IMAGE: quay.io/minio/minio:RELEASE.2025-02-28T09-55-16Z + RUSTFS_IMAGE: docker.io/rustfs/rustfs:1.0.0 run: | - ci/scripts/retry timeout 120 docker pull "$MINIO_IMAGE" \ - || echo "::warning::Could not pre-pull $MINIO_IMAGE, the tests will pull it themselves" + ci/scripts/retry timeout 120 docker pull "$RUSTFS_IMAGE" \ + || echo "::warning::Could not pre-pull $RUSTFS_IMAGE, the tests will pull it themselves" - name: Run tests (excluding doctests) env: RUST_BACKTRACE: 1 diff --git a/datafusion-cli/CONTRIBUTING.md b/datafusion-cli/CONTRIBUTING.md index 8be656ec4e..fa615028a4 100644 --- a/datafusion-cli/CONTRIBUTING.md +++ b/datafusion-cli/CONTRIBUTING.md @@ -35,7 +35,7 @@ cargo test --all-targets ## Running Storage Integration Tests -By default, storage integration tests are not run. These tests use the `testcontainers` crate to start up a local MinIO server using Docker on port 9000. +By default, storage integration tests are not run. These tests use the `testcontainers` crate to start up a local RustFS server using Docker with a dynamically assigned host port. To run them you will need to set `TEST_STORAGE_INTEGRATION`: @@ -47,8 +47,8 @@ For some of the tests, [snapshots](https://datafusion.apache.org/contributor-gui ### AWS -S3 integration is tested against [Minio](https://github.com/minio/minio) with [TestContainers](https://github.com/testcontainers/testcontainers-rs) -This requires Docker to be running on your machine and port 9000 to be free. +S3 integration is tested against [RustFS](https://github.com/rustfs/rustfs) with [TestContainers](https://github.com/testcontainers/testcontainers-rs) +This requires Docker to be running on your machine. If you see an error mentioning "failed to load IMDS session token" such as diff --git a/datafusion-cli/Cargo.toml b/datafusion-cli/Cargo.toml index aa508123d6..51a6109a01 100644 --- a/datafusion-cli/Cargo.toml +++ b/datafusion-cli/Cargo.toml @@ -78,7 +78,7 @@ ctor = { workspace = true } insta = { workspace = true } insta-cmd = "0.7.0" rstest = { workspace = true } -testcontainers-modules = { workspace = true, features = ["minio"] } +testcontainers-modules = { workspace = true } # Makes sure `test_display_pg_json` behaves in a consistent way regardless of # feature unification with dependencies serde_json = { workspace = true, features = ["preserve_order"] } diff --git a/datafusion-cli/tests/cli_integration.rs b/datafusion-cli/tests/cli_integration.rs index a9209bc598..e0fedf0d5b 100644 --- a/datafusion-cli/tests/cli_integration.rs +++ b/datafusion-cli/tests/cli_integration.rs @@ -20,17 +20,22 @@ use std::process::Command; use rstest::rstest; use async_trait::async_trait; +use futures::TryStreamExt; use insta::internals::SettingsBindDropGuard; use insta::{Settings, glob}; use insta_cmd::{assert_cmd_snapshot, get_cargo_bin}; +use object_store::{ + ObjectStore, ObjectStoreExt, aws::AmazonS3Builder, local::LocalFileSystem, +}; use std::path::PathBuf; use std::time::Duration; use std::{env, fs}; -use testcontainers_modules::minio; -use testcontainers_modules::testcontainers::core::{CmdWaitFor, ExecCommand, Mount}; +use testcontainers_modules::testcontainers::core::{ + CmdWaitFor, ExecCommand, IntoContainerPort, +}; use testcontainers_modules::testcontainers::runners::AsyncRunner; use testcontainers_modules::testcontainers::{ - ContainerAsync, Image, ImageExt, TestcontainersError, + ContainerAsync, GenericImage, ImageExt, TestcontainersError, }; fn cli() -> Command { @@ -46,43 +51,39 @@ fn make_settings() -> Settings { settings } -const MINIO_ROOT_USER: &str = "TEST-DataFusionLogin"; -const MINIO_ROOT_PASSWORD: &str = "TEST-DataFusionPassword"; +const RUSTFS_ACCESS_KEY: &str = "TEST-DataFusionLogin"; +const RUSTFS_SECRET_KEY: &str = "TEST-DataFusionPassword"; -/// Registry override for the image pinned by `testcontainers-modules`. -/// -/// MinIO withdrew `minio/minio` from Docker Hub on 2026-09-11. quay.io still -/// serves the same tag, so only the registry changes here. An unblock, not a -/// fix: see <https://github.com/apache/datafusion/issues/25215>. -const MINIO_IMAGE_NAME: &str = "quay.io/minio/minio"; +/// Pinned RustFS image, also pulled by the CLI CI job. +const RUSTFS_IMAGE_NAME: &str = "docker.io/rustfs/rustfs"; +const RUSTFS_IMAGE_TAG: &str = "1.0.0"; -/// How many times to try bringing up the MinIO container before failing. +/// How many times to try bringing up the RustFS container before failing. /// -/// Both the image pull and the `mc` calls that provision the bucket fail -/// intermittently on CI with transient errors such as -/// `bytes remaining on stream`. Retrying is much cheaper than a flaky run. -const MINIO_SETUP_ATTEMPTS: u32 = 3; +/// Image pulls and fixture uploads can fail with transient errors such as +/// `bytes remaining on stream`. Retry these failures before failing the test. +const RUSTFS_SETUP_ATTEMPTS: u32 = 3; -/// Delay before the first retry of the MinIO setup, doubled on each attempt. -const MINIO_SETUP_RETRY_DELAY: Duration = Duration::from_secs(5); +/// Delay before the first retry of the RustFS setup, doubled on each attempt. +const RUSTFS_SETUP_RETRY_DELAY: Duration = Duration::from_secs(5); -/// Time budget for a single MinIO setup attempt. A stalled image pull or `mc` +/// Time budget for a single RustFS setup attempt. A stalled image pull or `curl` /// invocation is retried instead of hanging the whole test run. -const MINIO_SETUP_TIMEOUT: Duration = Duration::from_mins(3); +const RUSTFS_SETUP_TIMEOUT: Duration = Duration::from_mins(3); -/// Starts a MinIO container preloaded with the test data, retrying transient +/// Starts a RustFS container preloaded with the test data, retrying transient /// Docker failures. /// /// Returns `None` when the test should be skipped, that is when /// `TEST_STORAGE_INTEGRATION` is unset or the registry is rate limiting the /// image pull. Panics if the container cannot be started for any other reason. -async fn start_minio_or_skip() -> Option<ContainerAsync<minio::MinIO>> { +async fn start_rustfs_or_skip() -> Option<ContainerAsync<GenericImage>> { if env::var("TEST_STORAGE_INTEGRATION").is_err() { eprintln!("Skipping external storages integration tests"); return None; } - match setup_minio_container().await { + match setup_rustfs_container().await { Ok(container) => Some(container), Err(e) if is_docker_pull_rate_limit(&e) => { eprintln!("Skipping test: Docker pull rate limit reached: {e}"); @@ -105,30 +106,30 @@ fn is_retryable(error: &str) -> bool { && !error.contains("failed to initialize a docker client") } -async fn setup_minio_container() -> Result<ContainerAsync<minio::MinIO>, String> { - let mut delay = MINIO_SETUP_RETRY_DELAY; - let mut last_error = String::from("MinIO container setup was not attempted at all"); +async fn setup_rustfs_container() -> Result<ContainerAsync<GenericImage>, String> { + let mut delay = RUSTFS_SETUP_RETRY_DELAY; + let mut last_error = String::from("RustFS container setup was not attempted at all"); - for attempt in 1..=MINIO_SETUP_ATTEMPTS { + for attempt in 1..=RUSTFS_SETUP_ATTEMPTS { last_error = match tokio::time::timeout( - MINIO_SETUP_TIMEOUT, - try_setup_minio_container(), + RUSTFS_SETUP_TIMEOUT, + try_setup_rustfs_container(), ) .await { Ok(Ok(container)) => return Ok(container), Ok(Err(e)) => e, Err(_) => format!( - "Timed out after {MINIO_SETUP_TIMEOUT:?} while starting the MinIO container" + "Timed out after {RUSTFS_SETUP_TIMEOUT:?} while starting the RustFS container" ), }; - if attempt == MINIO_SETUP_ATTEMPTS || !is_retryable(&last_error) { + if attempt == RUSTFS_SETUP_ATTEMPTS || !is_retryable(&last_error) { break; } eprintln!( - "MinIO container setup failed (attempt {attempt}/{MINIO_SETUP_ATTEMPTS}), \ + "RustFS container setup failed (attempt {attempt}/{RUSTFS_SETUP_ATTEMPTS}), \ retrying in {delay:?}: {last_error}" ); tokio::time::sleep(delay).await; @@ -138,69 +139,76 @@ async fn setup_minio_container() -> Result<ContainerAsync<minio::MinIO>, String> Err(last_error) } -/// A single attempt at starting and provisioning a MinIO container. +/// A single attempt at starting and provisioning a RustFS container. /// /// The container is removed again if provisioning fails, so that the next /// attempt starts from a clean state. -async fn try_setup_minio_container() -> Result<ContainerAsync<minio::MinIO>, String> { - let container = start_minio_container().await?; +async fn try_setup_rustfs_container() -> Result<ContainerAsync<GenericImage>, String> { + let container = start_rustfs_container().await?; - match provision_minio_container(&container).await { + match provision_rustfs_container(&container).await { Ok(()) => Ok(container), Err(e) => { if let Err(rm_error) = container.rm().await { - eprintln!("Failed to remove the MinIO container: {rm_error}"); + eprintln!("Failed to remove the RustFS container: {rm_error}"); } Err(e) } } } -async fn start_minio_container() -> Result<ContainerAsync<minio::MinIO>, String> { - let data_path = - PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../datafusion/core/tests/data"); - - let absolute_data_path = data_path - .canonicalize() - .expect("Failed to get absolute path for test data"); - - minio::MinIO::default() - .with_name(MINIO_IMAGE_NAME) - .with_env_var("MINIO_ROOT_USER", MINIO_ROOT_USER) - .with_env_var("MINIO_ROOT_PASSWORD", MINIO_ROOT_PASSWORD) - .with_mount(Mount::bind_mount( - absolute_data_path.to_str().unwrap(), - "/source", - )) +async fn start_rustfs_container() -> Result<ContainerAsync<GenericImage>, String> { + GenericImage::new(RUSTFS_IMAGE_NAME, RUSTFS_IMAGE_TAG) + .with_exposed_port(9000.tcp()) + .with_env_var("RUSTFS_ACCESS_KEY", RUSTFS_ACCESS_KEY) + .with_env_var("RUSTFS_SECRET_KEY", RUSTFS_SECRET_KEY) + .with_env_var("RUSTFS_CONSOLE_ENABLE", "false") .start() .await .map_err(|e| match e { TestcontainersError::Client(e) => format!( - "Failed to start MinIO container. Ensure Docker is running and accessible: {e}" + "Failed to start RustFS container. Ensure Docker is running and accessible: {e}" ), - e => format!("Failed to start MinIO container: {e}"), + e => format!("Failed to start RustFS container: {e}"), }) } -/// Waits for MinIO to be healthy and uploads the test files. +/// Waits for RustFS to be healthy and uploads the test files. /// -/// This is done via the `mc` CLI shipped in the image to avoid an s3 dependency. -async fn provision_minio_container( - container: &ContainerAsync<minio::MinIO>, +/// Use the image's `curl` to check readiness and create the bucket, then upload +/// the fixtures through `object_store`. +async fn provision_rustfs_container( + container: &ContainerAsync<GenericImage>, ) -> Result<(), String> { + let credentials = format!("{RUSTFS_ACCESS_KEY}:{RUSTFS_SECRET_KEY}"); let commands = [ - ExecCommand::new(["/usr/bin/mc", "ready", "local"]), ExecCommand::new([ - "/usr/bin/mc", - "alias", - "set", - "localminio", - "http://localhost:9000", - MINIO_ROOT_USER, - MINIO_ROOT_PASSWORD, + "curl", + "--fail", + "--silent", + "--show-error", + "--retry", + "60", + "--retry-delay", + "1", + "--retry-all-errors", + "--max-time", + "5", + "http://localhost:9000/health/ready", + ]), + ExecCommand::new([ + "curl", + "--fail", + "--silent", + "--show-error", + "--aws-sigv4", + "aws:amz:us-east-1:s3", + "--user", + &credentials, + "--request", + "PUT", + "http://localhost:9000/data", ]), - ExecCommand::new(["/usr/bin/mc", "mb", "localminio/data"]), - ExecCommand::new(["/usr/bin/mc", "cp", "-r", "/source/", "localminio/data/"]), ]; for command in commands { @@ -223,14 +231,45 @@ async fn provision_minio_container( } } + let port = container + .get_host_port_ipv4(9000) + .await + .map_err(|e| e.to_string())?; + let store = AmazonS3Builder::new() + .with_bucket_name("data") + .with_region("us-east-1") + .with_access_key_id(RUSTFS_ACCESS_KEY) + .with_secret_access_key(RUSTFS_SECRET_KEY) + .with_endpoint(format!("http://localhost:{port}")) + .with_allow_http(true) + .build() + .map_err(|e| e.to_string())?; + let data_path = + PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../datafusion/core/tests/data"); + let source = + LocalFileSystem::new_with_prefix(data_path).map_err(|e| e.to_string())?; + let mut files = source.list(None); + while let Some(file) = files.try_next().await.map_err(|e| e.to_string())? { + let data = source + .get(&file.location) + .await + .map_err(|e| e.to_string())? + .bytes() + .await + .map_err(|e| e.to_string())?; + store + .put(&file.location, data.into()) + .await + .map_err(|e| e.to_string())?; + } + Ok(()) } -/// CI pre-pulls the MinIO image so that the storage integration tests do not -/// have to pull it themselves. Guard against that pre-pull going stale when -/// `testcontainers-modules` bumps the image it uses. +/// CI pre-pulls the RustFS image so that the storage integration tests do not +/// have to pull it themselves. Keep the CI image and the test image in sync. #[test] -fn minio_image_matches_ci_prepull() { +fn rustfs_image_matches_ci_prepull() { let workflow = PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../.github/workflows/rust.yml"); @@ -239,25 +278,12 @@ fn minio_image_matches_ci_prepull() { return; }; - let image = minio::MinIO::default(); + let image_ref = format!("{RUSTFS_IMAGE_NAME}:{RUSTFS_IMAGE_TAG}"); - // The override only redirects the registry, so it would silently stop - // tracking upstream if the crate ever pinned a different image. assert!( - MINIO_IMAGE_NAME.ends_with(&format!("/{}", image.name())), - "`testcontainers-modules` now uses `{}`, which MINIO_IMAGE_NAME \ - (`{MINIO_IMAGE_NAME}`) no longer mirrors.", - image.name() - ); - - // Match the assignment: `minio/minio:<tag>` is a substring of the quay - // reference and would pass either way. - let image_ref = format!("{MINIO_IMAGE_NAME}:{}", image.tag()); - - assert!( - contents.contains(&format!("MINIO_IMAGE: {image_ref}")), - "{} does not pre-pull `{image_ref}`. Update MINIO_IMAGE in the \ - `Pre-pull MinIO image` step to match the image used by the tests.", + contents.contains(&format!("RUSTFS_IMAGE: {image_ref}")), + "{} does not pre-pull `{image_ref}`. Update RUSTFS_IMAGE in the \ + `Pre-pull RustFS image` step to match the image used by the tests.", workflow.display() ); } @@ -695,7 +721,7 @@ fn test_cli_wide_result_set_no_crash() { #[tokio::test] async fn test_cli() { - let Some(container) = start_minio_or_skip().await else { + let Some(container) = start_rustfs_or_skip().await else { return; }; @@ -709,8 +735,8 @@ async fn test_cli() { assert_cmd_snapshot!( cli() .env_clear() - .env("AWS_ACCESS_KEY_ID", MINIO_ROOT_USER) - .env("AWS_SECRET_ACCESS_KEY", MINIO_ROOT_PASSWORD) + .env("AWS_ACCESS_KEY_ID", RUSTFS_ACCESS_KEY) + .env("AWS_SECRET_ACCESS_KEY", RUSTFS_SECRET_KEY) .env("AWS_ENDPOINT", format!("http://localhost:{port}")) .env("AWS_ALLOW_HTTP", "true") .pass_stdin(input) @@ -722,7 +748,7 @@ async fn test_cli() { async fn test_aws_options() { // Separate test is needed to pass aws as options in sql and not via env - let Some(container) = start_minio_or_skip().await else { + let Some(container) = start_rustfs_or_skip().await else { return; }; @@ -736,8 +762,8 @@ async fn test_aws_options() { STORED AS CSV LOCATION 's3://data/cars.csv' OPTIONS( - 'aws.access_key_id' '{MINIO_ROOT_USER}', - 'aws.secret_access_key' '{MINIO_ROOT_PASSWORD}', + 'aws.access_key_id' '{RUSTFS_ACCESS_KEY}', + 'aws.secret_access_key' '{RUSTFS_SECRET_KEY}', 'aws.endpoint' 'http://localhost:{port}', 'aws.allow_http' 'true' ); @@ -812,7 +838,7 @@ fn test_backtrace_output(#[case] query: &str) { #[tokio::test] async fn test_s3_url_fallback() { - let Some(container) = start_minio_or_skip().await else { + let Some(container) = start_rustfs_or_skip().await else { return; }; @@ -833,13 +859,13 @@ OPTIONS ( SELECT * FROM partitioned_data ORDER BY column_1, column_2 LIMIT 5; "#; - assert_cmd_snapshot!(cli().with_minio(&container).await.pass_stdin(input)); + assert_cmd_snapshot!(cli().with_rustfs(&container).await.pass_stdin(input)); } /// Validate object store profiling output #[tokio::test] async fn test_object_store_profiling() { - let Some(container) = start_minio_or_skip().await else { + let Some(container) = start_rustfs_or_skip().await else { return; }; @@ -885,27 +911,29 @@ SELECT * from CARS LIMIT 1; SELECT * from CARS LIMIT 1; "#; - assert_cmd_snapshot!(cli().with_minio(&container).await.pass_stdin(input)); + assert_cmd_snapshot!(cli().with_rustfs(&container).await.pass_stdin(input)); } -/// Extension trait to Add the minio connection information to a Command +/// Add the RustFS connection information to a Command. #[async_trait] -trait MinioCommandExt { - async fn with_minio(&mut self, container: &ContainerAsync<minio::MinIO>) - -> &mut Self; +trait RustfsCommandExt { + async fn with_rustfs( + &mut self, + container: &ContainerAsync<GenericImage>, + ) -> &mut Self; } #[async_trait] -impl MinioCommandExt for Command { - async fn with_minio( +impl RustfsCommandExt for Command { + async fn with_rustfs( &mut self, - container: &ContainerAsync<minio::MinIO>, + container: &ContainerAsync<GenericImage>, ) -> &mut Self { let port = container.get_host_port_ipv4(9000).await.unwrap(); self.env_clear() - .env("AWS_ACCESS_KEY_ID", MINIO_ROOT_USER) - .env("AWS_SECRET_ACCESS_KEY", MINIO_ROOT_PASSWORD) + .env("AWS_ACCESS_KEY_ID", RUSTFS_ACCESS_KEY) + .env("AWS_SECRET_ACCESS_KEY", RUSTFS_SECRET_KEY) .env("AWS_ENDPOINT", format!("http://localhost:{port}")) .env("AWS_ALLOW_HTTP", "true") } --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
