diff --git a/Cargo.lock b/Cargo.lock index 64c5c61875..2e24b86e52 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2662,7 +2662,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.5.10", + "socket2 0.6.4", "system-configuration", "tokio", "tower-service", @@ -2925,6 +2925,26 @@ dependencies = [ "trybuild", ] +[[package]] +name = "iceberg-storage-object_store" +version = "0.10.1" +dependencies = [ + "async-trait", + "bytes", + "dashmap", + "futures", + "iceberg", + "iceberg_test_utils", + "object_store", + "percent-encoding", + "serde", + "serde_json", + "tokio", + "tracing", + "typetag", + "url", +] + [[package]] name = "iceberg-storage-opendal" version = "0.10.1" @@ -3876,6 +3896,44 @@ dependencies = [ "objc2-core-foundation", ] +[[package]] +name = "object_store" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "622acbc9100d3c10e2ee15804b0caa40e55c933d5aa53814cd520805b7958a49" +dependencies = [ + "async-trait", + "base64 0.22.1", + "bytes", + "chrono", + "form_urlencoded", + "futures-channel", + "futures-core", + "futures-util", + "http 1.5.0", + "http-body-util", + "humantime", + "hyper", + "itertools 0.14.0", + "md-5 0.10.6", + "parking_lot", + "percent-encoding", + "quick-xml 0.39.4", + "rand 0.10.2", + "reqwest 0.12.28", + "ring", + "serde", + "serde_json", + "serde_urlencoded", + "thiserror 2.0.18", + "tokio", + "tracing", + "url", + "walkdir", + "wasm-bindgen-futures", + "web-time", +] + [[package]] name = "once_cell" version = "1.21.4" @@ -3937,7 +3995,7 @@ dependencies = [ "md-5 0.11.0", "mea", "percent-encoding", - "quick-xml", + "quick-xml 0.41.0", "reqsign-core", "serde", "serde_json", @@ -4017,7 +4075,7 @@ dependencies = [ "mea", "opendal-core", "opendal-service-azure-common", - "quick-xml", + "quick-xml 0.41.0", "reqsign-azure-storage", "reqsign-core", "reqsign-file-read-tokio", @@ -4061,7 +4119,7 @@ dependencies = [ "log", "opendal-core", "percent-encoding", - "quick-xml", + "quick-xml 0.41.0", "reqsign-core", "reqsign-file-read-tokio", "reqsign-google", @@ -4096,7 +4154,7 @@ dependencies = [ "http 1.5.0", "log", "opendal-core", - "quick-xml", + "quick-xml 0.41.0", "reqsign-aliyun-oss", "reqsign-core", "reqsign-file-read-tokio", @@ -4116,7 +4174,7 @@ dependencies = [ "log", "md-5 0.11.0", "opendal-core", - "quick-xml", + "quick-xml 0.41.0", "reqsign-aws-v4", "reqsign-core", "reqsign-file-read-tokio", @@ -4614,6 +4672,16 @@ version = "0.2.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "5a651516ddc9168ebd67b24afd085a718be02f8858fe406591b013d101ce2f40" +[[package]] +name = "quick-xml" +version = "0.39.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cdcc8dd4e2f670d309a5f0e83fe36dfdc05af317008fea29144da1a2ac858e5e" +dependencies = [ + "memchr", + "serde", +] + [[package]] name = "quick-xml" version = "0.41.0" @@ -4637,7 +4705,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls", - "socket2 0.5.10", + "socket2 0.6.4", "thiserror 2.0.18", "tokio", "tracing", @@ -4676,7 +4744,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.5.10", + "socket2 0.6.4", "tracing", "windows-sys 0.61.2", ] @@ -4932,7 +5000,7 @@ dependencies = [ "http 1.5.0", "log", "percent-encoding", - "quick-xml", + "quick-xml 0.41.0", "reqsign-core", "rust-ini", "serde", @@ -4950,7 +5018,7 @@ dependencies = [ "bytes", "http 1.5.0", "log", - "quick-xml", + "quick-xml 0.41.0", "reqsign-aws-core", "reqsign-core", "serde", @@ -5040,20 +5108,27 @@ dependencies = [ "bytes", "futures-core", "futures-util", + "h2", "http 1.5.0", "http-body 1.0.1", "http-body-util", "hyper", + "hyper-rustls", "hyper-util", "js-sys", "log", "percent-encoding", "pin-project-lite", + "quinn", + "rustls", + "rustls-native-certs", + "rustls-pki-types", "serde", "serde_json", "serde_urlencoded", "sync_wrapper", "tokio", + "tokio-rustls", "tokio-util", "tower", "tower-http", diff --git a/Cargo.toml b/Cargo.toml index 661c56b094..5a10e378ab 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -75,6 +75,7 @@ bytes = "1.11" cfg-if = "1" chrono = "0.4.41" crc32fast = "1" +dashmap = "6.1" derive_builder = "0.20" expect-test = "1" fastnum = { version = "0.7", default-features = false, features = [ @@ -95,6 +96,7 @@ iceberg-catalog-rest = { version = "0.10.0", path = "./crates/catalog/rest" } iceberg-catalog-s3tables = { version = "0.10.0", path = "./crates/catalog/s3tables" } iceberg-catalog-sql = { version = "0.10.0", path = "./crates/catalog/sql" } iceberg-property-macro = { version = "0.10.0", path = "./crates/property-macro" } +iceberg-storage-object_store = { version = "0.10.0", path = "./crates/storage/object_store" } iceberg-storage-opendal = { version = "0.10.0", path = "./crates/storage/opendal" } itertools = "0.13" linkedbytes = "0.1.8" @@ -104,10 +106,12 @@ mockall = "0.13.1" mockito = "1" motore-macros = "0.4.3" murmur3 = "0.5.2" +object_store = "0.13" once_cell = "1.20" opendal = "0.58" ordered-float = "4" parquet = "59.2" +percent-encoding = "2.3" pilota = "0.11.10" pretty_assertions = "1.4" proc-macro2 = "1" @@ -118,6 +122,10 @@ regex = "1.11.3" reqwest = { version = "0.12.12", default-features = false, features = ["json"] } roaring = { version = "0.11" } rstest = "0.26" +rustls = "0.23.45" +rustls-native-certs = "0.8" +rustls-pki-types = "1.15" +rustls-webpki = "0.103.15" serde = { version = "1.0.219", features = ["rc"] } serde_bytes = "0.11.17" serde_derive = "1.0.219" diff --git a/crates/storage/object_store/Cargo.toml b/crates/storage/object_store/Cargo.toml new file mode 100644 index 0000000000..9b2c869a2c --- /dev/null +++ b/crates/storage/object_store/Cargo.toml @@ -0,0 +1,54 @@ +# 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] +edition = { workspace = true } +license = { workspace = true } +name = "iceberg-storage-object_store" +publish = false +repository = { workspace = true } +version = { workspace = true } + +categories = ["database"] +description = "Apache Iceberg object_store storage implementation" +keywords = ["iceberg", "object-store", "storage", "s3"] + +[features] +default = ["object_store-s3"] +object_store-s3 = ["object_store/aws"] + +[dependencies] +async-trait = { workspace = true } +bytes = { workspace = true } +dashmap = { workspace = true } +futures = { workspace = true } +iceberg = { workspace = true } +object_store = { workspace = true } +percent-encoding = { workspace = true } +serde = { workspace = true, features = ["derive"] } +tokio = { workspace = true } +tracing = { workspace = true } +typetag = { workspace = true } +url = { workspace = true } + +[dev-dependencies] +iceberg_test_utils = { path = "../../test_utils", features = ["tests"] } +serde_json = { workspace = true } +tokio = { workspace = true, features = ["macros", "rt-multi-thread"] } + +[lints] +workspace = true diff --git a/crates/storage/object_store/LICENSE b/crates/storage/object_store/LICENSE new file mode 120000 index 0000000000..5853aaea53 --- /dev/null +++ b/crates/storage/object_store/LICENSE @@ -0,0 +1 @@ +../../../LICENSE \ No newline at end of file diff --git a/crates/storage/object_store/NOTICE b/crates/storage/object_store/NOTICE new file mode 120000 index 0000000000..295f6bdb3a --- /dev/null +++ b/crates/storage/object_store/NOTICE @@ -0,0 +1 @@ +../../../NOTICE \ No newline at end of file diff --git a/crates/storage/object_store/public-api.txt b/crates/storage/object_store/public-api.txt new file mode 100644 index 0000000000..bba9e1c2d3 --- /dev/null +++ b/crates/storage/object_store/public-api.txt @@ -0,0 +1,44 @@ +pub mod iceberg_storage_object_store +pub enum iceberg_storage_object_store::ObjectStoreStorage +pub iceberg_storage_object_store::ObjectStoreStorage::S3(iceberg_storage_object_store::S3Storage) +impl core::clone::Clone for iceberg_storage_object_store::ObjectStoreStorage +pub fn iceberg_storage_object_store::ObjectStoreStorage::clone(&self) -> iceberg_storage_object_store::ObjectStoreStorage +impl core::fmt::Debug for iceberg_storage_object_store::ObjectStoreStorage +pub fn iceberg_storage_object_store::ObjectStoreStorage::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl iceberg::io::storage::Storage for iceberg_storage_object_store::ObjectStoreStorage +pub fn iceberg_storage_object_store::ObjectStoreStorage::delete<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::delete_prefix<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::delete_stream<'life0, 'async_trait>(&'life0 self, paths: futures_core::stream::BoxStream<'static, alloc::string::String>) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::exists<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::metadata<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::new_input(&self, path: &str) -> iceberg::error::Result +pub fn iceberg_storage_object_store::ObjectStoreStorage::new_output(&self, path: &str) -> iceberg::error::Result +pub fn iceberg_storage_object_store::ObjectStoreStorage::read<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::reader<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::write<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str, bs: bytes::bytes::Bytes) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_storage_object_store::ObjectStoreStorage::writer<'life0, 'life1, 'async_trait>(&'life0 self, path: &'life1 str) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +impl serde_core::ser::Serialize for iceberg_storage_object_store::ObjectStoreStorage +pub fn iceberg_storage_object_store::ObjectStoreStorage::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer +impl<'de> serde_core::de::Deserialize<'de> for iceberg_storage_object_store::ObjectStoreStorage +pub fn iceberg_storage_object_store::ObjectStoreStorage::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> +pub enum iceberg_storage_object_store::ObjectStoreStorageFactory +pub iceberg_storage_object_store::ObjectStoreStorageFactory::S3 +impl core::clone::Clone for iceberg_storage_object_store::ObjectStoreStorageFactory +pub fn iceberg_storage_object_store::ObjectStoreStorageFactory::clone(&self) -> iceberg_storage_object_store::ObjectStoreStorageFactory +impl core::fmt::Debug for iceberg_storage_object_store::ObjectStoreStorageFactory +pub fn iceberg_storage_object_store::ObjectStoreStorageFactory::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl iceberg::io::storage::StorageFactory for iceberg_storage_object_store::ObjectStoreStorageFactory +pub fn iceberg_storage_object_store::ObjectStoreStorageFactory::build(&self, config: &iceberg::io::storage::config::StorageConfig) -> iceberg::error::Result> +impl serde_core::ser::Serialize for iceberg_storage_object_store::ObjectStoreStorageFactory +pub fn iceberg_storage_object_store::ObjectStoreStorageFactory::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer +impl<'de> serde_core::de::Deserialize<'de> for iceberg_storage_object_store::ObjectStoreStorageFactory +pub fn iceberg_storage_object_store::ObjectStoreStorageFactory::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> +pub struct iceberg_storage_object_store::S3Storage +impl core::clone::Clone for iceberg_storage_object_store::S3Storage +pub fn iceberg_storage_object_store::S3Storage::clone(&self) -> iceberg_storage_object_store::S3Storage +impl core::fmt::Debug for iceberg_storage_object_store::S3Storage +pub fn iceberg_storage_object_store::S3Storage::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl serde_core::ser::Serialize for iceberg_storage_object_store::S3Storage +pub fn iceberg_storage_object_store::S3Storage::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer +impl<'de> serde_core::de::Deserialize<'de> for iceberg_storage_object_store::S3Storage +pub fn iceberg_storage_object_store::S3Storage::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> diff --git a/crates/storage/object_store/src/lib.rs b/crates/storage/object_store/src/lib.rs new file mode 100644 index 0000000000..aaae42cec7 --- /dev/null +++ b/crates/storage/object_store/src/lib.rs @@ -0,0 +1,971 @@ +// 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. + +//! `object_store`-based storage implementation for Apache Iceberg. +//! +//! This crate provides [`ObjectStoreStorage`] and [`ObjectStoreStorageFactory`], +//! which implement the [`Storage`] and +//! [`StorageFactory`] traits from the `iceberg` crate +//! using the [`object_store`](https://docs.rs/object_store) crate as the backend. +//! +//! Currently only S3 storage is supported (via the `object_store-s3` feature flag, +//! enabled by default). + +#[cfg(feature = "object_store-s3")] +mod s3; + +use std::ops::Range; +use std::sync::Arc; + +use async_trait::async_trait; +use bytes::{Bytes, BytesMut}; +use dashmap::DashMap; +use futures::stream::{BoxStream, FuturesUnordered}; +use futures::{StreamExt, TryStreamExt}; +#[cfg(feature = "object_store-s3")] +use iceberg::io::S3Config; +use iceberg::io::{ + FileMetadata, FileRead, FileWrite, InputFile, OutputFile, Storage, StorageConfig, + StorageFactory, +}; +use iceberg::{Error, ErrorKind, Result}; +use object_store::path::Path as ObjectStorePath; +use object_store::{MultipartUpload, ObjectStore, ObjectStoreExt, PutPayload, UploadPart}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; + +/// Convert an `object_store::Error` into an `iceberg::Error`, +/// dispatching known variants to their corresponding `ErrorKind`. +fn from_object_store_error(e: object_store::Error) -> Error { + let (kind, msg) = match &e { + object_store::Error::NotFound { path, .. } => { + (ErrorKind::Unexpected, format!("Object not found: {path}")) + } + object_store::Error::AlreadyExists { path, .. } => ( + ErrorKind::Unexpected, + format!("Object already exists: {path}"), + ), + object_store::Error::PermissionDenied { path, .. } => { + (ErrorKind::Unexpected, format!("Permission denied: {path}")) + } + object_store::Error::Unauthenticated { path, .. } => { + (ErrorKind::Unexpected, format!("Unauthenticated: {path}")) + } + object_store::Error::NotSupported { .. } => ( + ErrorKind::FeatureUnsupported, + "Operation not supported".to_string(), + ), + _ => ( + ErrorKind::Unexpected, + "Failure in doing io operation".to_string(), + ), + }; + Error::new(kind, msg).with_source(e) +} + +/// Property key for configuring S3 bulk delete batch size (matches Iceberg Java specification). +pub const S3_DELETE_BATCH_SIZE: &str = "s3.delete.batch-size"; +/// Default batch size for S3 bulk deletions (matches Iceberg Java specification). +pub const DEFAULT_DELETE_BATCH_SIZE: usize = 250; +/// Maximum batch size allowed by the AWS S3 DeleteObjects API specification. +pub const S3_MAX_DELETE_BATCH_SIZE: usize = 1000; + +fn parse_delete_batch_size(config: &StorageConfig) -> usize { + if let Some(val) = config.get(S3_DELETE_BATCH_SIZE) { + match val.parse::() { + Ok(parsed) if parsed > 0 => parsed.clamp(1, S3_MAX_DELETE_BATCH_SIZE), + _ => { + tracing::warn!( + val = %val, + "Invalid s3.delete.batch-size; falling back to default 250" + ); + DEFAULT_DELETE_BATCH_SIZE + } + } + } else { + DEFAULT_DELETE_BATCH_SIZE + } +} + +fn default_delete_batch_size() -> usize { + DEFAULT_DELETE_BATCH_SIZE +} + +/// Convert `object_store::ObjectMeta` into `iceberg::io::FileMetadata`. +fn to_file_metadata(meta: object_store::ObjectMeta) -> FileMetadata { + FileMetadata { size: meta.size } +} + +/// `object_store`-based storage factory. +/// +/// Use this factory with `FileIOBuilder::new(factory)` to create FileIO instances +/// backed by the `object_store` crate. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub enum ObjectStoreStorageFactory { + /// S3 storage factory. + #[cfg(feature = "object_store-s3")] + S3, +} + +#[typetag::serde(name = "ObjectStoreStorageFactory")] +impl StorageFactory for ObjectStoreStorageFactory { + #[allow(unused_variables)] + fn build(&self, config: &StorageConfig) -> Result> { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorageFactory::S3 => { + let s3_config = S3Config::try_from(config)?; + let delete_batch_size = parse_delete_batch_size(config); + tracing::info!( + batch_size = delete_batch_size, + "Initialized S3 storage with delete batch size {} (configure via '{}' to adjust)", + delete_batch_size, + S3_DELETE_BATCH_SIZE + ); + Ok(Arc::new(ObjectStoreStorage::S3(S3Storage { + config: Arc::new(s3_config), + delete_batch_size, + store_cache: Arc::new(DashMap::new()), + }))) + } + } + } +} + +type StoreCache = Arc>>; + +/// `object_store` S3 storage state. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct S3Storage { + config: Arc, + #[serde(default = "default_delete_batch_size")] + delete_batch_size: usize, + #[serde(skip, default)] + store_cache: StoreCache, +} + +/// `object_store`-based storage implementation. +/// +/// Stores are cached per bucket to avoid rebuilding the client on every operation. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub enum ObjectStoreStorage { + /// S3 storage variant. + #[cfg(feature = "object_store-s3")] + S3(S3Storage), +} + +struct StoreAndPath { + bucket: String, + store: Arc, + path: ObjectStorePath, +} + +/// Helper for batching deletions per bucket. +struct BucketBatch { + store: Arc, + locations: Vec, +} + +impl ObjectStoreStorage { + fn delete_batch_size(&self) -> usize { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => s3.delete_batch_size, + } + } + + /// Get or create a cached store and extract the relative `ObjectStorePath`. + fn get_store_and_path(&self, path: &str) -> Result { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => { + let parsed = parse_s3_url(path)?; + + let store = s3 + .store_cache + .entry(parsed.bucket.clone()) + .or_try_insert_with(|| build_s3_store(&s3.config, &parsed.bucket))? + .value() + .clone(); + + let object_path = + ObjectStorePath::from_url_path(&parsed.relative).map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid URL path: {}", parsed.relative), + ) + .with_source(e) + })?; + + Ok(StoreAndPath { + bucket: parsed.bucket, + store, + path: object_path, + }) + } + } + } +} + +#[typetag::serde(name = "ObjectStoreStorage")] +#[async_trait] +impl Storage for ObjectStoreStorage { + async fn exists(&self, path: &str) -> Result { + let target = self.get_store_and_path(path)?; + match target.store.head(&target.path).await { + Ok(_) => Ok(true), + Err(object_store::Error::NotFound { .. }) => Ok(false), + Err(e) => Err(from_object_store_error(e)), + } + } + + async fn metadata(&self, path: &str) -> Result { + let target = self.get_store_and_path(path)?; + let meta = target + .store + .head(&target.path) + .await + .map_err(from_object_store_error)?; + Ok(to_file_metadata(meta)) + } + + async fn read(&self, path: &str) -> Result { + let target = self.get_store_and_path(path)?; + let result = target + .store + .get(&target.path) + .await + .map_err(from_object_store_error)?; + result.bytes().await.map_err(from_object_store_error) + } + + async fn reader(&self, path: &str) -> Result> { + let target = self.get_store_and_path(path)?; + Ok(Box::new(ObjectStoreReader { + store: target.store, + path: target.path, + })) + } + + async fn write(&self, path: &str, bs: Bytes) -> Result<()> { + let target = self.get_store_and_path(path)?; + target + .store + .put(&target.path, PutPayload::from_bytes(bs)) + .await + .map_err(from_object_store_error)?; + Ok(()) + } + + async fn writer(&self, path: &str) -> Result> { + let target = self.get_store_and_path(path)?; + let upload = target + .store + .put_multipart(&target.path) + .await + .map_err(from_object_store_error)?; + Ok(Box::new(ObjectStoreWriter { + upload: Some(upload), + buffer: BytesMut::new(), + tasks: FuturesUnordered::new(), + parts_submitted: 0, + bytes_written: 0, + })) + } + + async fn delete(&self, path: &str) -> Result<()> { + let target = self.get_store_and_path(path)?; + target + .store + .delete(&target.path) + .await + .map_err(from_object_store_error)?; + Ok(()) + } + + async fn delete_prefix(&self, path: &str) -> Result<()> { + let target = self.get_store_and_path(path)?; + let locations = target + .store + .list(Some(&target.path)) + .map_ok(|m| m.location) + .boxed(); + target + .store + .delete_stream(locations) + .try_for_each(|_| async { Ok(()) }) + .await + .map_err(from_object_store_error)?; + Ok(()) + } + + async fn delete_stream(&self, paths: BoxStream<'static, String>) -> Result<()> { + let batch_size = self.delete_batch_size(); + let mut chunk_stream = paths.chunks(batch_size); + + while let Some(chunk) = chunk_stream.next().await { + let mut batches: std::collections::HashMap = + std::collections::HashMap::new(); + + for path in chunk { + let target = self.get_store_and_path(&path)?; + batches + .entry(target.bucket) + .or_insert_with(|| BucketBatch { + store: target.store, + locations: Vec::new(), + }) + .locations + .push(target.path); + } + + for (_bucket, batch) in batches { + let location_stream = + futures::stream::iter(batch.locations.into_iter().map(Ok)).boxed(); + batch + .store + .delete_stream(location_stream) + .try_for_each(|_| async { Ok(()) }) + .await + .map_err(from_object_store_error)?; + } + } + Ok(()) + } + + fn new_input(&self, path: &str) -> Result { + Ok(InputFile::new(Arc::new(self.clone()), path.to_string())) + } + + fn new_output(&self, path: &str) -> Result { + Ok(OutputFile::new(Arc::new(self.clone()), path.to_string())) + } +} + +/// Reader that implements `FileRead` using `object_store`. +struct ObjectStoreReader { + store: Arc, + path: ObjectStorePath, +} + +#[async_trait] +impl FileRead for ObjectStoreReader { + async fn read(&self, range: Range) -> Result { + if range.is_empty() { + return Ok(Bytes::new()); + } + let opts = object_store::GetOptions { + range: Some((range.start..range.end).into()), + ..Default::default() + }; + let result = self + .store + .get_opts(&self.path, opts) + .await + .map_err(from_object_store_error)?; + result.bytes().await.map_err(from_object_store_error) + } +} + +/// Minimum part size for S3 multipart upload (5 MiB). +const MIN_PART_SIZE: usize = 5 * 1024 * 1024; + +/// Default maximum concurrent in-flight part uploads. +const MAX_CONCURRENT_PART_UPLOADS: usize = 8; + +/// Writer that implements `FileWrite` using `object_store` multipart upload. +struct ObjectStoreWriter { + upload: Option>, + buffer: BytesMut, + tasks: FuturesUnordered, + parts_submitted: usize, + bytes_written: u64, +} + +impl ObjectStoreWriter { + /// Best-effort abort on upload error, converting the underlying + /// object_store error to an iceberg Error while logging abort failures. + async fn abort_and_wrap( + upload: &mut Box, + e: object_store::Error, + context: &'static str, + ) -> Error { + if let Err(abort_err) = upload.abort().await { + tracing::warn!( + error = %abort_err, + "Failed to abort multipart upload after {context}" + ); + } + from_object_store_error(e) + } + + /// Flushes any buffered bytes as an in-flight part upload to S3, + /// throttling concurrency so that at most `MAX_CONCURRENT_PART_UPLOADS` + /// uploads are in flight at any given time. + async fn flush_buffer( + buffer: &mut BytesMut, + tasks: &mut FuturesUnordered, + parts_submitted: &mut usize, + upload: &mut Box, + ) -> Result<()> { + if !buffer.is_empty() { + // Drain in-flight tasks down below cap before queuing the next chunk + while tasks.len() >= MAX_CONCURRENT_PART_UPLOADS { + if let Some(Err(e)) = tasks.next().await { + return Err(Self::abort_and_wrap(upload, e, "part upload failure").await); + } + } + + let part_data = std::mem::take(buffer).freeze(); + let part_fut = upload.put_part(PutPayload::from_bytes(part_data)); + tasks.push(part_fut); + *parts_submitted += 1; + } + Ok(()) + } + + /// Accumulates bytes into `buffer`, flushing 5 MiB parts when full. + async fn append_bytes( + buffer: &mut BytesMut, + tasks: &mut FuturesUnordered, + parts_submitted: &mut usize, + mut bs: Bytes, + upload: &mut Box, + ) -> Result<()> { + while !bs.is_empty() { + let remaining = MIN_PART_SIZE.saturating_sub(buffer.len()); + if remaining == 0 { + Self::flush_buffer(buffer, tasks, parts_submitted, upload).await?; + continue; + } + + if bs.len() < remaining { + buffer.extend_from_slice(&bs); + return Ok(()); + } + let chunk = bs.split_to(remaining); + buffer.extend_from_slice(&chunk); + Self::flush_buffer(buffer, tasks, parts_submitted, upload).await?; + } + Ok(()) + } +} + +impl Drop for ObjectStoreWriter { + fn drop(&mut self) { + if let Some(mut upload) = self.upload.take() { + if let Ok(handle) = tokio::runtime::Handle::try_current() { + handle.spawn(async move { + if let Err(e) = upload.abort().await { + tracing::warn!( + error = %e, + "Failed to abort multipart upload on drop" + ); + } + }); + } else { + tracing::warn!( + "ObjectStoreWriter dropped outside a Tokio runtime; multipart upload abort skipped" + ); + } + } + } +} + +#[async_trait] +impl FileWrite for ObjectStoreWriter { + async fn write(&mut self, bs: Bytes) -> Result<()> { + let upload = self.upload.as_mut().ok_or_else(|| { + Error::new( + ErrorKind::PreconditionFailed, + "Writer has already been closed", + ) + })?; + self.bytes_written += bs.len() as u64; + Self::append_bytes( + &mut self.buffer, + &mut self.tasks, + &mut self.parts_submitted, + bs, + upload, + ) + .await + } + + async fn close(&mut self) -> Result { + let mut upload = self.upload.take().ok_or_else(|| { + Error::new( + ErrorKind::PreconditionFailed, + "Writer has already been closed", + ) + })?; + + // Flush any remaining buffered data, or emit an empty part if 0 parts + // have been sent (S3 requires at least 1 part to complete a multipart upload). + if !self.buffer.is_empty() || self.parts_submitted == 0 { + let part_data = std::mem::take(&mut self.buffer).freeze(); + let part_fut = upload.put_part(PutPayload::from_bytes(part_data)); + self.tasks.push(part_fut); + self.parts_submitted += 1; + } + + // Await all in-flight part uploads; abort if any part fails. + while let Some(res) = self.tasks.next().await { + if let Err(e) = res { + return Err(Self::abort_and_wrap(&mut upload, e, "part upload failure").await); + } + } + + // Complete multipart upload; abort if S3 rejects completion. + if let Err(e) = upload.complete().await { + return Err(Self::abort_and_wrap(&mut upload, e, "complete failure").await); + } + + Ok(FileMetadata { + size: self.bytes_written, + }) + } +} + +#[cfg(test)] +mod tests { + use super::*; + + #[cfg(feature = "object_store-s3")] + fn make_s3_storage() -> ObjectStoreStorage { + ObjectStoreStorage::S3(S3Storage { + config: Arc::new(S3Config::default()), + delete_batch_size: DEFAULT_DELETE_BATCH_SIZE, + store_cache: Arc::new(DashMap::new()), + }) + } + + #[test] + fn test_parse_delete_batch_size_default() { + let config = StorageConfig::new(); + assert_eq!(parse_delete_batch_size(&config), DEFAULT_DELETE_BATCH_SIZE); + } + + #[test] + fn test_parse_delete_batch_size_custom() { + let mut props = std::collections::HashMap::new(); + props.insert(S3_DELETE_BATCH_SIZE.to_string(), "500".to_string()); + let config = StorageConfig::from_props(props); + assert_eq!(parse_delete_batch_size(&config), 500); + } + + #[test] + fn test_parse_delete_batch_size_exceeds_max() { + let mut props = std::collections::HashMap::new(); + props.insert(S3_DELETE_BATCH_SIZE.to_string(), "5000".to_string()); + let config = StorageConfig::from_props(props); + assert_eq!(parse_delete_batch_size(&config), S3_MAX_DELETE_BATCH_SIZE); + } + + #[test] + fn test_parse_delete_batch_size_invalid_fallback() { + let mut props = std::collections::HashMap::new(); + props.insert( + S3_DELETE_BATCH_SIZE.to_string(), + "invalid_number".to_string(), + ); + let config = StorageConfig::from_props(props); + assert_eq!(parse_delete_batch_size(&config), DEFAULT_DELETE_BATCH_SIZE); + + let mut props_zero = std::collections::HashMap::new(); + props_zero.insert(S3_DELETE_BATCH_SIZE.to_string(), "0".to_string()); + let config_zero = StorageConfig::from_props(props_zero); + assert_eq!( + parse_delete_batch_size(&config_zero), + DEFAULT_DELETE_BATCH_SIZE + ); + } + + #[cfg(feature = "object_store-s3")] + #[test] + fn test_store_cache_reuses_store() { + let storage = make_s3_storage(); + let target1 = storage + .get_store_and_path("s3://test-bucket/file1.parquet") + .unwrap(); + let target2 = storage + .get_store_and_path("s3://test-bucket/file2.parquet") + .unwrap(); + assert!(Arc::ptr_eq(&target1.store, &target2.store)); + } + + #[cfg(feature = "object_store-s3")] + #[test] + fn test_store_cache_different_buckets() { + let storage = make_s3_storage(); + let target1 = storage + .get_store_and_path("s3://bucket-a/file.parquet") + .unwrap(); + let target2 = storage + .get_store_and_path("s3://bucket-b/file.parquet") + .unwrap(); + assert!(!Arc::ptr_eq(&target1.store, &target2.store)); + } + + #[cfg(feature = "object_store-s3")] + #[test] + fn test_relative_path_extraction() { + let storage = make_s3_storage(); + let target = storage + .get_store_and_path("s3://my-bucket/data/file.parquet") + .unwrap(); + assert_eq!(target.path.as_ref(), "data/file.parquet"); + + let target_encoded = storage + .get_store_and_path("s3://my-bucket/data%20dir/file.parquet") + .unwrap(); + assert_eq!(target_encoded.path.as_ref(), "data dir/file.parquet"); + } + + #[cfg(feature = "object_store-s3")] + #[test] + fn test_storage_serialization_roundtrip() { + let storage = make_s3_storage(); + let serialized = serde_json::to_string(&storage).unwrap(); + let deserialized: ObjectStoreStorage = serde_json::from_str(&serialized).unwrap(); + match deserialized { + ObjectStoreStorage::S3(s3) => { + assert_eq!(s3.config, Arc::new(S3Config::default())); + } + } + } + + #[cfg(feature = "object_store-s3")] + #[test] + fn test_storage_factory_serialization_roundtrip() { + let factory = ObjectStoreStorageFactory::S3; + let serialized = serde_json::to_string(&factory).unwrap(); + let deserialized: ObjectStoreStorageFactory = serde_json::from_str(&serialized).unwrap(); + assert!(matches!(deserialized, ObjectStoreStorageFactory::S3)); + } + + #[cfg(feature = "object_store-s3")] + #[test] + fn test_file_io_serialization_roundtrip() { + use iceberg::io::FileIOBuilder; + let factory = Arc::new(ObjectStoreStorageFactory::S3); + let file_io = FileIOBuilder::new(factory).build(); + let bytes = file_io.serialize_all().unwrap(); + let deserialized = iceberg::io::FileIO::deserialize_all(&bytes).unwrap(); + assert_eq!(file_io.config(), deserialized.config()); + } + + #[tokio::test] + async fn test_reader_empty_range() { + let store = Arc::new(object_store::memory::InMemory::new()); + let path = ObjectStorePath::from("data/test.parquet"); + let reader = ObjectStoreReader { store, path }; + + let result = reader.read(0..0).await.unwrap(); + assert!(result.is_empty()); + + let result = reader.read(10..10).await.unwrap(); + assert!(result.is_empty()); + } + + #[tokio::test] + async fn test_writer_already_closed_errors() { + let mut writer = ObjectStoreWriter { + upload: None, + buffer: BytesMut::new(), + tasks: FuturesUnordered::new(), + parts_submitted: 0, + bytes_written: 0, + }; + let write_err = writer.write(Bytes::from_static(b"data")).await.unwrap_err(); + assert_eq!(write_err.kind(), ErrorKind::PreconditionFailed); + assert_eq!(write_err.message(), "Writer has already been closed"); + + let close_res = writer.close().await; + assert!(close_res.is_err()); + let close_err = close_res.err().unwrap(); + assert_eq!(close_err.kind(), ErrorKind::PreconditionFailed); + assert_eq!(close_err.message(), "Writer has already been closed"); + } + + use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; + + use object_store::PutResult; + + #[derive(Debug)] + struct MockMultipartUpload { + parts_count: Arc, + completed: Arc, + aborted: Arc, + fail_part: bool, + fail_complete: bool, + } + + #[async_trait] + impl MultipartUpload for MockMultipartUpload { + fn put_part(&mut self, _data: PutPayload) -> UploadPart { + let count = self.parts_count.clone(); + let fail = self.fail_part; + Box::pin(async move { + if fail { + Err(object_store::Error::Generic { + store: "mock", + source: "part failed".into(), + }) + } else { + count.fetch_add(1, Ordering::SeqCst); + Ok(()) + } + }) + } + + async fn complete(&mut self) -> object_store::Result { + if self.fail_complete { + Err(object_store::Error::Generic { + store: "mock", + source: "complete failed".into(), + }) + } else { + self.completed.store(true, Ordering::SeqCst); + Ok(PutResult { + e_tag: None, + version: None, + }) + } + } + + async fn abort(&mut self) -> object_store::Result<()> { + self.aborted.store(true, Ordering::SeqCst); + Ok(()) + } + } + + fn make_mock_writer( + fail_part: bool, + fail_complete: bool, + ) -> ( + ObjectStoreWriter, + Arc, + Arc, + Arc, + ) { + let parts_count = Arc::new(AtomicUsize::new(0)); + let completed = Arc::new(AtomicBool::new(false)); + let aborted = Arc::new(AtomicBool::new(false)); + + let upload = Box::new(MockMultipartUpload { + parts_count: parts_count.clone(), + completed: completed.clone(), + aborted: aborted.clone(), + fail_part, + fail_complete, + }); + + let writer = ObjectStoreWriter { + upload: Some(upload), + buffer: BytesMut::new(), + tasks: FuturesUnordered::new(), + parts_submitted: 0, + bytes_written: 0, + }; + + (writer, parts_count, completed, aborted) + } + + #[test] + fn test_writer_drop_outside_tokio_warns_and_does_not_panic() { + // A plain synchronous #[test] runs on an OS thread outside a Tokio runtime context. + // Passing an active upload verifies that Handle::try_current() failing logs a warning without panicking. + let (writer, _parts, completed, aborted) = make_mock_writer(false, false); + drop(writer); + + assert!(!completed.load(Ordering::SeqCst)); + assert!(!aborted.load(Ordering::SeqCst)); + } + + #[tokio::test] + async fn test_writer_close_success() { + let (mut writer, parts, completed, aborted) = make_mock_writer(false, false); + writer + .write(Bytes::from(vec![0u8; 6 * 1024 * 1024])) + .await + .unwrap(); + writer.close().await.unwrap(); + + assert_eq!(parts.load(Ordering::SeqCst), 2); // 5 MiB + 1 MiB leftover part + assert!(completed.load(Ordering::SeqCst)); + assert!(!aborted.load(Ordering::SeqCst)); + } + + #[tokio::test] + async fn test_writer_close_empty_file_emits_one_part() { + let (mut writer, parts, completed, aborted) = make_mock_writer(false, false); + // Closing a fresh writer with 0 bytes must emit exactly 1 part for S3 validity + writer.close().await.unwrap(); + + assert_eq!(parts.load(Ordering::SeqCst), 1); + assert!(completed.load(Ordering::SeqCst)); + assert!(!aborted.load(Ordering::SeqCst)); + } + + #[tokio::test] + async fn test_writer_part_failure_aborts() { + let (mut writer, _parts, completed, aborted) = make_mock_writer(true, false); + writer + .write(Bytes::from(vec![0u8; 6 * 1024 * 1024])) + .await + .unwrap(); + let close_res = writer.close().await; + assert!(close_res.is_err()); + let err = close_res.err().unwrap(); + + assert_eq!(err.kind(), ErrorKind::Unexpected); + assert!(!completed.load(Ordering::SeqCst)); + assert!( + aborted.load(Ordering::SeqCst), + "abort() MUST be called when part upload fails" + ); + } + + #[tokio::test] + async fn test_writer_complete_failure_aborts() { + let (mut writer, _parts, completed, aborted) = make_mock_writer(false, true); + writer.write(Bytes::from_static(b"hello")).await.unwrap(); + let close_res = writer.close().await; + assert!(close_res.is_err()); + let err = close_res.err().unwrap(); + + assert_eq!(err.kind(), ErrorKind::Unexpected); + assert!(!completed.load(Ordering::SeqCst)); + assert!( + aborted.load(Ordering::SeqCst), + "abort() MUST be called when complete() fails" + ); + } + + #[tokio::test] + async fn test_writer_bounded_concurrency_past_limit() { + let (mut writer, parts, completed, aborted) = make_mock_writer(false, false); + // Write 10 parts (50 MiB) which exceeds MAX_CONCURRENT_PART_UPLOADS (8) + writer + .write(Bytes::from(vec![0u8; 10 * 5 * 1024 * 1024])) + .await + .unwrap(); + writer.close().await.unwrap(); + + assert_eq!(parts.load(Ordering::SeqCst), 10); + assert!(completed.load(Ordering::SeqCst)); + assert!(!aborted.load(Ordering::SeqCst)); + } + + #[tokio::test] + async fn test_writer_fail_fast_during_write() { + let (mut writer, _parts, completed, aborted) = make_mock_writer(true, false); + // Writing past MAX_CONCURRENT_PART_UPLOADS triggers in-flight task awaiting inside write() + let err = writer + .write(Bytes::from(vec![0u8; 10 * 5 * 1024 * 1024])) + .await + .unwrap_err(); + + assert_eq!(err.kind(), ErrorKind::Unexpected); + assert!(!completed.load(Ordering::SeqCst)); + assert!( + aborted.load(Ordering::SeqCst), + "abort() MUST be called immediately on write() failure" + ); + } + + #[tokio::test] + async fn test_writer_drop_aborts() { + let (writer, _parts, completed, aborted) = make_mock_writer(false, false); + drop(writer); + + // Allow background tokio task to run + tokio::time::sleep(std::time::Duration::from_millis(50)).await; + assert!(!completed.load(Ordering::SeqCst)); + assert!( + aborted.load(Ordering::SeqCst), + "abort() MUST be called on drop" + ); + } + + #[tokio::test] + async fn test_writer_bounded_concurrency_in_flight_cap() { + let in_flight = Arc::new(AtomicUsize::new(0)); + let max_in_flight = Arc::new(AtomicUsize::new(0)); + + #[derive(Debug)] + struct TrackingMockUpload { + in_flight: Arc, + max_in_flight: Arc, + } + + #[async_trait] + impl MultipartUpload for TrackingMockUpload { + fn put_part(&mut self, _data: PutPayload) -> UploadPart { + let in_flight = self.in_flight.clone(); + let max_in_flight = self.max_in_flight.clone(); + Box::pin(async move { + let current = in_flight.fetch_add(1, Ordering::SeqCst) + 1; + max_in_flight.fetch_max(current, Ordering::SeqCst); + tokio::time::sleep(std::time::Duration::from_millis(5)).await; + in_flight.fetch_sub(1, Ordering::SeqCst); + Ok(()) + }) + } + + async fn complete(&mut self) -> object_store::Result { + Ok(PutResult { + e_tag: None, + version: None, + }) + } + + async fn abort(&mut self) -> object_store::Result<()> { + Ok(()) + } + } + + let upload = Box::new(TrackingMockUpload { + in_flight: in_flight.clone(), + max_in_flight: max_in_flight.clone(), + }); + + let mut writer = ObjectStoreWriter { + upload: Some(upload), + buffer: BytesMut::new(), + tasks: FuturesUnordered::new(), + parts_submitted: 0, + bytes_written: 0, + }; + + // Write 20 parts (100 MiB) in a single write call + writer + .write(Bytes::from(vec![0u8; 20 * 5 * 1024 * 1024])) + .await + .unwrap(); + writer.close().await.unwrap(); + + assert!( + max_in_flight.load(Ordering::SeqCst) <= MAX_CONCURRENT_PART_UPLOADS, + "max concurrent in-flight uploads exceeded {} (was {})", + MAX_CONCURRENT_PART_UPLOADS, + max_in_flight.load(Ordering::SeqCst) + ); + } +} diff --git a/crates/storage/object_store/src/s3.rs b/crates/storage/object_store/src/s3.rs new file mode 100644 index 0000000000..7547ccdaa8 --- /dev/null +++ b/crates/storage/object_store/src/s3.rs @@ -0,0 +1,375 @@ +// 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::str::FromStr; +use std::sync::Arc; + +use iceberg::io::S3Config; +use iceberg::{Error, ErrorKind, Result}; +use object_store::ObjectStore; +use object_store::aws::{AmazonS3Builder, AmazonS3ConfigKey}; +use percent_encoding::percent_decode_str; +use url::Url; + +/// Parsed components of an S3 URL. +#[derive(Debug, PartialEq, Eq)] +pub(crate) struct ParsedS3Url { + pub(crate) scheme: String, + pub(crate) bucket: String, + pub(crate) relative: String, +} + +/// Parse an absolute S3 URL into [`ParsedS3Url`]. +/// +/// Accepts `s3://`, `s3a://`, and `s3n://` schemes. +pub(crate) fn parse_s3_url(path: &str) -> Result { + let url = Url::parse(path).map_err(|e| { + Error::new(ErrorKind::DataInvalid, format!("Invalid URL: {path}")).with_source(e) + })?; + + let scheme = url.scheme(); + match scheme { + "s3" | "s3a" | "s3n" => {} + _ => { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + format!("Unsupported S3 scheme: {scheme} in url: {path}"), + )); + } + } + + let bucket = url.host_str().ok_or_else(|| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid s3 url: {path}, missing bucket"), + ) + })?; + + if bucket.is_empty() { + return Err(Error::new( + ErrorKind::DataInvalid, + format!("Empty s3 url: {path}, missing bucket"), + )); + } + + let bucket = percent_decode_str(bucket).decode_utf8().map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Invalid percent-encoded bucket in s3 url: {path}"), + ) + .with_source(e) + })?; + + let relative = url.path().strip_prefix('/').unwrap_or(url.path()); + + Ok(ParsedS3Url { + scheme: scheme.to_string(), + bucket: bucket.to_string(), + relative: relative.to_string(), + }) +} + +/// Parse a string into an [`AmazonS3ConfigKey`]. +fn parse_s3_config_key(key: &str) -> Result { + AmazonS3ConfigKey::from_str(key).map_err(|e| { + Error::new( + ErrorKind::DataInvalid, + format!("Failed to parse S3 config key: {key}"), + ) + .with_source(e) + }) +} + +/// Configure Server-Side Encryption on `AmazonS3Builder` from `S3Config`. +/// +/// Uses string-based `with_config(parse_s3_config_key(...), ...)` because +/// `S3EncryptionConfigKey` is not re-exported from `object_store::aws` in 0.13.x. +/// The string keys (`"aws_server_side_encryption"`) are stable and used in +/// object_store's own test suite. +fn configure_sse(mut builder: AmazonS3Builder, config: &S3Config) -> Result { + if let Some(ref sse) = config.server_side_encryption { + match sse.as_str() { + "aws:kms" => match &config.server_side_encryption_aws_kms_key_id { + Some(key) => { + builder = builder.with_sse_kms_encryption(key); + } + None => { + builder = builder.with_config( + parse_s3_config_key("aws_server_side_encryption")?, + "aws:kms", + ); + } + }, + "AES256" => { + builder = builder + .with_config(parse_s3_config_key("aws_server_side_encryption")?, "AES256"); + } + other => { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + format!("Unsupported server side encryption type: {other}"), + )); + } + } + } + + if let Some(ref custom_key) = config.server_side_encryption_customer_key { + // Note: object_store's with_ssec_encryption automatically computes and sets + // the x-amz-server-side-encryption-customer-key-MD5 header from the decoded key. + builder = builder.with_ssec_encryption(custom_key); + } + + Ok(builder) +} + +/// Build an `AmazonS3` store from iceberg's `S3Config` for a given bucket. +pub(crate) fn build_s3_store(config: &S3Config, bucket: &str) -> Result> { + if config.role_arn.is_some() { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + "S3 assume-role (role_arn) is not supported by object_store backend", + )); + } + if config.disable_ec2_metadata { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + "S3 disable_ec2_metadata is not supported by object_store backend", + )); + } + if config.disable_config_load { + return Err(Error::new( + ErrorKind::FeatureUnsupported, + "S3 disable_config_load is not supported by object_store backend", + )); + } + + let mut builder = AmazonS3Builder::new().with_bucket_name(bucket); + + if let Some(ref endpoint) = config.endpoint { + builder = builder.with_endpoint(endpoint); + if endpoint.starts_with("http://") { + builder = builder.with_allow_http(true); + } + } + if let Some(ref access_key_id) = config.access_key_id { + builder = builder.with_access_key_id(access_key_id); + } + if let Some(ref secret_access_key) = config.secret_access_key { + builder = builder.with_secret_access_key(secret_access_key); + } + if let Some(ref session_token) = config.session_token { + builder = builder.with_token(session_token); + } + if let Some(ref region) = config.region { + builder = builder.with_region(region); + } + builder = builder.with_virtual_hosted_style_request(config.enable_virtual_host_style); + if config.allow_anonymous { + builder = builder.with_skip_signature(true); + } + + builder = configure_sse(builder, config)?; + + let store = builder.build().map_err(|e| { + Error::new(ErrorKind::Unexpected, "Failed to build S3 object store").with_source(e) + })?; + Ok(Arc::new(store)) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_parse_s3_url() { + let parsed = parse_s3_url("s3://my-bucket/path/to/file.parquet").unwrap(); + assert_eq!(parsed.scheme, "s3"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, "path/to/file.parquet"); + } + + #[test] + fn test_parse_s3a_url() { + let parsed = parse_s3_url("s3a://my-bucket/path/to/file.parquet").unwrap(); + assert_eq!(parsed.scheme, "s3a"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, "path/to/file.parquet"); + } + + #[test] + fn test_parse_s3_url_unsupported_scheme() { + let err = parse_s3_url("gs://my-bucket/file.parquet").unwrap_err(); + assert_eq!(err.kind(), ErrorKind::FeatureUnsupported); + } + + #[test] + fn test_parse_s3_url_bucket_only() { + let parsed = parse_s3_url("s3://my-bucket/").unwrap(); + assert_eq!(parsed.scheme, "s3"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, ""); + } + + #[test] + fn test_parse_s3n_url() { + let parsed = parse_s3_url("s3n://my-bucket/path/to/file.parquet").unwrap(); + assert_eq!(parsed.scheme, "s3n"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, "path/to/file.parquet"); + } + + #[test] + fn test_parse_s3_url_uppercase_scheme() { + let parsed = parse_s3_url("S3://my-bucket/path/to/file.parquet").unwrap(); + assert_eq!(parsed.scheme, "s3"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, "path/to/file.parquet"); + } + + #[test] + fn test_parse_s3_url_percent_encoded_bucket() { + let parsed = parse_s3_url("s3://my%2Dbucket/path/to/file.parquet").unwrap(); + assert_eq!(parsed.scheme, "s3"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, "path/to/file.parquet"); + } + + #[test] + fn test_parse_s3_url_percent_encoded_path() { + let parsed = parse_s3_url("s3://my-bucket/path%20with%20spaces/file.parquet").unwrap(); + assert_eq!(parsed.scheme, "s3"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, "path%20with%20spaces/file.parquet"); + } + + #[test] + fn test_parse_s3_url_double_slash_path() { + let parsed = parse_s3_url("s3://my-bucket//path/to/file.parquet").unwrap(); + assert_eq!(parsed.scheme, "s3"); + assert_eq!(parsed.bucket, "my-bucket"); + assert_eq!(parsed.relative, "/path/to/file.parquet"); + } + + #[test] + fn test_parse_s3_url_empty_bucket() { + assert!(parse_s3_url("s3:///path/to/file.parquet").is_err()); + assert!(parse_s3_url("s3://").is_err()); + } + + #[test] + fn test_parse_s3_config_key() { + assert!(parse_s3_config_key("aws_server_side_encryption").is_ok()); + let err = parse_s3_config_key("invalid_config_key_foo_bar").unwrap_err(); + assert_eq!(err.kind(), ErrorKind::DataInvalid); + } + + #[test] + fn test_configure_sse_kms_default_key() { + let config = S3Config::builder() + .region("us-east-1") + .server_side_encryption("aws:kms") + .build(); + assert!(configure_sse(AmazonS3Builder::new(), &config).is_ok()); + } + + #[test] + fn test_configure_sse_kms_custom_key() { + let config = S3Config::builder() + .region("us-east-1") + .server_side_encryption("aws:kms") + .server_side_encryption_aws_kms_key_id("arn:aws:kms:us-east-1:123456789012:key/test") + .build(); + assert!(configure_sse(AmazonS3Builder::new(), &config).is_ok()); + } + + #[test] + fn test_configure_sse_aes256() { + let config = S3Config::builder() + .region("us-east-1") + .server_side_encryption("AES256") + .build(); + assert!(configure_sse(AmazonS3Builder::new(), &config).is_ok()); + } + + #[test] + fn test_configure_sse_ssec() { + let config = S3Config::builder() + .region("us-east-1") + .server_side_encryption_customer_key("MDEyMzQ1Njc4OTAxMjM0NTY3ODkwMTIzNDU2Nzg5MDE=") + .build(); + assert!(configure_sse(AmazonS3Builder::new(), &config).is_ok()); + } + + #[test] + fn test_configure_sse_unsupported() { + let config = S3Config::builder() + .region("us-east-1") + .server_side_encryption("unsupported_sse") + .build(); + let err = configure_sse(AmazonS3Builder::new(), &config).unwrap_err(); + assert_eq!(err.kind(), ErrorKind::FeatureUnsupported); + } + + #[test] + fn test_build_s3_store_kms_encryption() { + let config = S3Config::builder() + .region("us-east-1") + .server_side_encryption("aws:kms") + .server_side_encryption_aws_kms_key_id("arn:aws:kms:us-east-1:123456789012:key/test") + .build(); + assert!(build_s3_store(&config, "my-bucket").is_ok()); + } + + #[test] + fn test_build_s3_store_ssec_encryption() { + let config = S3Config::builder() + .region("us-east-1") + .server_side_encryption_customer_key("MDEyMzQ1Njc4OTAxMjM0NTY3ODkwMTIzNDU2Nzg5MDE=") + .build(); + assert!(build_s3_store(&config, "my-bucket").is_ok()); + } + + #[test] + fn test_build_s3_store_unsupported_role_arn() { + let config = S3Config::builder() + .region("us-east-1") + .role_arn("arn:aws:iam::123456789012:role/test-role") + .build(); + let err = build_s3_store(&config, "my-bucket").unwrap_err(); + assert_eq!(err.kind(), ErrorKind::FeatureUnsupported); + } + + #[test] + fn test_build_s3_store_unsupported_disable_ec2_metadata() { + let config = S3Config::builder() + .region("us-east-1") + .disable_ec2_metadata(true) + .build(); + let err = build_s3_store(&config, "my-bucket").unwrap_err(); + assert_eq!(err.kind(), ErrorKind::FeatureUnsupported); + } + + #[test] + fn test_build_s3_store_unsupported_disable_config_load() { + let config = S3Config::builder() + .region("us-east-1") + .disable_config_load(true) + .build(); + let err = build_s3_store(&config, "my-bucket").unwrap_err(); + assert_eq!(err.kind(), ErrorKind::FeatureUnsupported); + } +} diff --git a/crates/storage/object_store/tests/file_io_s3_test.rs b/crates/storage/object_store/tests/file_io_s3_test.rs new file mode 100644 index 0000000000..3b90669143 --- /dev/null +++ b/crates/storage/object_store/tests/file_io_s3_test.rs @@ -0,0 +1,511 @@ +// 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. + +//! Integration tests for FileIO S3 using object_store backend. +//! +//! These tests assume Docker containers are started externally via `make docker-up`. +//! Each test uses unique file paths based on module path to avoid conflicts. + +#[cfg(feature = "object_store-s3")] +mod tests { + use std::sync::Arc; + + use bytes::Bytes; + use futures::StreamExt; + use iceberg::io::{ + FileIO, FileIOBuilder, S3_ACCESS_KEY_ID, S3_ENDPOINT, S3_PATH_STYLE_ACCESS, S3_REGION, + S3_SECRET_ACCESS_KEY, S3_SSE_KEY, S3_SSE_TYPE, + }; + use iceberg_storage_object_store::ObjectStoreStorageFactory; + use iceberg_test_utils::{get_object_store_endpoint, normalize_test_name_with_parts, set_up}; + + async fn get_file_io() -> FileIO { + set_up(); + + let endpoint = get_object_store_endpoint(); + + FileIOBuilder::new(Arc::new(ObjectStoreStorageFactory::S3)) + .with_props(vec![ + (S3_ENDPOINT, endpoint), + (S3_ACCESS_KEY_ID, "admin".to_string()), + (S3_SECRET_ACCESS_KEY, "password".to_string()), + (S3_REGION, "us-east-1".to_string()), + (S3_PATH_STYLE_ACCESS, "true".to_string()), + ]) + .build() + } + + fn roundtrip_file_io(file_io: &FileIO) -> FileIO { + let serialized = file_io.serialize_all().unwrap(); + FileIO::deserialize_all(&serialized).unwrap() + } + + #[tokio::test] + async fn test_file_io_s3_serialization_roundtrip() { + let file_io = roundtrip_file_io(&get_file_io().await); + let path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_serialization_roundtrip") + ); + + let _ = file_io.delete(&path).await; + file_io + .new_output(&path) + .unwrap() + .write(Bytes::from_static(b"roundtrip")) + .await + .unwrap(); + assert_eq!( + file_io.new_input(&path).unwrap().read().await.unwrap(), + Bytes::from_static(b"roundtrip") + ); + file_io.delete(&path).await.unwrap(); + assert!(!file_io.exists(&path).await.unwrap()); + } + + #[tokio::test] + async fn test_file_io_s3_exists() { + let file_io = get_file_io().await; + assert!(!file_io.exists("s3://bucket2/any").await.unwrap()); + assert!(file_io.exists("s3://bucket1/").await.unwrap()); + } + + #[tokio::test] + async fn test_file_io_s3_output() { + let file_io = get_file_io().await; + let output_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_output") + ); + let _ = file_io.delete(&output_path).await; + assert!(!file_io.exists(&output_path).await.unwrap()); + let output_file = file_io.new_output(&output_path).unwrap(); + { + output_file.write("123".into()).await.unwrap(); + } + assert!(file_io.exists(&output_path).await.unwrap()); + } + + #[tokio::test] + async fn test_file_io_s3_input() { + let file_io = get_file_io().await; + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_input") + ); + let output_file = file_io.new_output(&file_path).unwrap(); + { + output_file.write("test_input".into()).await.unwrap(); + } + + let input_file = file_io.new_input(&file_path).unwrap(); + { + let buffer = input_file.read().await.unwrap(); + assert_eq!(buffer, "test_input".as_bytes()); + } + } + + #[tokio::test] + async fn test_file_io_s3_delete_stream() { + let file_io = get_file_io().await; + + let paths: Vec = (0..5) + .map(|i| { + format!( + "s3://bucket1/{}/file-{i}", + normalize_test_name_with_parts!("test_file_io_s3_delete_stream") + ) + }) + .collect(); + for path in &paths { + let _ = file_io.delete(path).await; + file_io + .new_output(path) + .unwrap() + .write("delete-me".into()) + .await + .unwrap(); + assert!(file_io.exists(path).await.unwrap()); + } + + let stream = futures::stream::iter(paths.clone()).boxed(); + file_io.delete_stream(stream).await.unwrap(); + + for path in &paths { + assert!(!file_io.exists(path).await.unwrap()); + } + } + + #[tokio::test] + async fn test_file_io_s3_delete_stream_empty() { + let file_io = get_file_io().await; + let stream = futures::stream::empty().boxed(); + file_io.delete_stream(stream).await.unwrap(); + } + + #[tokio::test] + async fn test_file_io_s3_delete_stream_invalid_url() { + let file_io = get_file_io().await; + let stream = futures::stream::iter(vec!["invalid-url".to_string()]).boxed(); + let res = file_io.delete_stream(stream).await; + assert!(res.is_err()); + } + + #[tokio::test] + async fn test_file_io_s3_multipart_writer() { + let file_io = get_file_io().await; + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_multipart_writer") + ); + let _ = file_io.delete(&file_path).await; + + let output_file = file_io.new_output(&file_path).unwrap(); + let mut writer = output_file.writer().await.unwrap(); + + let chunk1 = Bytes::from_static(b"hello "); + let chunk2 = Bytes::from_static(b"multipart "); + let chunk3 = Bytes::from_static(b"world!"); + + writer.write(chunk1).await.unwrap(); + writer.write(chunk2).await.unwrap(); + writer.write(chunk3).await.unwrap(); + writer.close().await.unwrap(); + + assert!(file_io.exists(&file_path).await.unwrap()); + let input_file = file_io.new_input(&file_path).unwrap(); + let content = input_file.read().await.unwrap(); + assert_eq!(content, Bytes::from_static(b"hello multipart world!")); + + file_io.delete(&file_path).await.unwrap(); + assert!(!file_io.exists(&file_path).await.unwrap()); + } + + #[tokio::test] + async fn test_file_io_s3_percent_encoded_bucket() { + let file_io = get_file_io().await; + let file_path = format!( + "s3://bucket%31/{}", + normalize_test_name_with_parts!("test_file_io_s3_percent_encoded_bucket") + ); + let canonical_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_percent_encoded_bucket") + ); + + let _ = file_io.delete(&file_path).await; + file_io + .new_output(&file_path) + .unwrap() + .write(Bytes::from_static(b"encoded-bucket-content")) + .await + .unwrap(); + + assert!(file_io.exists(&file_path).await.unwrap()); + assert!(file_io.exists(&canonical_path).await.unwrap()); + + let content = file_io + .new_input(&canonical_path) + .unwrap() + .read() + .await + .unwrap(); + assert_eq!(content, Bytes::from_static(b"encoded-bucket-content")); + + file_io.delete(&canonical_path).await.unwrap(); + assert!(!file_io.exists(&file_path).await.unwrap()); + } + + fn assert_minio_kms_rejection(err: &iceberg::Error) { + let err_str = err.to_string(); + assert!( + err_str.contains("501") + || err_str.contains("NotImplemented") + || err_str.contains("KMS is not configured") + || err_str.contains("kms") + || err_str.contains("ServerNotInitialized"), + "Expected KMS rejection from KES-less MinIO confirming SSE header was sent, got: {err_str}" + ); + } + + #[tokio::test] + async fn test_file_io_s3_sse_kms_default() { + set_up(); + let endpoint = get_object_store_endpoint(); + + let file_io = FileIOBuilder::new(Arc::new(ObjectStoreStorageFactory::S3)) + .with_props(vec![ + (S3_ENDPOINT, endpoint), + (S3_ACCESS_KEY_ID, "admin".to_string()), + (S3_SECRET_ACCESS_KEY, "password".to_string()), + (S3_REGION, "us-east-1".to_string()), + (S3_PATH_STYLE_ACCESS, "true".to_string()), + (S3_SSE_TYPE, "kms".to_string()), + ]) + .build(); + + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_sse_kms_default") + ); + + let _ = file_io.delete(&file_path).await; + // MinIO without KES does not configure KMS; asserting that MinIO returns 501 / KMS error + // verifies that the SSE-KMS header was attached (an unencrypted PUT would succeed). + let err = file_io + .new_output(&file_path) + .unwrap() + .write(Bytes::from_static(b"kms-encrypted-data")) + .await + .expect_err( + "MinIO without KES must reject SSE-KMS requests, proving SSE header was sent", + ); + assert_minio_kms_rejection(&err); + } + + #[tokio::test] + async fn test_file_io_s3_sse_kms_custom_key() { + set_up(); + let endpoint = get_object_store_endpoint(); + + let file_io = FileIOBuilder::new(Arc::new(ObjectStoreStorageFactory::S3)) + .with_props(vec![ + (S3_ENDPOINT, endpoint), + (S3_ACCESS_KEY_ID, "admin".to_string()), + (S3_SECRET_ACCESS_KEY, "password".to_string()), + (S3_REGION, "us-east-1".to_string()), + (S3_PATH_STYLE_ACCESS, "true".to_string()), + (S3_SSE_TYPE, "kms".to_string()), + ( + S3_SSE_KEY, + "arn:aws:kms:us-east-1:000000000000:key/a4644f9c-2149-414e-b6a6-a8e82cd6b69e" + .to_string(), + ), + ]) + .build(); + + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_sse_kms_custom_key") + ); + + let _ = file_io.delete(&file_path).await; + // MinIO without KES does not configure KMS; asserting that MinIO returns 501 / KMS error + // verifies that the SSE-KMS custom key header was attached. + let err = file_io + .new_output(&file_path) + .unwrap() + .write(Bytes::from_static(b"kms-custom-key-encrypted-data")) + .await + .expect_err( + "MinIO without KES must reject SSE-KMS requests, proving SSE custom key header was sent", + ); + assert_minio_kms_rejection(&err); + } + + /// Writes 12 MiB (3 × 4 MiB chunks) to exercise real S3 multipart uploads + /// past the 10 MiB `WriteMultipart` buffer threshold, then reads back and + /// verifies byte-for-byte integrity. + #[tokio::test] + async fn test_file_io_s3_multipart_writer_past_threshold() { + let file_io = get_file_io().await; + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_multipart_writer_past_threshold") + ); + let _ = file_io.delete(&file_path).await; + + // 4 MiB chunk with deterministic pattern (repeating 0..=255) + const CHUNK_SIZE: usize = 4 * 1024 * 1024; + let pattern: Vec = (0..CHUNK_SIZE).map(|i| (i % 256) as u8).collect(); + let chunk = Bytes::from(pattern.clone()); + + let output_file = file_io.new_output(&file_path).unwrap(); + let mut writer = output_file.writer().await.unwrap(); + + // Write 3 chunks = 12 MiB total (past 10 MiB threshold) + for _ in 0..3 { + writer.write(chunk.clone()).await.unwrap(); + } + writer.close().await.unwrap(); + + // Read back and verify + let content = file_io.new_input(&file_path).unwrap().read().await.unwrap(); + assert_eq!(content.len(), 3 * CHUNK_SIZE); + for i in 0..3 { + assert_eq!( + &content[i * CHUNK_SIZE..(i + 1) * CHUNK_SIZE], + &pattern[..], + "chunk {i} mismatch" + ); + } + + file_io.delete(&file_path).await.unwrap(); + } + + /// Creates 15 files under `test_prefix/` and 2 under `other_prefix/`, + /// calls `delete_prefix` on `test_prefix/`, then asserts all 15 are gone + /// and the 2 outside the prefix are untouched. + #[tokio::test] + async fn test_file_io_s3_delete_prefix_bulk() { + let file_io = get_file_io().await; + let base = normalize_test_name_with_parts!("test_file_io_s3_delete_prefix_bulk"); + let target_prefix = format!("s3://bucket1/{base}/test_prefix"); + let other_prefix = format!("s3://bucket1/{base}/other_prefix"); + + // Create 15 files under test_prefix/ + for i in 0..15 { + let path = format!("{target_prefix}/file_{i}"); + let _ = file_io.delete(&path).await; + file_io + .new_output(&path) + .unwrap() + .write(Bytes::from(format!("data-{i}"))) + .await + .unwrap(); + } + + // Create 2 files under other_prefix/ + let keep_paths: Vec = (0..2).map(|i| format!("{other_prefix}/keep_{i}")).collect(); + for path in &keep_paths { + let _ = file_io.delete(path).await; + file_io + .new_output(path) + .unwrap() + .write(Bytes::from_static(b"keep-me")) + .await + .unwrap(); + } + + // Bulk delete under test_prefix/ + file_io.delete_prefix(&target_prefix).await.unwrap(); + + // Assert all 15 test_prefix files are gone + for i in 0..15 { + let path = format!("{target_prefix}/file_{i}"); + assert!( + !file_io.exists(&path).await.unwrap(), + "file_{i} should be deleted" + ); + } + + // Assert the 2 other_prefix files still exist + for path in &keep_paths { + assert!( + file_io.exists(path).await.unwrap(), + "{path} should still exist" + ); + } + + // Clean up + for path in &keep_paths { + file_io.delete(path).await.unwrap(); + } + } + + /// Writes a 1024-byte payload and verifies range reads return exact slices. + #[tokio::test] + async fn test_file_io_s3_range_reader() { + let file_io = get_file_io().await; + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_range_reader") + ); + let _ = file_io.delete(&file_path).await; + + // 1024 bytes: 0..=255 repeated 4 times + let payload: Vec = (0..1024).map(|i| (i % 256) as u8).collect(); + file_io + .new_output(&file_path) + .unwrap() + .write(Bytes::from(payload.clone())) + .await + .unwrap(); + + let reader = file_io + .new_input(&file_path) + .unwrap() + .reader() + .await + .unwrap(); + + // Test various ranges + let r1 = reader.read(0..10).await.unwrap(); + assert_eq!(&r1[..], &payload[0..10]); + + let r2 = reader.read(10..50).await.unwrap(); + assert_eq!(&r2[..], &payload[10..50]); + + let r3 = reader.read(100..200).await.unwrap(); + assert_eq!(&r3[..], &payload[100..200]); + + // Cross the 256-byte pattern boundary + let r4 = reader.read(250..260).await.unwrap(); + assert_eq!(&r4[..], &payload[250..260]); + + // Last 10 bytes + let r5 = reader.read(1014..1024).await.unwrap(); + assert_eq!(&r5[..], &payload[1014..1024]); + + file_io.delete(&file_path).await.unwrap(); + } + + #[tokio::test] + async fn test_file_io_s3_sse_s3_aes256() { + set_up(); + let endpoint = get_object_store_endpoint(); + + let file_io = FileIOBuilder::new(Arc::new(ObjectStoreStorageFactory::S3)) + .with_props(vec![ + (S3_ENDPOINT, endpoint), + (S3_ACCESS_KEY_ID, "admin".to_string()), + (S3_SECRET_ACCESS_KEY, "password".to_string()), + (S3_REGION, "us-east-1".to_string()), + (S3_PATH_STYLE_ACCESS, "true".to_string()), + (S3_SSE_TYPE, "s3".to_string()), + ]) + .build(); + + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_sse_s3_aes256") + ); + + let _ = file_io.delete(&file_path).await; + match file_io + .new_output(&file_path) + .unwrap() + .write(Bytes::from_static(b"aes256-encrypted-data")) + .await + { + Ok(_) => { + assert!(file_io.exists(&file_path).await.unwrap()); + let content = file_io.new_input(&file_path).unwrap().read().await.unwrap(); + assert_eq!(content, Bytes::from_static(b"aes256-encrypted-data")); + file_io.delete(&file_path).await.unwrap(); + } + Err(e) + if e.to_string().contains("501") + || e.to_string().contains("NotImplemented") + || e.to_string().contains("KMS is not configured") => + { + // MinIO without KES does not configure server-side encryption; passing 501 verifies header was sent + } + Err(e) => panic!("Unexpected error: {e:?}"), + } + } +}