laskoviymishka commented on code in PR #3165: URL: https://github.com/apache/iceberg-rust/pull/3165#discussion_r4187755404
########## crates/storage/object_store/src/lib.rs: ########## @@ -0,0 +1,1111 @@ +// 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; +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}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinSet; + +/// 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::<usize>() { + 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<Arc<dyn Storage>> { + 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<DashMap<String, Arc<dyn ObjectStore>>>; + +/// `object_store` S3 storage state. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct S3Storage { + config: Arc<S3Config>, + #[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<dyn ObjectStore>, + path: ObjectStorePath, +} + +/// Helper for batching deletions per bucket. +struct BucketBatch { + store: Arc<dyn ObjectStore>, + locations: Vec<ObjectStorePath>, +} + +impl ObjectStoreStorage { + fn delete_batch_size(&self) -> usize { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => s3.delete_batch_size.clamp(1, S3_MAX_DELETE_BATCH_SIZE), + } + } + + /// Get or create a cached store and extract the relative `ObjectStorePath`. + fn get_store_and_path(&self, path: &str) -> Result<StoreAndPath> { + 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(parsed.relative.as_str()); Review Comment: Build the path with `ObjectStorePath::parse(parsed.relative.as_str())` instead of `Path::from` — `from` percent-encodes each segment and `%` is in its INVALID set, so `dt=a%2Fb` becomes `dt=a%252Fb`. object_store's S3 client then encodes once more on the wire, so we store and look up `dt=a%252Fb` where Java and PyIceberg use `dt=a%2Fb`. Any table written elsewhere with a `%XX` in a partition value or filename 404s through this backend, and files we write here land at a key the manifest doesn't record — position deletes and `delete_prefix` miss too. The unit test at line 682 pins the wrong value (`dt=a%252Fb/file.parquet`), and the integration test only round-trips through the same client, so the double-encode cancels and the bug stays hidden — worth asserting the raw key via a second client or a bucket listing. `parse` treats the input as already-encoded, and it also errors on the empty segment in `s3://b//x` instead of silently collapsing it to `x` the way `from` does today. ########## crates/storage/object_store/src/lib.rs: ########## @@ -0,0 +1,1111 @@ +// 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; +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}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinSet; + +/// 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::<usize>() { + 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<Arc<dyn Storage>> { + 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<DashMap<String, Arc<dyn ObjectStore>>>; + +/// `object_store` S3 storage state. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct S3Storage { + config: Arc<S3Config>, + #[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 { Review Comment: I'd add `#[non_exhaustive]` to this enum and to `ObjectStoreStorageFactory` (line 120) before the first publish, and regenerate the snapshot. Once GCS/Azure variants land, anyone matching on these exhaustively breaks — and since this is the base crate the other backends snapshot against, it's cheap now and a semver bump later. ########## crates/storage/object_store/tests/file_io_s3_test.rs: ########## @@ -0,0 +1,563 @@ +// 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; + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_exists") + ); + + let _ = file_io.delete(&file_path).await; + assert!(!file_io.exists(&file_path).await.unwrap()); + + let output_file = file_io.new_output(&file_path).unwrap(); + output_file + .write(Bytes::from_static(b"test_exists")) + .await + .unwrap(); + assert!(file_io.exists(&file_path).await.unwrap()); + + file_io.delete(&file_path).await.unwrap(); + assert!(!file_io.exists(&file_path).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<String> = (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()); + } + + #[tokio::test] + async fn test_file_io_s3_encoded_partition_path() { + 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()), + ]) + .build(); + + let file_path = format!( + "s3://bucket1/{}/dt=a%2Fb/data.parquet", + normalize_test_name_with_parts!("test_file_io_s3_encoded_partition_path") + ); + + let _ = file_io.delete(&file_path).await; + file_io + .new_output(&file_path) + .unwrap() + .write(Bytes::from_static(b"partition-data")) + .await + .unwrap(); + + 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"partition-data")); + + file_io.delete(&file_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<u8> = (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<String> = (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<u8> = (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(_) => { Review Comment: Same gap the KMS tests had before round 5: the `Ok(_)` arm read-backs and passes, so if the SSE header were dropped the PUT would just succeed on KES-less MinIO and this stays green — the inline comment about 501 proving the header was sent contradicts the branch that accepts success. Pick one deterministic expectation for CI MinIO and assert it the way the KMS tests do, or record the request headers with a mock server. ########## crates/storage/object_store/src/lib.rs: ########## @@ -0,0 +1,1111 @@ +// 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; +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}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinSet; + +/// 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"; Review Comment: These three are `pub` from a backend crate, and `S3_DELETE_BATCH_SIZE` is a property key — those live in `iceberg::io` so both backends share one definition. I'd make them `pub(crate)` or move the key up before publish, and regenerate public-api.txt either way. ########## crates/storage/object_store/tests/file_io_s3_test.rs: ########## @@ -0,0 +1,563 @@ +// 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; + let file_path = format!( + "s3://bucket1/{}", + normalize_test_name_with_parts!("test_file_io_s3_exists") + ); + + let _ = file_io.delete(&file_path).await; + assert!(!file_io.exists(&file_path).await.unwrap()); + + let output_file = file_io.new_output(&file_path).unwrap(); + output_file + .write(Bytes::from_static(b"test_exists")) + .await + .unwrap(); + assert!(file_io.exists(&file_path).await.unwrap()); + + file_io.delete(&file_path).await.unwrap(); + assert!(!file_io.exists(&file_path).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<String> = (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()); + } + + #[tokio::test] + async fn test_file_io_s3_encoded_partition_path() { + 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()), + ]) + .build(); + + let file_path = format!( + "s3://bucket1/{}/dt=a%2Fb/data.parquet", + normalize_test_name_with_parts!("test_file_io_s3_encoded_partition_path") + ); + + let _ = file_io.delete(&file_path).await; + file_io + .new_output(&file_path) + .unwrap() + .write(Bytes::from_static(b"partition-data")) + .await + .unwrap(); + + 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"partition-data")); + + file_io.delete(&file_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") Review Comment: I credited the `expect_err` tightening as closing this last round, but the matcher's too loose to prove anything. `assert_minio_kms_rejection` passes on any error whose string contains "kms", and the `iceberg::Error` display includes the request URL — which carries the test key (`..._sse_kms_default`). So a connection refusal, 403, or missing bucket all stringify with "kms" and pass without the SSE header ever being sent. Drop the bare "kms" arm and match only MinIO's KES-less signals (501 / NotImplemented / "KMS is not configured" / ServerNotInitialized), or downcast and check the status. ########## crates/storage/object_store/src/lib.rs: ########## @@ -0,0 +1,1111 @@ +// 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; +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}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinSet; + +/// 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::<usize>() { + 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<Arc<dyn Storage>> { + 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<DashMap<String, Arc<dyn ObjectStore>>>; + +/// `object_store` S3 storage state. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct S3Storage { + config: Arc<S3Config>, + #[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<dyn ObjectStore>, + path: ObjectStorePath, +} + +/// Helper for batching deletions per bucket. +struct BucketBatch { + store: Arc<dyn ObjectStore>, + locations: Vec<ObjectStorePath>, +} + +impl ObjectStoreStorage { + fn delete_batch_size(&self) -> usize { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => s3.delete_batch_size.clamp(1, S3_MAX_DELETE_BATCH_SIZE), + } + } + + /// Get or create a cached store and extract the relative `ObjectStorePath`. + fn get_store_and_path(&self, path: &str) -> Result<StoreAndPath> { + 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(parsed.relative.as_str()); + + 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<bool> { + 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<FileMetadata> { + 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<Bytes> { + 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<Box<dyn FileRead>> { + 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<Box<dyn FileWrite>> { + 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: JoinSet::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<String, BucketBatch> = + 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<InputFile> { + Ok(InputFile::new(Arc::new(self.clone()), path.to_string())) + } + + fn new_output(&self, path: &str) -> Result<OutputFile> { + Ok(OutputFile::new(Arc::new(self.clone()), path.to_string())) + } +} + +/// Reader that implements `FileRead` using `object_store`. +struct ObjectStoreReader { + store: Arc<dyn ObjectStore>, + path: ObjectStorePath, +} + +#[async_trait] +impl FileRead for ObjectStoreReader { + async fn read(&self, range: Range<u64>) -> Result<Bytes> { + 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<Box<dyn MultipartUpload>>, + buffer: BytesMut, + tasks: JoinSet<object_store::Result<()>>, + 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<dyn MultipartUpload>, + 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) + } + + /// Drains in-flight tasks down below `MAX_CONCURRENT_PART_UPLOADS`. + async fn drain_capacity( + tasks: &mut JoinSet<object_store::Result<()>>, + upload: &mut Box<dyn MultipartUpload>, + ) -> Result<()> { + while tasks.len() >= MAX_CONCURRENT_PART_UPLOADS { + if let Some(res) = tasks.join_next().await { + match res { + Ok(Err(e)) => { + return Err(Self::abort_and_wrap(upload, e, "part upload failure").await); + } + Err(join_err) => { + if let Err(abort_err) = upload.abort().await { + tracing::warn!( + error = %abort_err, + "Failed to abort multipart upload after task panic" + ); + } + return Err( + Error::new(ErrorKind::Unexpected, "Part upload task panicked") + .with_source(join_err), + ); + } + Ok(Ok(())) => {} + } + } + } + Ok(()) + } + + /// 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 JoinSet<object_store::Result<()>>, + parts_submitted: &mut usize, + upload: &mut Box<dyn MultipartUpload>, + ) -> Result<()> { + if !buffer.is_empty() { + Self::drain_capacity(tasks, upload).await?; + + let part_data = std::mem::take(buffer).freeze(); + let part_fut = upload.put_part(PutPayload::from_bytes(part_data)); + tasks.spawn(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 JoinSet<object_store::Result<()>>, + parts_submitted: &mut usize, + mut bs: Bytes, + upload: &mut Box<dyn MultipartUpload>, + ) -> 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; + if let Err(e) = Self::append_bytes( + &mut self.buffer, + &mut self.tasks, + &mut self.parts_submitted, + bs, + upload, + ) + .await + { + let _ = self.upload.take(); + return Err(e); + } + Ok(()) + } + + async fn close(&mut self) -> Result<FileMetadata> { + let mut upload = self.upload.take().ok_or_else(|| { Review Comment: `close()` moves `upload` out of `self.upload` into a local, so if the `close()` future is dropped mid-await (a `select!`/timeout upstream) the local drops with no abort and `self.upload` is already `None` — Drop sees nothing and the upload leaks. `write()` doesn't have this hole because it borrows. I'd keep it in `self.upload` and borrow with `as_mut()` until `complete()` succeeds, or at least document that `close()` isn't cancel-safe. ########## crates/storage/object_store/src/lib.rs: ########## @@ -0,0 +1,1111 @@ +// 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; +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}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinSet; + +/// 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::<usize>() { + 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<Arc<dyn Storage>> { + 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<DashMap<String, Arc<dyn ObjectStore>>>; + +/// `object_store` S3 storage state. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct S3Storage { + config: Arc<S3Config>, + #[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<dyn ObjectStore>, + path: ObjectStorePath, +} + +/// Helper for batching deletions per bucket. +struct BucketBatch { + store: Arc<dyn ObjectStore>, + locations: Vec<ObjectStorePath>, +} + +impl ObjectStoreStorage { + fn delete_batch_size(&self) -> usize { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => s3.delete_batch_size.clamp(1, S3_MAX_DELETE_BATCH_SIZE), + } + } + + /// Get or create a cached store and extract the relative `ObjectStorePath`. + fn get_store_and_path(&self, path: &str) -> Result<StoreAndPath> { + 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(parsed.relative.as_str()); + + 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<bool> { + 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<FileMetadata> { + 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<Bytes> { + 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<Box<dyn FileRead>> { + 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<Box<dyn FileWrite>> { + 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: JoinSet::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<String, BucketBatch> = + 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<InputFile> { + Ok(InputFile::new(Arc::new(self.clone()), path.to_string())) + } + + fn new_output(&self, path: &str) -> Result<OutputFile> { + Ok(OutputFile::new(Arc::new(self.clone()), path.to_string())) + } +} + +/// Reader that implements `FileRead` using `object_store`. +struct ObjectStoreReader { + store: Arc<dyn ObjectStore>, + path: ObjectStorePath, +} + +#[async_trait] +impl FileRead for ObjectStoreReader { + async fn read(&self, range: Range<u64>) -> Result<Bytes> { + 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<Box<dyn MultipartUpload>>, + buffer: BytesMut, + tasks: JoinSet<object_store::Result<()>>, + 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<dyn MultipartUpload>, + 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) + } + + /// Drains in-flight tasks down below `MAX_CONCURRENT_PART_UPLOADS`. + async fn drain_capacity( + tasks: &mut JoinSet<object_store::Result<()>>, + upload: &mut Box<dyn MultipartUpload>, + ) -> Result<()> { + while tasks.len() >= MAX_CONCURRENT_PART_UPLOADS { + if let Some(res) = tasks.join_next().await { + match res { + Ok(Err(e)) => { + return Err(Self::abort_and_wrap(upload, e, "part upload failure").await); + } + Err(join_err) => { + if let Err(abort_err) = upload.abort().await { + tracing::warn!( + error = %abort_err, + "Failed to abort multipart upload after task panic" + ); + } + return Err( + Error::new(ErrorKind::Unexpected, "Part upload task panicked") + .with_source(join_err), + ); + } + Ok(Ok(())) => {} + } + } + } + Ok(()) + } + + /// 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 JoinSet<object_store::Result<()>>, + parts_submitted: &mut usize, + upload: &mut Box<dyn MultipartUpload>, + ) -> Result<()> { + if !buffer.is_empty() { + Self::drain_capacity(tasks, upload).await?; + + let part_data = std::mem::take(buffer).freeze(); + let part_fut = upload.put_part(PutPayload::from_bytes(part_data)); + tasks.spawn(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 JoinSet<object_store::Result<()>>, + parts_submitted: &mut usize, + mut bs: Bytes, + upload: &mut Box<dyn MultipartUpload>, + ) -> 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; + if let Err(e) = Self::append_bytes( + &mut self.buffer, + &mut self.tasks, + &mut self.parts_submitted, + bs, + upload, + ) + .await + { + let _ = self.upload.take(); Review Comment: The `take()` here is what stops Drop from firing a second `abort()` — but nothing pins it. The mock's `aborted` is a bool, so a second abort is invisible, and if this line were removed every write-path test would still pass. I'd make the mock count aborts (`Arc<AtomicUsize>`), then `drop(writer); yield_now().await;` and assert the count is 1 in the fail-fast and panic tests. ########## crates/storage/object_store/src/lib.rs: ########## @@ -0,0 +1,1111 @@ +// 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; +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}; +#[cfg(feature = "object_store-s3")] +use s3::{build_s3_store, parse_s3_url}; +use serde::{Deserialize, Serialize}; +use tokio::task::JoinSet; + +/// 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::<usize>() { + 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<Arc<dyn Storage>> { + 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<DashMap<String, Arc<dyn ObjectStore>>>; + +/// `object_store` S3 storage state. +#[derive(Clone, Debug, Serialize, Deserialize)] +pub struct S3Storage { + config: Arc<S3Config>, + #[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<dyn ObjectStore>, + path: ObjectStorePath, +} + +/// Helper for batching deletions per bucket. +struct BucketBatch { + store: Arc<dyn ObjectStore>, + locations: Vec<ObjectStorePath>, +} + +impl ObjectStoreStorage { + fn delete_batch_size(&self) -> usize { + match self { + #[cfg(feature = "object_store-s3")] + ObjectStoreStorage::S3(s3) => s3.delete_batch_size.clamp(1, S3_MAX_DELETE_BATCH_SIZE), + } + } + + /// Get or create a cached store and extract the relative `ObjectStorePath`. + fn get_store_and_path(&self, path: &str) -> Result<StoreAndPath> { + 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(parsed.relative.as_str()); + + 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<bool> { + 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<FileMetadata> { + 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<Bytes> { + 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<Box<dyn FileRead>> { + 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<Box<dyn FileWrite>> { + 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: JoinSet::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<String, BucketBatch> = + 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<InputFile> { + Ok(InputFile::new(Arc::new(self.clone()), path.to_string())) + } + + fn new_output(&self, path: &str) -> Result<OutputFile> { + Ok(OutputFile::new(Arc::new(self.clone()), path.to_string())) + } +} + +/// Reader that implements `FileRead` using `object_store`. +struct ObjectStoreReader { + store: Arc<dyn ObjectStore>, + path: ObjectStorePath, +} + +#[async_trait] +impl FileRead for ObjectStoreReader { + async fn read(&self, range: Range<u64>) -> Result<Bytes> { + 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; Review Comment: The fixed 5 MiB part size caps a single object at ~48.8 GiB (10,000 parts), with no growth, so a large compaction output fails at part 10,001 with an opaque error. Java defaults 32 MiB and makes it configurable (`s3.multipart.part-size-bytes`). Growing the part size as the count climbs is the real fix; at minimum I'd document the ceiling. Not urgent — flagging so it's a known limit, not a surprise. ########## crates/storage/object_store/src/s3.rs: ########## @@ -0,0 +1,384 @@ +// 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<ParsedS3Url> { + 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 = match path.find("://") { + Some(scheme_end) => { + let after_scheme = &path[scheme_end + 3..]; + match after_scheme.find('/') { + Some(idx) => after_scheme[idx..].strip_prefix('/').unwrap_or(""), + None => "", + } + } + None => 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> { + 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<AmazonS3Builder> { + 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<Arc<dyn ObjectStore>> { + 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); Review Comment: `AmazonS3Builder::new()` doesn't read `AWS_*` env the way `from_env()` does, so with no static creds object_store falls straight to web-identity/container/IMDS — someone migrating from the opendal backend (reqsign loads env + profile by default) silently loses `AWS_ACCESS_KEY_ID`/`AWS_REGION`/profile and ends up with IMDS timeouts or unauthenticated. I'd seed from `from_env()` and let explicit `S3Config` fields override. Worth confirming against object_store's builder — I'm going off the docs here, not a re-read of the source. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
