This is an automated email from the ASF dual-hosted git repository.

hubcio pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git


The following commit(s) were added to refs/heads/master by this push:
     new 1951e654c chore(integration): replace MinIO  with Floci (#4290)
1951e654c is described below

commit 1951e654c11021905dc46b0ce3696899761de696
Author: GRPM <[email protected]>
AuthorDate: Sat Sep 26 13:01:23 2026 +0400

    chore(integration): replace MinIO  with Floci (#4290)
---
 core/connectors/sinks/redshift_sink/README.md      |   2 +-
 .../tests/connectors/fixtures/delta/fixture.rs     | 158 ++++------------
 .../integration/tests/connectors/fixtures/floci.rs | 159 ++++++++++++++++
 .../tests/connectors/fixtures/iceberg/container.rs | 166 ++++-------------
 core/integration/tests/connectors/fixtures/mod.rs  |   1 +
 .../connectors/fixtures/redshift/container.rs      |  73 +-------
 .../tests/connectors/fixtures/redshift/mod.rs      |   2 +-
 .../tests/connectors/fixtures/redshift/sink.rs     | 203 +++++----------------
 .../tests/connectors/fixtures/s3/fixture.rs        | 145 ++-------------
 core/integration/tests/connectors/s3/s3_sink.rs    |  33 +++-
 10 files changed, 324 insertions(+), 618 deletions(-)

diff --git a/core/connectors/sinks/redshift_sink/README.md 
b/core/connectors/sinks/redshift_sink/README.md
index 40869de3e..9a1f8ba52 100644
--- a/core/connectors/sinks/redshift_sink/README.md
+++ b/core/connectors/sinks/redshift_sink/README.md
@@ -193,7 +193,7 @@ COPY {staging_table} ({columns})
 FROM STDIN BINARY
 ```
 
-The `s3_path` is parsed and used to fetch the object from the MinIO instance 
backing the mock container, with access key and secret key supplied to the 
container via environment variables rather than an IAM role. Instead of 
Redshift pulling directly from S3, the connector reads the object itself and 
streams it into the mock over `COPY ... FROM STDIN BINARY`, so the 
`CREDENTIALS`, `FORMAT AS PARQUET`, and `REGION` clauses have no equivalent 
here.
+In integration tests, the mock uses the `s3_path` to fetch objects from Floci 
with test credentials. Instead of Redshift pulling directly from S3, the mock 
reads the object and streams it into PostgreSQL over `COPY ... FROM STDIN 
BINARY`, so the `CREDENTIALS`, `FORMAT AS PARQUET`, and `REGION` clauses have 
no equivalent here.
 
 ## 5. Staging → target insert (idempotent upsert)
 
diff --git a/core/integration/tests/connectors/fixtures/delta/fixture.rs 
b/core/integration/tests/connectors/fixtures/delta/fixture.rs
index 90087b455..9fa467091 100644
--- a/core/integration/tests/connectors/fixtures/delta/fixture.rs
+++ b/core/integration/tests/connectors/fixtures/delta/fixture.rs
@@ -15,19 +15,19 @@
 // specific language governing permissions and limitations
 // under the License.
 
-use crate::connectors::fixtures;
+use std::{collections::HashMap, path::PathBuf};
+
 use async_trait::async_trait;
 use deltalake::kernel::{DataType, PrimitiveType, StructField};
 use deltalake::operations::create::CreateBuilder;
 use integration::harness::{TestBinaryError, TestFixture};
-use std::collections::HashMap;
-use std::path::PathBuf;
 use tempfile::TempDir;
-use testcontainers_modules::testcontainers::core::{IntoContainerPort, WaitFor};
-use testcontainers_modules::testcontainers::runners::AsyncRunner;
-use testcontainers_modules::testcontainers::{ContainerAsync, GenericImage, 
ImageExt};
 use tracing::info;
-use uuid::Uuid;
+
+use crate::connectors::fixtures::{
+    self,
+    floci::{self, ACCESS_KEY, FlociContainer, REGION, SECRET_KEY},
+};
 
 const ENV_SINK_TABLE_URI: &str = 
"IGGY_CONNECTORS_SINK_DELTA_PLUGIN_CONFIG_TABLE_URI";
 const ENV_SINK_PATH: &str = "IGGY_CONNECTORS_SINK_DELTA_PATH";
@@ -43,13 +43,7 @@ const ENV_SINK_AWS_S3_ENDPOINT_URL: &str =
 const ENV_SINK_AWS_S3_ALLOW_HTTP: &str =
     "IGGY_CONNECTORS_SINK_DELTA_PLUGIN_CONFIG_AWS_S3_ALLOW_HTTP";
 
-const MINIO_IMAGE: &str = "quay.io/minio/minio";
-const MINIO_TAG: &str = "RELEASE.2025-09-07T16-13-09Z";
-const MINIO_PORT: u16 = 9000;
-const MINIO_CONSOLE_PORT: u16 = 9001;
-const MINIO_ACCESS_KEY: &str = "admin";
-const MINIO_SECRET_KEY: &str = "password";
-const MINIO_BUCKET: &str = "delta-warehouse";
+const TEST_BUCKET: &str = "delta-warehouse";
 
 pub struct DeltaFixture {
     _temp_dir: TempDir,
@@ -202,100 +196,19 @@ impl TestFixture for DeltaFixture {
 
 pub struct DeltaS3Fixture {
     #[allow(dead_code)]
-    minio: ContainerAsync<GenericImage>,
-    minio_endpoint: String,
+    floci: FlociContainer,
+    floci_endpoint: String,
 }
 
 impl DeltaS3Fixture {
-    async fn start_minio(
-        network: &str,
-        container_name: &str,
-    ) -> Result<(ContainerAsync<GenericImage>, String), TestBinaryError> {
-        let container = GenericImage::new(MINIO_IMAGE, MINIO_TAG)
-            .with_exposed_port(MINIO_PORT.tcp())
-            .with_exposed_port(MINIO_CONSOLE_PORT.tcp())
-            .with_wait_for(WaitFor::message_on_stderr("API:"))
-            .with_network(network)
-            .with_container_name(container_name)
-            .with_env_var("MINIO_ROOT_USER", MINIO_ACCESS_KEY)
-            .with_env_var("MINIO_ROOT_PASSWORD", MINIO_SECRET_KEY)
-            .with_cmd(vec!["server", "/data", "--console-address", ":9001"])
-            .with_mapped_port(0, MINIO_PORT.tcp())
-            .with_mapped_port(0, MINIO_CONSOLE_PORT.tcp())
-            .start()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "DeltaS3Fixture".to_string(),
-                message: format!("Failed to start MinIO container: {error}"),
-            })?;
-
-        info!("Started MinIO container for Delta S3 tests");
-
-        let mapped_port = container
-            .ports()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "DeltaS3Fixture".to_string(),
-                message: format!("Failed to get ports: {error}"),
-            })?
-            .map_to_host_port_ipv4(MINIO_PORT)
-            .ok_or_else(|| TestBinaryError::FixtureSetup {
-                fixture_type: "DeltaS3Fixture".to_string(),
-                message: "No mapping for MinIO port".to_string(),
-            })?;
-
-        let endpoint = format!("http://localhost:{mapped_port}";);
-        info!("MinIO container available at {endpoint}");
-
-        Ok((container, endpoint))
-    }
-
-    async fn create_bucket(minio_endpoint: &str) -> Result<(), 
TestBinaryError> {
-        use tokio::process::Command;
-
-        let host = minio_endpoint.trim_start_matches("http://";);
-        let mc_host = format!("http://{}:{}@{}";, MINIO_ACCESS_KEY, 
MINIO_SECRET_KEY, host);
-
-        let output = Command::new("docker")
-            .args([
-                "run",
-                "--rm",
-                "--network=host",
-                "-e",
-                &format!("MC_HOST_minio={}", mc_host),
-                "quay.io/minio/mc",
-                "mb",
-                "--ignore-existing",
-                &format!("minio/{}", MINIO_BUCKET),
-            ])
-            .output()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "DeltaS3Fixture".to_string(),
-                message: format!("Failed to run mc command: {error}"),
-            })?;
-
-        if !output.status.success() {
-            let stderr = String::from_utf8_lossy(&output.stderr);
-            let stdout = String::from_utf8_lossy(&output.stdout);
-            return Err(TestBinaryError::FixtureSetup {
-                fixture_type: "DeltaS3Fixture".to_string(),
-                message: format!("Failed to create bucket: stderr={stderr}, 
stdout={stdout}"),
-            });
-        }
-
-        info!("Created MinIO bucket: {MINIO_BUCKET}");
-        Ok(())
-    }
-
-    async fn create_table(minio_endpoint: &str) -> Result<(), TestBinaryError> 
{
-        let table_uri = format!("s3://{MINIO_BUCKET}/delta_table");
+    async fn create_table(floci_endpoint: &str) -> Result<(), TestBinaryError> 
{
+        let table_uri = format!("s3://{TEST_BUCKET}/delta_table");
         let columns = table_columns();
         let storage_options = HashMap::from([
-            ("AWS_ACCESS_KEY_ID".into(), MINIO_ACCESS_KEY.into()),
-            ("AWS_SECRET_ACCESS_KEY".into(), MINIO_SECRET_KEY.into()),
-            ("AWS_REGION".into(), "us-east-1".into()),
-            ("AWS_ENDPOINT_URL".into(), minio_endpoint.into()),
+            ("AWS_ACCESS_KEY_ID".into(), ACCESS_KEY.into()),
+            ("AWS_SECRET_ACCESS_KEY".into(), SECRET_KEY.into()),
+            ("AWS_REGION".into(), REGION.into()),
+            ("AWS_ENDPOINT_URL".into(), floci_endpoint.into()),
             ("AWS_ALLOW_HTTP".into(), "true".into()),
             ("AWS_S3_ALLOW_HTTP".into(), "true".into()),
         ]);
@@ -306,7 +219,7 @@ impl DeltaS3Fixture {
             .await
             .map_err(|error| TestBinaryError::FixtureSetup {
                 fixture_type: "DeltaS3Fixture".to_string(),
-                message: format!("Failed to create Delta table in MinIO: 
{error}"),
+                message: format!("Failed to create Delta table in Floci: 
{error}"),
             })?;
         Ok(())
     }
@@ -318,16 +231,16 @@ impl DeltaS3Fixture {
         interval_ms: u64,
     ) -> Result<usize, TestBinaryError> {
         let table_uri =
-            
url::Url::parse(&format!("s3://{MINIO_BUCKET}/delta_table")).map_err(|e| {
+            
url::Url::parse(&format!("s3://{TEST_BUCKET}/delta_table")).map_err(|e| {
                 TestBinaryError::InvalidState {
                     message: format!("Failed to parse table URI: {e}"),
                 }
             })?;
         let storage_options = HashMap::from([
-            ("AWS_ACCESS_KEY_ID".into(), MINIO_ACCESS_KEY.into()),
-            ("AWS_SECRET_ACCESS_KEY".into(), MINIO_SECRET_KEY.into()),
-            ("AWS_REGION".into(), "us-east-1".into()),
-            ("AWS_ENDPOINT_URL".into(), self.minio_endpoint.clone()),
+            ("AWS_ACCESS_KEY_ID".into(), ACCESS_KEY.into()),
+            ("AWS_SECRET_ACCESS_KEY".into(), SECRET_KEY.into()),
+            ("AWS_REGION".into(), REGION.into()),
+            ("AWS_ENDPOINT_URL".into(), self.floci_endpoint.clone()),
             ("AWS_ALLOW_HTTP".into(), "true".into()),
             ("AWS_S3_ALLOW_HTTP".into(), "true".into()),
         ]);
@@ -345,24 +258,23 @@ impl DeltaS3Fixture {
 #[async_trait]
 impl TestFixture for DeltaS3Fixture {
     async fn setup() -> Result<Self, TestBinaryError> {
-        let id = Uuid::new_v4();
-        let network = format!("iggy-delta-s3-{id}");
-        let minio_name = fixtures::unique_container_name("minio-delta");
+        let floci_name = fixtures::unique_container_name("floci-delta");
 
-        let (minio, minio_endpoint) = Self::start_minio(&network, 
&minio_name).await?;
-        Self::create_bucket(&minio_endpoint).await?;
-        Self::create_table(&minio_endpoint).await?;
+        let floci = FlociContainer::start(None, &floci_name).await?;
+        let floci_endpoint = floci.endpoint.clone();
+        floci::create_bucket(&floci_endpoint, TEST_BUCKET).await?;
+        Self::create_table(&floci_endpoint).await?;
 
-        info!("Delta S3 fixture ready with MinIO at {minio_endpoint}");
+        info!("Delta S3 fixture ready with Floci at {floci_endpoint}");
 
         Ok(Self {
-            minio,
-            minio_endpoint,
+            floci,
+            floci_endpoint,
         })
     }
 
     fn connectors_runtime_envs(&self) -> HashMap<String, String> {
-        let table_uri = format!("s3://{MINIO_BUCKET}/delta_table");
+        let table_uri = format!("s3://{TEST_BUCKET}/delta_table");
 
         let mut envs = HashMap::new();
         envs.insert(
@@ -373,16 +285,16 @@ impl TestFixture for DeltaS3Fixture {
         envs.insert(ENV_SINK_STORAGE_BACKEND_TYPE.to_string(), 
"s3".to_string());
         envs.insert(
             ENV_SINK_AWS_S3_ACCESS_KEY.to_string(),
-            MINIO_ACCESS_KEY.to_string(),
+            ACCESS_KEY.to_string(),
         );
         envs.insert(
             ENV_SINK_AWS_S3_SECRET_KEY.to_string(),
-            MINIO_SECRET_KEY.to_string(),
+            SECRET_KEY.to_string(),
         );
-        envs.insert(ENV_SINK_AWS_S3_REGION.to_string(), 
"us-east-1".to_string());
+        envs.insert(ENV_SINK_AWS_S3_REGION.to_string(), REGION.to_string());
         envs.insert(
             ENV_SINK_AWS_S3_ENDPOINT_URL.to_string(),
-            self.minio_endpoint.clone(),
+            self.floci_endpoint.clone(),
         );
         envs.insert(ENV_SINK_AWS_S3_ALLOW_HTTP.to_string(), 
"true".to_string());
         envs
diff --git a/core/integration/tests/connectors/fixtures/floci.rs 
b/core/integration/tests/connectors/fixtures/floci.rs
new file mode 100644
index 000000000..95baf1334
--- /dev/null
+++ b/core/integration/tests/connectors/fixtures/floci.rs
@@ -0,0 +1,159 @@
+// 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.
+
+use std::time::Duration;
+
+use integration::harness::TestBinaryError;
+use s3::{
+    Bucket, BucketConfiguration, Region,
+    command::Command,
+    creds::Credentials,
+    request::{Request, tokio_backend::ReqwestRequest},
+};
+use testcontainers_modules::testcontainers::{
+    ContainerAsync, GenericImage, ImageExt,
+    core::{IntoContainerPort, WaitFor, wait::HttpWaitStrategy},
+    runners::AsyncRunner,
+};
+
+pub const ACCESS_KEY: &str = "test";
+pub const SECRET_KEY: &str = "test";
+pub const REGION: &str = "us-east-1";
+
+const IMAGE: &str = "floci/floci";
+const TAG: &str = "2.1.0";
+const PORT: u16 = 4566;
+const BUCKET_CREATE_ATTEMPTS: usize = 30;
+const BUCKET_CREATE_RETRY_DELAY: Duration = Duration::from_secs(1);
+
+pub struct FlociContainer {
+    #[allow(dead_code)]
+    container: ContainerAsync<GenericImage>,
+    pub endpoint: String,
+    pub internal_endpoint: String,
+}
+
+impl FlociContainer {
+    pub async fn start(
+        network: Option<&str>,
+        container_name: &str,
+    ) -> Result<Self, TestBinaryError> {
+        let request = GenericImage::new(IMAGE, TAG)
+            .with_exposed_port(PORT.tcp())
+            .with_wait_for(WaitFor::http(
+                HttpWaitStrategy::new("/_localstack/health")
+                    .with_port(PORT.tcp())
+                    .with_expected_status_code(200u16),
+            ))
+            .with_container_name(container_name)
+            .with_mapped_port(0, PORT.tcp());
+        let request = if let Some(network) = network {
+            request.with_network(network)
+        } else {
+            request
+        };
+        let container = request
+            .start()
+            .await
+            .map_err(|error| TestBinaryError::FixtureSetup {
+                fixture_type: "FlociContainer".to_string(),
+                message: format!("Failed to start container: {error}"),
+            })?;
+
+        let mapped_port = 
container.get_host_port_ipv4(PORT).await.map_err(|error| {
+            TestBinaryError::FixtureSetup {
+                fixture_type: "FlociContainer".to_string(),
+                message: format!("Failed to get port: {error}"),
+            }
+        })?;
+
+        let endpoint = format!("http://localhost:{mapped_port}";);
+        let internal_endpoint = format!("http://{container_name}:{PORT}";);
+        tracing::info!("Floci available at {endpoint}");
+
+        Ok(Self {
+            container,
+            endpoint,
+            internal_endpoint,
+        })
+    }
+}
+
+pub async fn create_bucket(
+    endpoint: &str,
+    bucket_name: &str,
+) -> Result<Box<Bucket>, TestBinaryError> {
+    let region = Region::Custom {
+        region: REGION.to_string(),
+        endpoint: endpoint.to_string(),
+    };
+    let credentials = Credentials::new(Some(ACCESS_KEY), Some(SECRET_KEY), 
None, None, None)
+        .map_err(|error| TestBinaryError::FixtureSetup {
+            fixture_type: "FlociContainer".to_string(),
+            message: format!("Failed to create credentials: {error}"),
+        })?;
+
+    // rust-s3's bucket creation helper adds a LocationConstraint for custom
+    // endpoints, but us-east-1 bucket creation requires an empty body.
+    let mut bucket = Bucket::new(bucket_name, region, 
credentials).map_err(|error| {
+        TestBinaryError::FixtureSetup {
+            fixture_type: "FlociContainer".to_string(),
+            message: format!("Failed to open bucket '{bucket_name}': {error}"),
+        }
+    })?;
+    bucket.set_path_style();
+    let mut last_status = 0;
+    let mut last_response = String::new();
+    for _ in 0..BUCKET_CREATE_ATTEMPTS {
+        let request = ReqwestRequest::new(
+            &bucket,
+            "",
+            Command::CreateBucket {
+                config: BucketConfiguration::default(),
+            },
+        )
+        .await
+        .map_err(|error| TestBinaryError::FixtureSetup {
+            fixture_type: "FlociContainer".to_string(),
+            message: format!("Failed to prepare bucket creation for 
'{bucket_name}': {error}"),
+        })?;
+        let response =
+            request
+                .response_data(false)
+                .await
+                .map_err(|error| TestBinaryError::FixtureSetup {
+                    fixture_type: "FlociContainer".to_string(),
+                    message: format!("Failed to create bucket '{bucket_name}': 
{error}"),
+                })?;
+        last_status = response.status_code();
+        if (200..300).contains(&last_status) || last_status == 409 {
+            return Ok(bucket);
+        }
+        last_response = 
String::from_utf8_lossy(response.as_slice()).into_owned();
+        if (400..500).contains(&last_status) {
+            break;
+        }
+        tokio::time::sleep(BUCKET_CREATE_RETRY_DELAY).await;
+    }
+
+    Err(TestBinaryError::FixtureSetup {
+        fixture_type: "FlociContainer".to_string(),
+        message: format!(
+            "Bucket '{bucket_name}' not creatable (last status: {last_status}, 
response: {last_response})"
+        ),
+    })
+}
diff --git a/core/integration/tests/connectors/fixtures/iceberg/container.rs 
b/core/integration/tests/connectors/fixtures/iceberg/container.rs
index b699ccb70..6e1c879bb 100644
--- a/core/integration/tests/connectors/fixtures/iceberg/container.rs
+++ b/core/integration/tests/connectors/fixtures/iceberg/container.rs
@@ -15,13 +15,13 @@
 // specific language governing permissions and limitations
 // under the License.
 
-use crate::connectors::fixtures;
+use std::collections::HashMap;
+
 use async_trait::async_trait;
 use integration::harness::{TestBinaryError, TestFixture};
 use reqwest_middleware::ClientWithMiddleware as HttpClient;
 use reqwest_retry::RetryTransientMiddleware;
 use reqwest_retry::policies::ExponentialBackoff;
-use std::collections::HashMap;
 use testcontainers_modules::testcontainers::core::wait::HttpWaitStrategy;
 use testcontainers_modules::testcontainers::core::{IntoContainerPort, WaitFor};
 use testcontainers_modules::testcontainers::runners::AsyncRunner;
@@ -29,17 +29,16 @@ use 
testcontainers_modules::testcontainers::{ContainerAsync, GenericImage, Image
 use tracing::info;
 use uuid::Uuid;
 
-const MINIO_IMAGE: &str = "quay.io/minio/minio";
-const MINIO_TAG: &str = "RELEASE.2025-09-07T16-13-09Z";
-const MINIO_PORT: u16 = 9000;
-const MINIO_CONSOLE_PORT: u16 = 9001;
+use crate::connectors::fixtures::{
+    self,
+    floci::{self, ACCESS_KEY, FlociContainer, REGION, SECRET_KEY},
+};
+
 const ICEBERG_REST_IMAGE: &str = "docker.io/apache/iceberg-rest-fixture";
 const ICEBERG_REST_TAG: &str = "latest";
 const ICEBERG_REST_PORT: u16 = 8181;
 
-pub const MINIO_ACCESS_KEY: &str = "admin";
-pub const MINIO_SECRET_KEY: &str = "password";
-pub const MINIO_BUCKET: &str = "warehouse";
+pub const TEST_BUCKET: &str = "warehouse";
 
 pub const ENV_SINK_URI: &str = 
"IGGY_CONNECTORS_SINK_ICEBERG_PLUGIN_CONFIG_URI";
 pub const ENV_SINK_WAREHOUSE: &str = 
"IGGY_CONNECTORS_SINK_ICEBERG_PLUGIN_CONFIG_WAREHOUSE";
@@ -54,64 +53,6 @@ pub const ENV_SINK_PATH: &str = 
"IGGY_CONNECTORS_SINK_ICEBERG_PATH";
 pub const ENV_AWS_ACCESS_KEY_ID: &str = "AWS_ACCESS_KEY_ID";
 pub const ENV_AWS_SECRET_ACCESS_KEY: &str = "AWS_SECRET_ACCESS_KEY";
 
-pub struct MinioContainer {
-    #[allow(dead_code)]
-    container: ContainerAsync<GenericImage>,
-    pub endpoint: String,
-    pub internal_endpoint: String,
-}
-
-impl MinioContainer {
-    pub async fn start(network: &str, container_name: &str) -> Result<Self, 
TestBinaryError> {
-        let container = GenericImage::new(MINIO_IMAGE, MINIO_TAG)
-            .with_exposed_port(MINIO_PORT.tcp())
-            .with_exposed_port(MINIO_CONSOLE_PORT.tcp())
-            .with_wait_for(WaitFor::http(
-                HttpWaitStrategy::new("/minio/health/live")
-                    .with_port(MINIO_PORT.tcp())
-                    .with_expected_status_code(200u16),
-            ))
-            .with_network(network)
-            .with_container_name(container_name)
-            .with_env_var("MINIO_ROOT_USER", MINIO_ACCESS_KEY)
-            .with_env_var("MINIO_ROOT_PASSWORD", MINIO_SECRET_KEY)
-            .with_cmd(vec!["server", "/data", "--console-address", ":9001"])
-            .with_mapped_port(0, MINIO_PORT.tcp())
-            .with_mapped_port(0, MINIO_CONSOLE_PORT.tcp())
-            .start()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "MinioContainer".to_string(),
-                message: format!("Failed to start container: {error}"),
-            })?;
-
-        info!("Started MinIO container");
-
-        let mapped_port = container
-            .ports()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "MinioContainer".to_string(),
-                message: format!("Failed to get ports: {error}"),
-            })?
-            .map_to_host_port_ipv4(MINIO_PORT)
-            .ok_or_else(|| TestBinaryError::FixtureSetup {
-                fixture_type: "MinioContainer".to_string(),
-                message: "No mapping for MinIO port".to_string(),
-            })?;
-
-        let endpoint = format!("http://localhost:{mapped_port}";);
-        let internal_endpoint = 
format!("http://{container_name}:{MINIO_PORT}";);
-        info!("MinIO container available at {endpoint} (internal: 
{internal_endpoint})");
-
-        Ok(Self {
-            container,
-            endpoint,
-            internal_endpoint,
-        })
-    }
-}
-
 pub struct IcebergRestContainer {
     #[allow(dead_code)]
     container: ContainerAsync<GenericImage>,
@@ -121,9 +62,9 @@ pub struct IcebergRestContainer {
 impl IcebergRestContainer {
     pub async fn start(
         network: &str,
-        minio_internal_endpoint: &str,
+        floci_internal_endpoint: &str,
     ) -> Result<Self, TestBinaryError> {
-        let warehouse_path = format!("s3://{MINIO_BUCKET}/");
+        let warehouse_path = format!("s3://{TEST_BUCKET}/");
 
         let container = GenericImage::new(ICEBERG_REST_IMAGE, ICEBERG_REST_TAG)
             .with_exposed_port(ICEBERG_REST_PORT.tcp())
@@ -141,13 +82,13 @@ impl IcebergRestContainer {
             .with_env_var("CATALOG_URI", "jdbc:sqlite:/tmp/iceberg_catalog.db")
             .with_env_var("CATALOG_WAREHOUSE", &warehouse_path)
             .with_env_var("CATALOG_IO__IMPL", 
"org.apache.iceberg.aws.s3.S3FileIO")
-            .with_env_var("CATALOG_S3_ENDPOINT", minio_internal_endpoint)
-            .with_env_var("CATALOG_S3_ACCESS__KEY__ID", MINIO_ACCESS_KEY)
-            .with_env_var("CATALOG_S3_SECRET__ACCESS__KEY", MINIO_SECRET_KEY)
+            .with_env_var("CATALOG_S3_ENDPOINT", floci_internal_endpoint)
+            .with_env_var("CATALOG_S3_ACCESS__KEY__ID", ACCESS_KEY)
+            .with_env_var("CATALOG_S3_SECRET__ACCESS__KEY", SECRET_KEY)
             .with_env_var("CATALOG_S3_PATH__STYLE__ACCESS", "true")
-            .with_env_var("AWS_REGION", "us-east-1")
-            .with_env_var("AWS_ACCESS_KEY_ID", MINIO_ACCESS_KEY)
-            .with_env_var("AWS_SECRET_ACCESS_KEY", MINIO_SECRET_KEY)
+            .with_env_var("AWS_REGION", REGION)
+            .with_env_var("AWS_ACCESS_KEY_ID", ACCESS_KEY)
+            .with_env_var("AWS_SECRET_ACCESS_KEY", SECRET_KEY)
             .with_mapped_port(0, ICEBERG_REST_PORT.tcp())
             
.with_container_name(fixtures::unique_container_name("iceberg-rest"))
             .start()
@@ -184,51 +125,12 @@ impl IcebergRestContainer {
 
 pub struct IcebergFixture {
     #[allow(dead_code)]
-    minio: MinioContainer,
+    floci: FlociContainer,
     #[allow(dead_code)]
     iceberg_rest: IcebergRestContainer,
     http_client: HttpClient,
     pub catalog_url: String,
-    pub minio_endpoint: String,
-}
-
-impl IcebergFixture {
-    async fn create_bucket(minio_endpoint: &str) -> Result<(), 
TestBinaryError> {
-        use std::process::Command;
-
-        let host = minio_endpoint.trim_start_matches("http://";);
-        let mc_host = format!("http://{}:{}@{}";, MINIO_ACCESS_KEY, 
MINIO_SECRET_KEY, host);
-
-        let output = Command::new("docker")
-            .args([
-                "run",
-                "--rm",
-                "--network=host",
-                "-e",
-                &format!("MC_HOST_minio={}", mc_host),
-                "quay.io/minio/mc",
-                "mb",
-                "--ignore-existing",
-                &format!("minio/{}", MINIO_BUCKET),
-            ])
-            .output()
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "IcebergFixture".to_string(),
-                message: format!("Failed to run mc command: {error}"),
-            })?;
-
-        if !output.status.success() {
-            let stderr = String::from_utf8_lossy(&output.stderr);
-            let stdout = String::from_utf8_lossy(&output.stdout);
-            return Err(TestBinaryError::FixtureSetup {
-                fixture_type: "IcebergFixture".to_string(),
-                message: format!("Failed to create bucket: stderr={stderr}, 
stdout={stdout}"),
-            });
-        }
-
-        info!("Created MinIO bucket: {MINIO_BUCKET}");
-        Ok(())
-    }
+    pub floci_endpoint: String,
 }
 
 pub trait IcebergOps: Sync {
@@ -454,20 +356,20 @@ impl TestFixture for IcebergFixture {
     async fn setup() -> Result<Self, TestBinaryError> {
         let id = Uuid::new_v4();
         let network = format!("iggy-iceberg-{id}");
-        let minio_name = fixtures::unique_container_name("minio-iceberg");
+        let floci_name = fixtures::unique_container_name("floci-iceberg");
 
-        let minio = MinioContainer::start(&network, &minio_name).await?;
+        let floci = FlociContainer::start(Some(&network), &floci_name).await?;
 
-        Self::create_bucket(&minio.endpoint).await?;
+        floci::create_bucket(&floci.endpoint, TEST_BUCKET).await?;
 
-        let iceberg_rest = IcebergRestContainer::start(&network, 
&minio.internal_endpoint).await?;
+        let iceberg_rest = IcebergRestContainer::start(&network, 
&floci.internal_endpoint).await?;
 
         let http_client = create_http_client();
 
         Ok(Self {
             catalog_url: iceberg_rest.catalog_url.clone(),
-            minio_endpoint: minio.endpoint.clone(),
-            minio,
+            floci_endpoint: floci.endpoint.clone(),
+            floci,
             iceberg_rest,
             http_client,
         })
@@ -476,17 +378,14 @@ impl TestFixture for IcebergFixture {
     fn connectors_runtime_envs(&self) -> HashMap<String, String> {
         let mut envs = HashMap::new();
         envs.insert(ENV_SINK_URI.to_string(), self.catalog_url.clone());
-        envs.insert(ENV_SINK_WAREHOUSE.to_string(), MINIO_BUCKET.to_string());
-        envs.insert(ENV_SINK_STORE_URL.to_string(), 
self.minio_endpoint.clone());
+        envs.insert(ENV_SINK_WAREHOUSE.to_string(), TEST_BUCKET.to_string());
+        envs.insert(ENV_SINK_STORE_URL.to_string(), 
self.floci_endpoint.clone());
         envs.insert(
             ENV_SINK_STORE_ACCESS_KEY.to_string(),
-            MINIO_ACCESS_KEY.to_string(),
-        );
-        envs.insert(
-            ENV_SINK_STORE_SECRET.to_string(),
-            MINIO_SECRET_KEY.to_string(),
+            ACCESS_KEY.to_string(),
         );
-        envs.insert(ENV_SINK_STORE_REGION.to_string(), 
"us-east-1".to_string());
+        envs.insert(ENV_SINK_STORE_SECRET.to_string(), SECRET_KEY.to_string());
+        envs.insert(ENV_SINK_STORE_REGION.to_string(), REGION.to_string());
         envs.insert(
             ENV_SINK_PATH.to_string(),
             "../../target/debug/libiggy_connector_iceberg_sink".to_string(),
@@ -624,13 +523,10 @@ impl TestFixture for IcebergEnvAuthFixture {
         envs.remove(ENV_SINK_STORE_ACCESS_KEY);
         envs.remove(ENV_SINK_STORE_SECRET);
         // Inject standard AWS env vars to test the default credential 
provider chain.
-        envs.insert(
-            ENV_AWS_ACCESS_KEY_ID.to_string(),
-            MINIO_ACCESS_KEY.to_string(),
-        );
+        envs.insert(ENV_AWS_ACCESS_KEY_ID.to_string(), ACCESS_KEY.to_string());
         envs.insert(
             ENV_AWS_SECRET_ACCESS_KEY.to_string(),
-            MINIO_SECRET_KEY.to_string(),
+            SECRET_KEY.to_string(),
         );
         envs
     }
diff --git a/core/integration/tests/connectors/fixtures/mod.rs 
b/core/integration/tests/connectors/fixtures/mod.rs
index 05c7c824e..cddee53ab 100644
--- a/core/integration/tests/connectors/fixtures/mod.rs
+++ b/core/integration/tests/connectors/fixtures/mod.rs
@@ -21,6 +21,7 @@ mod clickhouse;
 mod delta;
 mod doris;
 mod elasticsearch;
+mod floci;
 mod http;
 mod iceberg;
 mod influxdb;
diff --git a/core/integration/tests/connectors/fixtures/redshift/container.rs 
b/core/integration/tests/connectors/fixtures/redshift/container.rs
index b15d1dfff..c13fb3b69 100644
--- a/core/integration/tests/connectors/fixtures/redshift/container.rs
+++ b/core/integration/tests/connectors/fixtures/redshift/container.rs
@@ -22,13 +22,14 @@ use pgwire::tokio::process_socket;
 use sqlx::{Pool, Postgres, postgres::PgPoolOptions};
 use testcontainers::{
     ContainerAsync, GenericImage, ImageExt,
-    core::{IntoContainerPort, WaitFor, wait::HttpWaitStrategy},
+    core::{IntoContainerPort, WaitFor},
     runners::AsyncRunner,
 };
 use tokio::{net::TcpListener, task::JoinHandle};
 
 use crate::connectors::fixtures::{
     self,
+    floci::{ACCESS_KEY, SECRET_KEY},
     redshift::redshift_mock::{handler::RedshiftHandlerFactory, load::S3Client},
 };
 
@@ -38,14 +39,7 @@ const POSTGRES_PORT: u16 = 5432;
 const POSTGRES_DB: &str = "postgres";
 const POSTGRES_USER: &str = "postgres";
 const POSTGRES_PASSWORD: &str = "postgres";
-const MINIO_IMAGE: &str = "quay.io/minio/minio";
-const MINIO_TAG: &str = "RELEASE.2025-09-07T16-13-09Z";
-const MINIO_PORT: u16 = 9000;
-const MINIO_CONSOLE_PORT: u16 = 9001;
-
-pub const MINIO_ACCESS_KEY: &str = "admin";
-pub const MINIO_SECRET_KEY: &str = "password";
-pub const MINIO_BUCKET: &str = "iggystaging";
+pub const STAGING_BUCKET: &str = "iggystaging";
 pub const DEFAULT_SINK_TABLE: &str = "iggy_messages";
 pub const STAGING_REGION: &str = "us-east-1";
 pub const STAGING_PREFIX: &str = "iggy/messages";
@@ -78,61 +72,6 @@ pub const DEFAULT_TEST_TOPIC: &str = "test_topic";
 pub const DEFAULT_POLL_ATTEMPTS: usize = 100;
 pub const DEFAULT_POLL_INTERVAL_MS: u64 = 50;
 
-pub struct MinioContainer {
-    #[allow(dead_code)]
-    container: ContainerAsync<GenericImage>,
-    pub endpoint: String,
-}
-
-impl MinioContainer {
-    pub async fn start(network: &str, container_name: &str) -> Result<Self, 
TestBinaryError> {
-        let container = GenericImage::new(MINIO_IMAGE, MINIO_TAG)
-            .with_exposed_port(MINIO_PORT.tcp())
-            .with_exposed_port(MINIO_CONSOLE_PORT.tcp())
-            .with_wait_for(WaitFor::http(
-                HttpWaitStrategy::new("/minio/health/live")
-                    .with_port(MINIO_PORT.tcp())
-                    .with_expected_status_code(200u16),
-            ))
-            .with_network(network)
-            .with_container_name(container_name)
-            .with_env_var("MINIO_ROOT_USER", MINIO_ACCESS_KEY)
-            .with_env_var("MINIO_ROOT_PASSWORD", MINIO_SECRET_KEY)
-            .with_cmd(vec!["server", "/data", "--console-address", ":9001"])
-            .with_mapped_port(0, MINIO_PORT.tcp())
-            .with_mapped_port(0, MINIO_CONSOLE_PORT.tcp())
-            .start()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "MinioContainer".to_string(),
-                message: format!("Failed to start container: {error}"),
-            })?;
-
-        tracing::info!("Started MinIO container");
-
-        let mapped_port = container
-            .ports()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "MinioContainer".to_string(),
-                message: format!("Failed to get ports: {error}"),
-            })?
-            .map_to_host_port_ipv4(MINIO_PORT)
-            .ok_or_else(|| TestBinaryError::FixtureSetup {
-                fixture_type: "MinioContainer".to_string(),
-                message: "No mapping for MinIO port".to_string(),
-            })?;
-
-        let endpoint = format!("http://localhost:{mapped_port}";);
-        tracing::info!("MinIO container available at {endpoint}");
-
-        Ok(Self {
-            container,
-            endpoint,
-        })
-    }
-}
-
 /// Base container management for PostgreSQL fixtures.
 pub struct PostgresContainer {
     #[allow(dead_code)]
@@ -198,10 +137,10 @@ impl RedshiftContainer {
         s3_endpoint: String,
     ) -> Result<Self, TestBinaryError> {
         let s3_client = S3Client::new(
-            MINIO_BUCKET,
+            STAGING_BUCKET,
             &s3_endpoint,
-            MINIO_ACCESS_KEY,
-            MINIO_SECRET_KEY,
+            ACCESS_KEY,
+            SECRET_KEY,
             STAGING_REGION,
         )
         .await
diff --git a/core/integration/tests/connectors/fixtures/redshift/mod.rs 
b/core/integration/tests/connectors/fixtures/redshift/mod.rs
index 0bfaa2f48..4a7b2d2fa 100644
--- a/core/integration/tests/connectors/fixtures/redshift/mod.rs
+++ b/core/integration/tests/connectors/fixtures/redshift/mod.rs
@@ -19,7 +19,7 @@ mod container;
 mod redshift_mock;
 mod sink;
 
-pub use container::{MinioContainer, PostgresContainer, RedshiftContainer};
+pub use container::{PostgresContainer, RedshiftContainer};
 pub use sink::{
     RedshiftSinkFixture, RedshiftSinkJsonFixture, RedshiftSinkNoArchiveFixture,
     RedshiftSinkVarbyteFixture,
diff --git a/core/integration/tests/connectors/fixtures/redshift/sink.rs 
b/core/integration/tests/connectors/fixtures/redshift/sink.rs
index 664b14e85..55532d235 100644
--- a/core/integration/tests/connectors/fixtures/redshift/sink.rs
+++ b/core/integration/tests/connectors/fixtures/redshift/sink.rs
@@ -19,13 +19,14 @@ use std::{collections::HashMap, time::Duration};
 
 use async_trait::async_trait;
 use integration::harness::{TestBinaryError, TestFixture};
+use s3::Bucket;
 use sqlx::{Pool, Postgres};
-use uuid::Uuid;
 
 use crate::connectors::fixtures::{
     self,
+    floci::{self, ACCESS_KEY, FlociContainer, SECRET_KEY},
     redshift::{
-        MinioContainer, PostgresContainer, RedshiftContainer,
+        PostgresContainer, RedshiftContainer,
         container::{
             AWS_IAM_ROLE, DEFAULT_POLL_ATTEMPTS, DEFAULT_POLL_INTERVAL_MS, 
DEFAULT_SINK_TABLE,
             DEFAULT_TEST_STREAM, DEFAULT_TEST_TOPIC, ENV_SINK_ARCHIVE, 
ENV_SINK_AWS_IAM_ROLE,
@@ -33,21 +34,21 @@ use crate::connectors::fixtures::{
             ENV_SINK_S3_ENDPOINT, ENV_SINK_S3_PREFIX, 
ENV_SINK_STAGING_ACCESS_KEY,
             ENV_SINK_STAGING_REGION, ENV_SINK_STAGING_SECRET, 
ENV_SINK_STREAMS_0_CONSUMER_GROUP,
             ENV_SINK_STREAMS_0_SCHEMA, ENV_SINK_STREAMS_0_STREAM, 
ENV_SINK_STREAMS_0_TOPICS,
-            ENV_SINK_TARGET_TABLE, MINIO_ACCESS_KEY, MINIO_BUCKET, 
MINIO_SECRET_KEY,
-            STAGING_PREFIX, STAGING_REGION, SinkPayloadFormat, SinkSchema,
+            ENV_SINK_TARGET_TABLE, STAGING_BUCKET, STAGING_PREFIX, 
STAGING_REGION,
+            SinkPayloadFormat, SinkSchema,
         },
     },
 };
 
 pub struct RedshiftSinkFixture {
     #[allow(dead_code)]
-    minio: MinioContainer,
+    floci: FlociContainer,
     #[allow(dead_code)]
     redshift: RedshiftContainer,
     postgres: PostgresContainer,
     payload_format: SinkPayloadFormat,
     schema: SinkSchema,
-    pub minio_endpoint: String,
+    pub floci_endpoint: String,
 }
 
 impl RedshiftSinkFixture {
@@ -101,22 +102,19 @@ impl TestFixture for RedshiftSinkFixture {
     async fn setup() -> Result<Self, TestBinaryError> {
         let postgres = PostgresContainer::start().await?;
 
-        let id = Uuid::now_v7();
-        let network = format!("iggy-redshift-{id}");
+        let floci_name = fixtures::unique_container_name("floci-redshift");
 
-        let minio_name = fixtures::unique_container_name("minio-redshift");
+        let floci = FlociContainer::start(None, &floci_name).await?;
 
-        let minio = MinioContainer::start(&network, &minio_name).await?;
-
-        create_bucket(&minio.endpoint)?;
+        floci::create_bucket(&floci.endpoint, STAGING_BUCKET).await?;
 
         let redshift =
-            RedshiftContainer::start(postgres.connection_string.clone(), 
minio.endpoint.clone())
+            RedshiftContainer::start(postgres.connection_string.clone(), 
floci.endpoint.clone())
                 .await?;
 
         Ok(Self {
-            minio_endpoint: minio.endpoint.clone(),
-            minio,
+            floci_endpoint: floci.endpoint.clone(),
+            floci,
             redshift,
             postgres,
             payload_format: SinkPayloadFormat::default(),
@@ -140,20 +138,17 @@ impl TestFixture for RedshiftSinkFixture {
 
         envs.insert(
             ENV_SINK_STAGING_ACCESS_KEY.to_string(),
-            MINIO_ACCESS_KEY.to_string(),
+            ACCESS_KEY.to_string(),
         );
 
-        envs.insert(
-            ENV_SINK_STAGING_SECRET.to_string(),
-            MINIO_SECRET_KEY.to_string(),
-        );
+        envs.insert(ENV_SINK_STAGING_SECRET.to_string(), 
SECRET_KEY.to_string());
 
-        envs.insert(ENV_SINK_S3_BUCKET.to_string(), MINIO_BUCKET.to_string());
+        envs.insert(ENV_SINK_S3_BUCKET.to_string(), 
STAGING_BUCKET.to_string());
         envs.insert(ENV_SINK_S3_PREFIX.to_string(), 
STAGING_PREFIX.to_string());
 
         envs.insert(
             ENV_SINK_S3_ENDPOINT.to_string(),
-            self.minio_endpoint.clone(),
+            self.floci_endpoint.clone(),
         );
 
         envs.insert(
@@ -215,23 +210,20 @@ impl TestFixture for RedshiftSinkVarbyteFixture {
     async fn setup() -> Result<Self, TestBinaryError> {
         let postgres = PostgresContainer::start().await?;
 
-        let id = Uuid::now_v7();
-        let network = format!("iggy-redshift-{id}");
-
-        let minio_name = fixtures::unique_container_name("minio-redshift");
+        let floci_name = fixtures::unique_container_name("floci-redshift");
 
-        let minio = MinioContainer::start(&network, &minio_name).await?;
+        let floci = FlociContainer::start(None, &floci_name).await?;
 
-        create_bucket(&minio.endpoint)?;
+        floci::create_bucket(&floci.endpoint, STAGING_BUCKET).await?;
 
         let redshift =
-            RedshiftContainer::start(postgres.connection_string.clone(), 
minio.endpoint.clone())
+            RedshiftContainer::start(postgres.connection_string.clone(), 
floci.endpoint.clone())
                 .await?;
 
         Ok(Self {
             inner: RedshiftSinkFixture {
-                minio_endpoint: minio.endpoint.clone(),
-                minio,
+                floci_endpoint: floci.endpoint.clone(),
+                floci,
                 redshift,
                 postgres,
                 payload_format: SinkPayloadFormat::Varbyte,
@@ -262,23 +254,20 @@ impl TestFixture for RedshiftSinkJsonFixture {
     async fn setup() -> Result<Self, TestBinaryError> {
         let postgres = PostgresContainer::start().await?;
 
-        let id = Uuid::now_v7();
-        let network = format!("iggy-redshift-{id}");
+        let floci_name = fixtures::unique_container_name("floci-redshift");
 
-        let minio_name = fixtures::unique_container_name("minio-redshift");
+        let floci = FlociContainer::start(None, &floci_name).await?;
 
-        let minio = MinioContainer::start(&network, &minio_name).await?;
-
-        create_bucket(&minio.endpoint)?;
+        floci::create_bucket(&floci.endpoint, STAGING_BUCKET).await?;
 
         let redshift =
-            RedshiftContainer::start(postgres.connection_string.clone(), 
minio.endpoint.clone())
+            RedshiftContainer::start(postgres.connection_string.clone(), 
floci.endpoint.clone())
                 .await?;
 
         Ok(Self {
             inner: RedshiftSinkFixture {
-                minio_endpoint: minio.endpoint.clone(),
-                minio,
+                floci_endpoint: floci.endpoint.clone(),
+                floci,
                 redshift,
                 postgres,
                 payload_format: SinkPayloadFormat::Text,
@@ -295,11 +284,12 @@ impl TestFixture for RedshiftSinkJsonFixture {
 /// redshift sink fixture for bytea payload format.
 pub struct RedshiftSinkNoArchiveFixture {
     inner: RedshiftSinkFixture,
+    bucket: Box<Bucket>,
 }
 
 impl RedshiftSinkNoArchiveFixture {
     pub async fn confirm_empty_bucket(&self) -> Result<bool, TestBinaryError> {
-        bucket_empty(&self.minio_endpoint).await
+        bucket_empty(&self.bucket).await
     }
 }
 
@@ -315,23 +305,21 @@ impl TestFixture for RedshiftSinkNoArchiveFixture {
     async fn setup() -> Result<Self, TestBinaryError> {
         let postgres = PostgresContainer::start().await?;
 
-        let id = Uuid::now_v7();
-        let network = format!("iggy-redshift-{id}");
-
-        let minio_name = fixtures::unique_container_name("minio-redshift");
+        let floci_name = fixtures::unique_container_name("floci-redshift");
 
-        let minio = MinioContainer::start(&network, &minio_name).await?;
+        let floci = FlociContainer::start(None, &floci_name).await?;
 
-        create_bucket(&minio.endpoint)?;
+        let bucket = floci::create_bucket(&floci.endpoint, 
STAGING_BUCKET).await?;
 
         let redshift =
-            RedshiftContainer::start(postgres.connection_string.clone(), 
minio.endpoint.clone())
+            RedshiftContainer::start(postgres.connection_string.clone(), 
floci.endpoint.clone())
                 .await?;
 
         Ok(Self {
+            bucket,
             inner: RedshiftSinkFixture {
-                minio_endpoint: minio.endpoint.clone(),
-                minio,
+                floci_endpoint: floci.endpoint.clone(),
+                floci,
                 redshift,
                 postgres,
                 payload_format: SinkPayloadFormat::Text,
@@ -365,116 +353,17 @@ impl TestFixture for RedshiftSinkNoArchiveFixture {
     }
 }
 
-fn create_bucket(minio_endpoint: &str) -> Result<(), TestBinaryError> {
-    use std::process::Command;
-
-    let host = minio_endpoint.trim_start_matches("http://";);
-    let mc_host = format!("http://{}:{}@{}";, MINIO_ACCESS_KEY, 
MINIO_SECRET_KEY, host);
-
-    let output = Command::new("docker")
-        .args([
-            "run",
-            "--rm",
-            "--network=host",
-            "-e",
-            &format!("MC_HOST_minio={}", mc_host),
-            "quay.io/minio/mc",
-            "mb",
-            "--ignore-existing",
-            &format!("minio/{}", MINIO_BUCKET),
-        ])
-        .output()
-        .map_err(|error| TestBinaryError::FixtureSetup {
-            fixture_type: "RedshiftFixture".to_string(),
-            message: format!("Failed to run mc command: {error}"),
-        })?;
-
-    if !output.status.success() {
-        let stderr = String::from_utf8_lossy(&output.stderr);
-        let stdout = String::from_utf8_lossy(&output.stdout);
-        return Err(TestBinaryError::FixtureSetup {
-            fixture_type: "IcebergFixture".to_string(),
-            message: format!("Failed to create bucket: stderr={stderr}, 
stdout={stdout}"),
-        });
-    }
-
-    tracing::info!("Created MinIO bucket: {MINIO_BUCKET}");
-    Ok(())
-}
-
-async fn bucket_empty(minio_endpoint: &str) -> Result<bool, TestBinaryError> {
-    use std::process::Command;
-
-    let host = minio_endpoint.trim_start_matches("http://";);
-    let mc_host = format!("http://{}:{}@{}";, MINIO_ACCESS_KEY, 
MINIO_SECRET_KEY, host);
-
-    let mut result = false;
-
+async fn bucket_empty(bucket: &Bucket) -> Result<bool, TestBinaryError> {
     for _ in 0..DEFAULT_POLL_ATTEMPTS {
-        let output = Command::new("docker")
-            .args([
-                "run",
-                "--rm",
-                "--network=host",
-                "-e",
-                &format!("MC_HOST_minio={}", mc_host),
-                "quay.io/minio/mc",
-                "ls",
-                &format!("minio/{}", MINIO_BUCKET),
-            ])
-            .output()
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "RedshiftFixture".to_string(),
-                message: format!("Failed to run mc command: {error}"),
-            })?;
-
-        if !output.status.success() {
-            let stderr = String::from_utf8_lossy(&output.stderr);
-            let stdout = String::from_utf8_lossy(&output.stdout);
-            return Err(TestBinaryError::FixtureSetup {
-                fixture_type: "IcebergFixture".to_string(),
-                message: format!("Failed to create bucket: stderr={stderr}, 
stdout={stdout}"),
-            });
-        }
-
-        tracing::info!("Checking bucket contents");
-
-        // When testing live S3
-
-        // let credentials = s3::creds::Credentials::new(
-        //     Some("AKIARI2XXXXXX"),
-        //     Some("X82flw9cR2iTYDS5nEFgXXXXX"),
-        //     None,
-        //     None,
-        //     None,
-        // )
-        // .expect("Failed to create credentials");
-
-        // let region = s3::Region::from_str(STAGING_REGION).expect("Failed to 
parse region");
-
-        // let bucket = s3::Bucket::new(MINIO_BUCKET, region, 
credentials).unwrap();
-
-        // let (results, _code) = bucket
-        //     .list_page(String::from(STAGING_PREFIX), None, None, None, 
Some(1))
-        //     .await
-        //     .unwrap();
-
-        // tracing::info!("Bucket contents: {:?}", results.contents);
-
-        // results.contents.is_empty();
-
-        match output.stdout.is_empty() {
-            true => {
-                return Ok(true);
-            }
-
-            false => {
-                result = false;
+        let listings = bucket.list(String::new(), None).await.map_err(|error| {
+            TestBinaryError::InvalidState {
+                message: format!("Failed to list staging bucket: {error}"),
             }
+        })?;
+        if listings.iter().all(|listing| listing.contents.is_empty()) {
+            return Ok(true);
         }
-
         
tokio::time::sleep(Duration::from_millis(DEFAULT_POLL_INTERVAL_MS)).await;
     }
-
-    Ok(result)
+    Ok(false)
 }
diff --git a/core/integration/tests/connectors/fixtures/s3/fixture.rs 
b/core/integration/tests/connectors/fixtures/s3/fixture.rs
index 876053248..9f641968e 100644
--- a/core/integration/tests/connectors/fixtures/s3/fixture.rs
+++ b/core/integration/tests/connectors/fixtures/s3/fixture.rs
@@ -15,31 +15,20 @@
 // specific language governing permissions and limitations
 // under the License.
 
+use std::collections::HashMap;
+
 use async_trait::async_trait;
 use integration::harness::seeds;
 use integration::harness::{TestBinaryError, TestFixture};
-use s3::creds::Credentials;
-use s3::{Bucket, Region};
-use std::collections::HashMap;
-use std::time::Duration;
-use testcontainers_modules::testcontainers::core::wait::HttpWaitStrategy;
-use testcontainers_modules::testcontainers::core::{IntoContainerPort, WaitFor};
-use testcontainers_modules::testcontainers::runners::AsyncRunner;
-use testcontainers_modules::testcontainers::{ContainerAsync, GenericImage, 
ImageExt};
+use s3::Bucket;
 use tracing::info;
-use uuid::Uuid;
 
-const MINIO_IMAGE: &str = "quay.io/minio/minio";
-const MINIO_TAG: &str = "RELEASE.2025-09-07T16-13-09Z";
-const MINIO_PORT: u16 = 9000;
-const MINIO_CONSOLE_PORT: u16 = 9001;
+use crate::connectors::fixtures::{
+    self,
+    floci::{self, ACCESS_KEY, FlociContainer, REGION, SECRET_KEY},
+};
 
-const MINIO_ACCESS_KEY: &str = "admin";
-const MINIO_SECRET_KEY: &str = "password";
-const MINIO_BUCKET: &str = "iggy-s3-test";
-/// Bounds the wait for MinIO's S3 API to come up behind its health endpoint.
-const BUCKET_CREATE_ATTEMPTS: u32 = 30;
-const BUCKET_CREATE_RETRY_DELAY: Duration = Duration::from_secs(1);
+const TEST_BUCKET: &str = "iggy-s3-test";
 
 const ENV_SINK_PATH: &str = "IGGY_CONNECTORS_SINK_S3_PATH";
 const ENV_SINK_STREAMS_0_STREAM: &str = 
"IGGY_CONNECTORS_SINK_S3_STREAMS_0_STREAM";
@@ -134,7 +123,7 @@ pub trait S3SinkOps: Sync {
 
 pub struct S3SinkFixture {
     #[allow(dead_code)]
-    container: ContainerAsync<GenericImage>,
+    floci: FlociContainer,
     bucket: Box<Bucket>,
     endpoint: String,
 }
@@ -152,111 +141,13 @@ impl S3SinkOps for S3SinkFixture {
 #[async_trait]
 impl TestFixture for S3SinkFixture {
     async fn setup() -> Result<Self, TestBinaryError> {
-        let id = Uuid::new_v4();
-        let container_name = format!("minio-s3-{id}");
-
-        let container = GenericImage::new(MINIO_IMAGE, MINIO_TAG)
-            .with_exposed_port(MINIO_PORT.tcp())
-            .with_exposed_port(MINIO_CONSOLE_PORT.tcp())
-            .with_wait_for(WaitFor::http(
-                HttpWaitStrategy::new("/minio/health/live")
-                    .with_port(MINIO_PORT.tcp())
-                    .with_expected_status_code(200u16),
-            ))
-            .with_container_name(&container_name)
-            .with_env_var("MINIO_ROOT_USER", MINIO_ACCESS_KEY)
-            .with_env_var("MINIO_ROOT_PASSWORD", MINIO_SECRET_KEY)
-            .with_cmd(vec!["server", "/data", "--console-address", ":9001"])
-            .with_mapped_port(0, MINIO_PORT.tcp())
-            .with_mapped_port(0, MINIO_CONSOLE_PORT.tcp())
-            .start()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "S3SinkFixture".to_string(),
-                message: format!("Failed to start MinIO container: {error}"),
-            })?;
-
-        let mapped_port = container
-            .ports()
-            .await
-            .map_err(|error| TestBinaryError::FixtureSetup {
-                fixture_type: "S3SinkFixture".to_string(),
-                message: format!("Failed to get ports: {error}"),
-            })?
-            .map_to_host_port_ipv4(MINIO_PORT)
-            .ok_or_else(|| TestBinaryError::FixtureSetup {
-                fixture_type: "S3SinkFixture".to_string(),
-                message: "No mapping for MinIO port".to_string(),
-            })?;
-
-        let endpoint = format!("http://localhost:{mapped_port}";);
-        info!("MinIO container for S3 sink available at {endpoint}");
-
-        let region = Region::Custom {
-            region: "us-east-1".to_string(),
-            endpoint: endpoint.clone(),
-        };
-        let credentials = Credentials::new(
-            Some(MINIO_ACCESS_KEY),
-            Some(MINIO_SECRET_KEY),
-            None,
-            None,
-            None,
-        )
-        .map_err(|e| TestBinaryError::FixtureSetup {
-            fixture_type: "S3SinkFixture".to_string(),
-            message: format!("Failed to create credentials: {e}"),
-        })?;
-
-        // MinIO answers on its health endpoint before it can serve the S3 API,
-        // so bucket creation can come back 503 while it finishes starting. The
-        // call itself is `Ok` in that case -- the status lives in the response
-        // -- so taking it as success left the bucket absent, and the first
-        // `list_objects` then parsed an S3 error document as a listing and
-        // failed with the unhelpful `missing field 'Name'`. Retry until the
-        // status is a real one: 2xx created, 409 already owned by us.
-        let mut last_status = 0;
-        let mut created = false;
-        for _ in 0..BUCKET_CREATE_ATTEMPTS {
-            let response = Bucket::create_with_path_style(
-                MINIO_BUCKET,
-                region.clone(),
-                credentials.clone(),
-                s3::BucketConfiguration::default(),
-            )
-            .await
-            .map_err(|e| TestBinaryError::FixtureSetup {
-                fixture_type: "S3SinkFixture".to_string(),
-                message: format!("Failed to create bucket: {e}"),
-            })?;
-            last_status = response.response_code;
-            if (200..300).contains(&last_status) || last_status == 409 {
-                created = true;
-                break;
-            }
-            tokio::time::sleep(BUCKET_CREATE_RETRY_DELAY).await;
-        }
-        if !created {
-            return Err(TestBinaryError::FixtureSetup {
-                fixture_type: "S3SinkFixture".to_string(),
-                message: format!(
-                    "Bucket '{MINIO_BUCKET}' not creatable after \
-                     {BUCKET_CREATE_ATTEMPTS} attempts (last status: 
{last_status})"
-                ),
-            });
-        }
-        info!("S3 bucket '{MINIO_BUCKET}' ready (status: {last_status})");
-
-        let mut bucket = Bucket::new(MINIO_BUCKET, region, 
credentials).map_err(|e| {
-            TestBinaryError::FixtureSetup {
-                fixture_type: "S3SinkFixture".to_string(),
-                message: format!("Failed to create bucket handle: {e}"),
-            }
-        })?;
-        bucket.set_path_style();
+        let floci =
+            FlociContainer::start(None, 
&fixtures::unique_container_name("floci-s3")).await?;
+        let endpoint = floci.endpoint.clone();
+        let bucket = floci::create_bucket(&endpoint, TEST_BUCKET).await?;
 
         Ok(Self {
-            container,
+            floci,
             bucket,
             endpoint,
         })
@@ -277,17 +168,17 @@ impl TestFixture for S3SinkFixture {
             format!("[{}]", seeds::names::TOPIC),
         );
         envs.insert(ENV_SINK_STREAMS_0_SCHEMA.to_string(), "json".to_string());
-        envs.insert(ENV_SINK_PLUGIN_BUCKET.to_string(), 
MINIO_BUCKET.to_string());
-        envs.insert(ENV_SINK_PLUGIN_REGION.to_string(), 
"us-east-1".to_string());
+        envs.insert(ENV_SINK_PLUGIN_BUCKET.to_string(), 
TEST_BUCKET.to_string());
+        envs.insert(ENV_SINK_PLUGIN_REGION.to_string(), REGION.to_string());
         envs.insert(ENV_SINK_PLUGIN_ENDPOINT.to_string(), 
self.endpoint.clone());
         envs.insert(ENV_SINK_PLUGIN_PREFIX.to_string(), String::new());
         envs.insert(
             ENV_SINK_PLUGIN_ACCESS_KEY.to_string(),
-            MINIO_ACCESS_KEY.to_string(),
+            ACCESS_KEY.to_string(),
         );
         envs.insert(
             ENV_SINK_PLUGIN_SECRET_KEY.to_string(),
-            MINIO_SECRET_KEY.to_string(),
+            SECRET_KEY.to_string(),
         );
         envs.insert(
             ENV_SINK_PLUGIN_FILE_ROTATION.to_string(),
diff --git a/core/integration/tests/connectors/s3/s3_sink.rs 
b/core/integration/tests/connectors/s3/s3_sink.rs
index a3352d06e..9ede317a1 100644
--- a/core/integration/tests/connectors/s3/s3_sink.rs
+++ b/core/integration/tests/connectors/s3/s3_sink.rs
@@ -15,8 +15,6 @@
 // specific language governing permissions and limitations
 // under the License.
 
-use crate::connectors::create_test_messages;
-use crate::connectors::fixtures::{S3SinkFixture, S3SinkOps, 
S3SinkRotationFixture};
 use bytes::Bytes;
 use iggy::prelude::{IggyMessage, Partitioning};
 use iggy_common::Identifier;
@@ -25,6 +23,13 @@ use iggy_connector_sdk::api::{ConnectorStatus, 
SinkInfoResponse};
 use integration::harness::seeds;
 use integration::iggy_harness;
 use reqwest::Client;
+use s3::{
+    command::Command,
+    request::{Request, tokio_backend::ReqwestRequest},
+};
+
+use crate::connectors::create_test_messages;
+use crate::connectors::fixtures::{S3SinkFixture, S3SinkOps, 
S3SinkRotationFixture};
 
 const API_KEY: &str = "test-api-key";
 const S3_SINK_KEY: &str = "s3";
@@ -60,14 +65,28 @@ async fn s3_sink_initializes_and_runs(harness: 
&TestHarness, fixture: S3SinkFixt
         keys.is_empty(),
         "Startup must not publish probe objects: {keys:?}"
     );
-    let uploads = fixture
-        .bucket()
-        .list_multiparts_uploads(None, None)
+    // Floci 2.1.0 omits IsTruncated from empty listings, which rust-s3 
requires.
+    let uploads_request = ReqwestRequest::new(
+        fixture.bucket(),
+        "/",
+        Command::ListMultipartUploads {
+            prefix: None,
+            delimiter: None,
+            key_marker: None,
+            max_uploads: None,
+        },
+    )
+    .await
+    .expect("Prepare multipart listing");
+    let uploads_response = uploads_request
+        .response_data(false)
         .await
         .expect("List multipart uploads");
+    assert_eq!(uploads_response.status_code(), 200);
+    let uploads_xml = uploads_response.as_str().expect("Read multipart 
uploads");
     assert!(
-        uploads.iter().all(|page| page.uploads.is_empty()),
-        "Startup must abort its multipart probe: {uploads:?}"
+        !uploads_xml.contains("<Upload"),
+        "Startup must abort its multipart probe: {uploads_xml}"
     );
 
     drop(fixture);

Reply via email to