From 2b74bc688ef95fee9d1043339cf03cfb4d0dbaba Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Mon, 15 Jun 2026 18:28:21 +0800 Subject: [PATCH 01/10] feat(rest): use vended storage credentials from LoadTable response --- crates/catalog/rest/src/catalog.rs | 43 +++++++++++++++++++++--------- 1 file changed, 30 insertions(+), 13 deletions(-) diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index 03fda00bea..3660ecf79c 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -46,7 +46,7 @@ use crate::token::{ use crate::types::{ CatalogConfig, CommitTableRequest, CommitTableResponse, CreateNamespaceRequest, CreateTableRequest, ListNamespaceResponse, ListTablesResponse, LoadTableResult, - NamespaceResponse, RegisterTableRequest, RenameTableRequest, + NamespaceResponse, RegisterTableRequest, RenameTableRequest, StorageCredential, }; /// REST catalog URI @@ -903,6 +903,7 @@ impl Catalog for RestCatalog { let request = context .client .request(Method::POST, context.config.tables_endpoint(namespace)) + .header("X-Iceberg-Access-Delegation", "vended-credentials") .json(&CreateTableRequest { name: creation.name, location: creation.location, @@ -946,11 +947,7 @@ impl Catalog for RestCatalog { "Metadata location missing in `create_table` response!", ))?; - let config = response - .config - .into_iter() - .chain(self.user_config.props.clone()) - .collect(); + let config = table_file_io_config(&response, &self.user_config.props); let file_io = self .load_file_io(Some(metadata_location), Some(config)) @@ -980,6 +977,8 @@ impl Catalog for RestCatalog { let request = context .client .request(Method::GET, context.config.table_endpoint(table_ident)) + // Opt in to vended storage credentials. + .header("X-Iceberg-Access-Delegation", "vended-credentials") .build()?; let http_response = context.client.query_catalog(request).await?; @@ -1003,11 +1002,7 @@ impl Catalog for RestCatalog { } }; - let config = response - .config - .into_iter() - .chain(self.user_config.props.clone()) - .collect(); + let config = table_file_io_config(&response, &self.user_config.props); let file_io = self .load_file_io(response.metadata_location.as_deref(), Some(config)) @@ -1217,9 +1212,12 @@ impl Catalog for RestCatalog { } }; + // Reload for a credentialed FileIO (commit response carries no credentials). let file_io = self - .load_file_io(Some(&response.metadata_location), None) - .await?; + .load_table(commit.identifier()) + .await? + .file_io() + .clone(); Table::builder() .identifier(commit.identifier().clone()) @@ -1231,6 +1229,23 @@ impl Catalog for RestCatalog { } } +/// FileIO props: server `config`, then vended `storage_credentials` (longest prefix wins), then user props. +fn table_file_io_config( + response: &LoadTableResult, + user_props: &HashMap, +) -> HashMap { + let mut config: HashMap = response.config.clone(); + if let Some(creds) = response.storage_credentials.as_ref() { + let mut sorted: Vec<&StorageCredential> = creds.iter().collect(); + sorted.sort_by_key(|c| c.prefix.len()); + for cred in sorted { + config.extend(cred.config.clone()); + } + } + config.extend(user_props.clone()); + config +} + #[cfg(test)] mod tests { use std::fs::File; @@ -2898,6 +2913,7 @@ mod tests { let config_mock = create_config_mock(&mut server).await; + // GET hit twice: commit refresh + post-commit reload. let load_table_mock = server .mock("GET", "/v1/namespaces/ns1/tables/test1") .with_status(200) @@ -2906,6 +2922,7 @@ mod tests { env!("CARGO_MANIFEST_DIR"), "load_table_response.json" )) + .expect(2) .create_async() .await; From e829e8a97e235c715c0e12e9475984b758721bc8 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Tue, 16 Jun 2026 00:29:12 +0800 Subject: [PATCH 02/10] feat(rest): per-prefix vended credentials with opt-in delegation header --- crates/catalog/rest/src/catalog.rs | 155 ++++++++++---- .../load_table_response_with_credentials.json | 78 ++++++++ crates/iceberg/public-api.txt | 1 + crates/iceberg/src/io/file_io.rs | 189 +++++++++++++++--- crates/iceberg/src/table.rs | 6 + crates/iceberg/src/transaction/mod.rs | 4 +- 6 files changed, 375 insertions(+), 58 deletions(-) create mode 100644 crates/catalog/rest/testdata/load_table_response_with_credentials.json diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index 3660ecf79c..f9c7f2e2a8 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -603,6 +603,7 @@ impl RestCatalog { &self, metadata_location: Option<&str>, extra_config: Option>, + storage_credentials: Option<&[StorageCredential]>, ) -> Result { let mut props = self.context().await?.config.props.clone(); if let Some(config) = extra_config { @@ -635,9 +636,19 @@ impl RestCatalog { ) })?; - let file_io = FileIOBuilder::new(factory).with_props(props).build(); + let mut builder = FileIOBuilder::new(factory).with_props(props.clone()); - Ok(file_io) + // Vended credentials are scoped per location prefix: give each its own + // storage so reads/writes use the matching credentials. + if let Some(creds) = storage_credentials { + for cred in creds { + let mut prefixed = props.clone(); + prefixed.extend(cred.config.clone()); + builder = builder.with_prefixed_props(cred.prefix.clone(), prefixed); + } + } + + Ok(builder.build()) } /// Invalidate the current token without generating a new one. On the next request, the client @@ -903,7 +914,6 @@ impl Catalog for RestCatalog { let request = context .client .request(Method::POST, context.config.tables_endpoint(namespace)) - .header("X-Iceberg-Access-Delegation", "vended-credentials") .json(&CreateTableRequest { name: creation.name, location: creation.location, @@ -947,10 +957,15 @@ impl Catalog for RestCatalog { "Metadata location missing in `create_table` response!", ))?; - let config = table_file_io_config(&response, &self.user_config.props); + let mut base_config = response.config.clone(); + base_config.extend(self.user_config.props.clone()); let file_io = self - .load_file_io(Some(metadata_location), Some(config)) + .load_file_io( + Some(metadata_location), + Some(base_config), + response.storage_credentials.as_deref(), + ) .await?; let table_builder = Table::builder() @@ -974,11 +989,12 @@ impl Catalog for RestCatalog { async fn load_table(&self, table_ident: &TableIdent) -> Result { let context = self.context().await?; + // Vended credentials are opt-in via a `header.X-Iceberg-Access-Delegation` + // catalog property (applied to every request like the Iceberg Java client); + // any returned `storage_credentials` are wired into the FileIO below. let request = context .client .request(Method::GET, context.config.table_endpoint(table_ident)) - // Opt in to vended storage credentials. - .header("X-Iceberg-Access-Delegation", "vended-credentials") .build()?; let http_response = context.client.query_catalog(request).await?; @@ -1002,10 +1018,15 @@ impl Catalog for RestCatalog { } }; - let config = table_file_io_config(&response, &self.user_config.props); + let mut base_config = response.config.clone(); + base_config.extend(self.user_config.props.clone()); let file_io = self - .load_file_io(response.metadata_location.as_deref(), Some(config)) + .load_file_io( + response.metadata_location.as_deref(), + Some(base_config), + response.storage_credentials.as_deref(), + ) .await?; let table_builder = Table::builder() @@ -1141,7 +1162,9 @@ impl Catalog for RestCatalog { "Metadata location missing in `register_table` response!", ))?; - let file_io = self.load_file_io(Some(metadata_location), None).await?; + let file_io = self + .load_file_io(Some(metadata_location), None, None) + .await?; Table::builder() .identifier(table_ident.clone()) @@ -1212,12 +1235,11 @@ impl Catalog for RestCatalog { } }; - // Reload for a credentialed FileIO (commit response carries no credentials). + // The commit response carries no credentials, so build a plain FileIO; + // the transaction layer reuses the credentialed one it already holds. let file_io = self - .load_table(commit.identifier()) - .await? - .file_io() - .clone(); + .load_file_io(Some(&response.metadata_location), None, None) + .await?; Table::builder() .identifier(commit.identifier().clone()) @@ -1229,23 +1251,6 @@ impl Catalog for RestCatalog { } } -/// FileIO props: server `config`, then vended `storage_credentials` (longest prefix wins), then user props. -fn table_file_io_config( - response: &LoadTableResult, - user_props: &HashMap, -) -> HashMap { - let mut config: HashMap = response.config.clone(); - if let Some(creds) = response.storage_credentials.as_ref() { - let mut sorted: Vec<&StorageCredential> = creds.iter().collect(); - sorted.sort_by_key(|c| c.prefix.len()); - for cred in sorted { - config.extend(cred.config.clone()); - } - } - config.extend(user_props.clone()); - config -} - #[cfg(test)] mod tests { use std::fs::File; @@ -2698,6 +2703,88 @@ mod tests { rename_table_mock.assert_async().await; } + #[tokio::test] + async fn test_load_table_uses_vended_credentials() { + let mut server = Server::new_async().await; + + let config_mock = create_config_mock(&mut server).await; + + // Vended credentials are opt-in via a `header.*` catalog property (like the + // Java client). With it configured, the header is sent and the response's + // `storage-credentials` are accepted (the FileIO builds). + let load_table_mock = server + .mock("GET", "/v1/namespaces/ns1/tables/test1") + .match_header("x-iceberg-access-delegation", "vended-credentials") + .with_status(200) + .with_body_from_file(format!( + "{}/testdata/{}", + env!("CARGO_MANIFEST_DIR"), + "load_table_response_with_credentials.json" + )) + .create_async() + .await; + + let props = HashMap::from([( + "header.X-Iceberg-Access-Delegation".to_string(), + "vended-credentials".to_string(), + )]); + let catalog = RestCatalog::new( + RestCatalogConfig::builder() + .uri(server.url()) + .props(props) + .build(), + Some(Arc::new(LocalFsStorageFactory)), + Runtime::current(), + ); + + let table = catalog + .load_table(&TableIdent::from_strs(["ns1", "test1"]).unwrap()) + .await + .unwrap(); + + assert_eq!( + "s3://warehouse/database/table/metadata/00001-5f2f8166-244c-4eae-ac36-384ecdec81fc.gz.metadata.json", + table.metadata_location().unwrap() + ); + + config_mock.assert_async().await; + load_table_mock.assert_async().await; + } + + #[tokio::test] + async fn test_load_table_omits_delegation_header_by_default() { + let mut server = Server::new_async().await; + + let config_mock = create_config_mock(&mut server).await; + + // No delegation header is hardcoded: without a `header.*` prop, none is sent. + let load_table_mock = server + .mock("GET", "/v1/namespaces/ns1/tables/test1") + .match_header("x-iceberg-access-delegation", mockito::Matcher::Missing) + .with_status(200) + .with_body_from_file(format!( + "{}/testdata/{}", + env!("CARGO_MANIFEST_DIR"), + "load_table_response.json" + )) + .create_async() + .await; + + let catalog = RestCatalog::new( + RestCatalogConfig::builder().uri(server.url()).build(), + Some(Arc::new(LocalFsStorageFactory)), + Runtime::current(), + ); + + catalog + .load_table(&TableIdent::from_strs(["ns1", "test1"]).unwrap()) + .await + .unwrap(); + + config_mock.assert_async().await; + load_table_mock.assert_async().await; + } + #[tokio::test] async fn test_create_table() { let mut server = Server::new_async().await; @@ -2913,7 +3000,7 @@ mod tests { let config_mock = create_config_mock(&mut server).await; - // GET hit twice: commit refresh + post-commit reload. + // GET hit once: the transaction refreshes the table before committing. let load_table_mock = server .mock("GET", "/v1/namespaces/ns1/tables/test1") .with_status(200) @@ -2922,7 +3009,7 @@ mod tests { env!("CARGO_MANIFEST_DIR"), "load_table_response.json" )) - .expect(2) + .expect(1) .create_async() .await; diff --git a/crates/catalog/rest/testdata/load_table_response_with_credentials.json b/crates/catalog/rest/testdata/load_table_response_with_credentials.json new file mode 100644 index 0000000000..c0a540fbac --- /dev/null +++ b/crates/catalog/rest/testdata/load_table_response_with_credentials.json @@ -0,0 +1,78 @@ +{ + "metadata-location": "s3://warehouse/database/table/metadata/00001-5f2f8166-244c-4eae-ac36-384ecdec81fc.gz.metadata.json", + "metadata": { + "format-version": 1, + "table-uuid": "b55d9dda-6561-423a-8bfc-787980ce421f", + "location": "s3://warehouse/database/table", + "last-updated-ms": 1646787054459, + "last-column-id": 2, + "schema": { + "type": "struct", + "schema-id": 0, + "fields": [ + {"id": 1, "name": "id", "required": false, "type": "int"}, + {"id": 2, "name": "data", "required": false, "type": "string"} + ] + }, + "current-schema-id": 0, + "schemas": [ + { + "type": "struct", + "schema-id": 0, + "fields": [ + {"id": 1, "name": "id", "required": false, "type": "int"}, + {"id": 2, "name": "data", "required": false, "type": "string"} + ] + } + ], + "partition-spec": [], + "default-spec-id": 0, + "partition-specs": [{"spec-id": 0, "fields": []}], + "last-partition-id": 999, + "default-sort-order-id": 0, + "sort-orders": [{"order-id": 0, "fields": []}], + "properties": {"owner": "bryan", "write.metadata.compression-codec": "gzip"}, + "current-snapshot-id": 3497810964824022504, + "refs": {"main": {"snapshot-id": 3497810964824022504, "type": "branch"}}, + "snapshots": [ + { + "snapshot-id": 3497810964824022504, + "timestamp-ms": 1646787054459, + "summary": { + "operation": "append", + "spark.app.id": "local-1646787004168", + "added-data-files": "1", + "added-records": "1", + "added-files-size": "697", + "changed-partition-count": "1", + "total-records": "1", + "total-files-size": "697", + "total-data-files": "1", + "total-delete-files": "0", + "total-position-deletes": "0", + "total-equality-deletes": "0" + }, + "manifest-list": "s3://warehouse/database/table/metadata/snap-3497810964824022504-1-c4f68204-666b-4e50-a9df-b10c34bf6b82.avro", + "schema-id": 0 + } + ], + "snapshot-log": [{"timestamp-ms": 1646787054459, "snapshot-id": 3497810964824022504}], + "metadata-log": [ + { + "timestamp-ms": 1646787031514, + "metadata-file": "s3://warehouse/database/table/metadata/00000-88484a1c-00e5-4a07-a787-c0e7aeffa805.gz.metadata.json" + } + ] + }, + "config": {"client.factory": "io.tabular.iceberg.catalog.TabularAwsClientFactory", "region": "us-west-2"}, + "storage-credentials": [ + { + "prefix": "s3://warehouse/database/table", + "config": { + "s3.access-key-id": "vended-key-id", + "s3.secret-access-key": "vended-secret", + "s3.session-token": "vended-token" + } + } + ] +} diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 41d4d7691b..5959f10c74 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -701,6 +701,7 @@ impl iceberg::io::FileIOBuilder pub fn iceberg::io::FileIOBuilder::build(self) -> iceberg::io::FileIO pub fn iceberg::io::FileIOBuilder::config(&self) -> &iceberg::io::StorageConfig pub fn iceberg::io::FileIOBuilder::new(factory: alloc::sync::Arc) -> Self +pub fn iceberg::io::FileIOBuilder::with_prefixed_props(self, prefix: impl core::convert::Into, props: impl core::iter::traits::collect::IntoIterator) -> Self pub fn iceberg::io::FileIOBuilder::with_prop(self, key: impl alloc::string::ToString, value: impl alloc::string::ToString) -> Self pub fn iceberg::io::FileIOBuilder::with_props(self, args: impl core::iter::traits::collect::IntoIterator) -> Self impl core::clone::Clone for iceberg::io::FileIOBuilder diff --git a/crates/iceberg/src/io/file_io.rs b/crates/iceberg/src/io/file_io.rs index 227d8f4d5b..0338bca566 100644 --- a/crates/iceberg/src/io/file_io.rs +++ b/crates/iceberg/src/io/file_io.rs @@ -15,11 +15,12 @@ // specific language governing permissions and limitations // under the License. +use std::collections::HashMap; use std::ops::Range; use std::sync::{Arc, OnceLock}; use bytes::Bytes; -use futures::{Stream, StreamExt}; +use futures::{Stream, StreamExt, stream}; use super::storage::{ LocalFsStorageFactory, MemoryStorageFactory, Storage, StorageConfig, StorageFactory, @@ -67,6 +68,17 @@ pub struct FileIO { factory: Arc, /// Cached storage instance (lazily initialized) storage: Arc>>, + /// Per-prefix storages (longest prefix first) for tables that vend distinct + /// credentials per location prefix. Paths matching none use `storage` above. + prefixed: Arc>, +} + +/// A storage scoped to a location `prefix`, lazily built from its own config. +#[derive(Debug)] +struct PrefixedStorage { + prefix: String, + config: StorageConfig, + storage: OnceLock>, } impl FileIO { @@ -78,6 +90,7 @@ impl FileIO { config: StorageConfig::new(), factory: Arc::new(MemoryStorageFactory), storage: Arc::new(OnceLock::new()), + prefixed: Arc::new(Vec::new()), } } @@ -89,6 +102,7 @@ impl FileIO { config: StorageConfig::new(), factory: Arc::new(LocalFsStorageFactory), storage: Arc::new(OnceLock::new()), + prefixed: Arc::new(Vec::new()), } } @@ -97,24 +111,31 @@ impl FileIO { &self.config } - /// Get or create the storage instance. - /// - /// The factory is invoked on first access and the result is cached - /// for all subsequent operations. - fn get_storage(&self) -> Result> { - // Check if already initialized - if let Some(storage) = self.storage.get() { - return Ok(storage.clone()); + /// Get or create the storage for `path`, routing to the longest-matching + /// prefix storage if any, else the default. Built once, then cached. + fn get_storage(&self, path: &str) -> Result> { + // `prefixed` is sorted longest-first, so the first match is most specific. + for ps in self.prefixed.iter() { + if path.starts_with(&ps.prefix) { + return Self::get_or_build(&ps.storage, &self.factory, &ps.config); + } } + Self::get_or_build(&self.storage, &self.factory, &self.config) + } - // Build the storage - let storage = self.factory.build(&self.config)?; - - // Try to set it (another thread might have set it first) - let _ = self.storage.set(storage.clone()); - - // Return whatever is in the cell (either ours or another thread's) - Ok(self.storage.get().unwrap().clone()) + /// Get a cached storage from `cell`, building it from `config` on first use. + fn get_or_build( + cell: &OnceLock>, + factory: &Arc, + config: &StorageConfig, + ) -> Result> { + if let Some(storage) = cell.get() { + return Ok(storage.clone()); + } + let storage = factory.build(config)?; + // Another thread might have set it first; keep whatever ends up in the cell. + let _ = cell.set(storage); + Ok(cell.get().unwrap().clone()) } /// Deletes file. @@ -123,7 +144,7 @@ impl FileIO { /// /// * path: It should be *absolute* path starting with scheme string used to construct [`FileIO`]. pub async fn delete(&self, path: impl AsRef) -> Result<()> { - self.get_storage()?.delete(path.as_ref()).await + self.get_storage(path.as_ref())?.delete(path.as_ref()).await } /// Remove the path and all nested dirs and files recursively. @@ -138,7 +159,9 @@ impl FileIO { /// - If the path is a empty directory, this function will remove the directory itself. /// - If the path is a non-empty directory, this function will remove the directory and all nested files and directories. pub async fn delete_prefix(&self, path: impl AsRef) -> Result<()> { - self.get_storage()?.delete_prefix(path.as_ref()).await + self.get_storage(path.as_ref())? + .delete_prefix(path.as_ref()) + .await } /// Delete multiple files from a stream of paths. @@ -150,7 +173,29 @@ impl FileIO { &self, paths: impl Stream + Send + 'static, ) -> Result<()> { - self.get_storage()?.delete_stream(paths.boxed()).await + // No per-prefix storages: delete the whole batch on the default storage. + if self.prefixed.is_empty() { + return self.get_storage("")?.delete_stream(paths.boxed()).await; + } + + // Otherwise group paths by routed storage, then batch-delete each group. + let mut groups: HashMap> = HashMap::new(); + let mut paths = paths.boxed(); + while let Some(path) = paths.next().await { + let key = self + .prefixed + .iter() + .find(|ps| path.starts_with(&ps.prefix)) + .map(|ps| ps.prefix.clone()) + .unwrap_or_default(); + groups.entry(key).or_default().push(path); + } + + for batch in groups.into_values() { + let storage = self.get_storage(&batch[0])?; + storage.delete_stream(stream::iter(batch).boxed()).await?; + } + Ok(()) } /// Check file exists. @@ -159,7 +204,7 @@ impl FileIO { /// /// * path: It should be *absolute* path starting with scheme string used to construct [`FileIO`]. pub async fn exists(&self, path: impl AsRef) -> Result { - self.get_storage()?.exists(path.as_ref()).await + self.get_storage(path.as_ref())?.exists(path.as_ref()).await } /// Creates input file. @@ -168,7 +213,7 @@ impl FileIO { /// /// * path: It should be *absolute* path starting with scheme string used to construct [`FileIO`]. pub fn new_input(&self, path: impl AsRef) -> Result { - self.get_storage()?.new_input(path.as_ref()) + self.get_storage(path.as_ref())?.new_input(path.as_ref()) } /// Creates output file. @@ -177,7 +222,7 @@ impl FileIO { /// /// * path: It should be *absolute* path starting with scheme string used to construct [`FileIO`]. pub fn new_output(&self, path: impl AsRef) -> Result { - self.get_storage()?.new_output(path.as_ref()) + self.get_storage(path.as_ref())?.new_output(path.as_ref()) } } @@ -191,6 +236,8 @@ pub struct FileIOBuilder { factory: Arc, /// Storage configuration config: StorageConfig, + /// Per-location-prefix configs (prefix, config). + prefixed: Vec<(String, StorageConfig)>, } impl FileIOBuilder { @@ -199,6 +246,7 @@ impl FileIOBuilder { Self { factory, config: StorageConfig::new(), + prefixed: Vec::new(), } } @@ -219,6 +267,23 @@ impl FileIOBuilder { self } + /// Add a per-prefix storage config. Paths starting with `prefix` (longest + /// match wins) use these props instead of the default config. + pub fn with_prefixed_props( + mut self, + prefix: impl Into, + props: impl IntoIterator, + ) -> Self { + let config = StorageConfig::from_props( + props + .into_iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + ); + self.prefixed.push((prefix.into(), config)); + self + } + /// Get the storage configuration. pub fn config(&self) -> &StorageConfig { &self.config @@ -226,10 +291,22 @@ impl FileIOBuilder { /// Builds [`FileIO`]. pub fn build(self) -> FileIO { + let mut prefixed: Vec = self + .prefixed + .into_iter() + .map(|(prefix, config)| PrefixedStorage { + prefix, + config, + storage: OnceLock::new(), + }) + .collect(); + // Longest prefix first so routing picks the most specific match. + prefixed.sort_by(|a, b| b.prefix.len().cmp(&a.prefix.len())); FileIO { config: self.config, factory: self.factory, storage: Arc::new(OnceLock::new()), + prefixed: Arc::new(prefixed), } } } @@ -544,4 +621,70 @@ mod tests { assert_eq!(file_io.config().get("key1"), Some(&"value1".to_string())); assert_eq!(file_io.config().get("key2"), Some(&"value2".to_string())); } + + #[tokio::test] + async fn test_prefixed_props_sorted_by_descending_prefix_length() { + let factory = Arc::new(MemoryStorageFactory); + let file_io = FileIOBuilder::new(factory) + .with_prefixed_props("memory://a/", [("k", "short")]) + .with_prefixed_props("memory://a/longer/", [("k", "long")]) + .build(); + + // Longest prefix first so the most specific match wins at routing time. + let prefixes: Vec<&str> = file_io.prefixed.iter().map(|p| p.prefix.as_str()).collect(); + assert_eq!(prefixes, vec!["memory://a/longer/", "memory://a/"]); + } + + #[tokio::test] + async fn test_get_storage_routes_by_prefix() { + let factory = Arc::new(MemoryStorageFactory); + let file_io = FileIOBuilder::new(factory) + .with_prop("scope", "default") + .with_prefixed_props("memory://creds/", [("scope", "prefixed")]) + .build(); + + let default_a = file_io.get_storage("memory://other/x").unwrap(); + let default_b = file_io.get_storage("memory://other/y").unwrap(); + let prefixed_a = file_io.get_storage("memory://creds/x").unwrap(); + let prefixed_b = file_io.get_storage("memory://creds/y").unwrap(); + + // Repeated routing to the same bucket returns the memoized storage... + assert!(Arc::ptr_eq(&default_a, &default_b)); + assert!(Arc::ptr_eq(&prefixed_a, &prefixed_b)); + // ...and a prefix-matching path resolves to a distinct storage from the default. + assert!(!Arc::ptr_eq(&default_a, &prefixed_a)); + } + + #[tokio::test] + async fn test_delete_stream_routes_by_prefix() { + let factory = Arc::new(MemoryStorageFactory); + let file_io = FileIOBuilder::new(factory) + .with_prefixed_props("memory:/creds/", [("k", "v")]) + .build(); + + // One file under each routing bucket (default vs prefixed storage). + let default_path = "memory:/other/a.txt"; + let prefixed_path = "memory:/creds/b.txt"; + for path in [default_path, prefixed_path] { + file_io + .new_output(path) + .unwrap() + .write("x".into()) + .await + .unwrap(); + assert!(file_io.exists(path).await.unwrap()); + } + + // delete_stream must route each path to the storage that holds it. + file_io + .delete_stream(futures::stream::iter(vec![ + default_path.to_string(), + prefixed_path.to_string(), + ])) + .await + .unwrap(); + + assert!(!file_io.exists(default_path).await.unwrap()); + assert!(!file_io.exists(prefixed_path).await.unwrap()); + } } diff --git a/crates/iceberg/src/table.rs b/crates/iceberg/src/table.rs index 31feade038..2885e33644 100644 --- a/crates/iceberg/src/table.rs +++ b/crates/iceberg/src/table.rs @@ -220,6 +220,12 @@ impl Table { self } + /// Sets the [`Table`] `FileIO` and returns an updated instance. + pub(crate) fn with_file_io(mut self, file_io: FileIO) -> Self { + self.file_io = file_io; + self + } + /// Returns a TableBuilder to build a table pub fn builder() -> TableBuilder { TableBuilder::new() diff --git a/crates/iceberg/src/transaction/mod.rs b/crates/iceberg/src/transaction/mod.rs index 1b50c98a0e..50f74a8178 100644 --- a/crates/iceberg/src/transaction/mod.rs +++ b/crates/iceberg/src/transaction/mod.rs @@ -258,7 +258,9 @@ impl Transaction { .requirements(existing_requirements) .build(); - catalog.update_table(table_commit).await + let committed = catalog.update_table(table_commit).await?; + // Reuse the credentialed FileIO we already hold; the commit response has none. + Ok(committed.with_file_io(self.table.file_io().clone())) } } From 024f9b4805bc51f705436633bd6f51f4beb4c52c Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Wed, 17 Jun 2026 17:05:34 +0800 Subject: [PATCH 03/10] fix(io): bound delete_stream memory by flushing per-prefix batches --- crates/iceberg/src/io/file_io.rs | 86 ++++++++++++++++++++++++++++++-- 1 file changed, 82 insertions(+), 4 deletions(-) diff --git a/crates/iceberg/src/io/file_io.rs b/crates/iceberg/src/io/file_io.rs index 0338bca566..8666f6c24a 100644 --- a/crates/iceberg/src/io/file_io.rs +++ b/crates/iceberg/src/io/file_io.rs @@ -178,7 +178,9 @@ impl FileIO { return self.get_storage("")?.delete_stream(paths.boxed()).await; } - // Otherwise group paths by routed storage, then batch-delete each group. + // Route by prefix, flushing bounded batches as we iterate so memory stays + // bounded on large streams (like Java's `S3FileIO.deleteFiles`). + const DELETE_BATCH_SIZE: usize = 1000; let mut groups: HashMap> = HashMap::new(); let mut paths = paths.boxed(); while let Some(path) = paths.next().await { @@ -188,12 +190,24 @@ impl FileIO { .find(|ps| path.starts_with(&ps.prefix)) .map(|ps| ps.prefix.clone()) .unwrap_or_default(); - groups.entry(key).or_default().push(path); + let buf = groups.entry(key).or_default(); + buf.push(path); + if buf.len() >= DELETE_BATCH_SIZE { + let full = std::mem::take(buf); + self.get_storage(&full[0])? + .delete_stream(stream::iter(full).boxed()) + .await?; + } } + // Flush remainders. for batch in groups.into_values() { - let storage = self.get_storage(&batch[0])?; - storage.delete_stream(stream::iter(batch).boxed()).await?; + if batch.is_empty() { + continue; + } + self.get_storage(&batch[0])? + .delete_stream(stream::iter(batch).boxed()) + .await?; } Ok(()) } @@ -635,6 +649,39 @@ mod tests { assert_eq!(prefixes, vec!["memory://a/longer/", "memory://a/"]); } + #[tokio::test] + async fn test_prefixed_config_carries_credential_values() { + // Prefix config gets the vended credentials; default config keeps only base props. + let factory = Arc::new(MemoryStorageFactory); + let file_io = FileIOBuilder::new(factory) + .with_prop("s3.region", "us-east-1") + .with_prefixed_props("s3://bucket/table", [ + ("s3.region", "us-east-1"), + ("s3.access-key-id", "vended-key"), + ("s3.secret-access-key", "vended-secret"), + ]) + .build(); + + // Default: base props, no credentials. + assert_eq!( + file_io.config().get("s3.region"), + Some(&"us-east-1".to_string()) + ); + assert_eq!(file_io.config().get("s3.access-key-id"), None); + + // Prefix: base props + vended credentials. + let prefixed = &file_io.prefixed[0].config; + assert_eq!(prefixed.get("s3.region"), Some(&"us-east-1".to_string())); + assert_eq!( + prefixed.get("s3.access-key-id"), + Some(&"vended-key".to_string()) + ); + assert_eq!( + prefixed.get("s3.secret-access-key"), + Some(&"vended-secret".to_string()) + ); + } + #[tokio::test] async fn test_get_storage_routes_by_prefix() { let factory = Arc::new(MemoryStorageFactory); @@ -687,4 +734,35 @@ mod tests { assert!(!file_io.exists(default_path).await.unwrap()); assert!(!file_io.exists(prefixed_path).await.unwrap()); } + + #[tokio::test] + async fn test_delete_stream_flushes_across_batches() { + // More than the flush threshold (1000): exercises mid-stream flush + remainder. + let factory = Arc::new(MemoryStorageFactory); + let file_io = FileIOBuilder::new(factory) + .with_prefixed_props("memory:/creds/", [("k", "v")]) + .build(); + + let n = 1050; + let mut paths = Vec::with_capacity(n); + for i in 0..n { + let p = format!("memory:/creds/f{i}.txt"); + file_io + .new_output(&p) + .unwrap() + .write("x".into()) + .await + .unwrap(); + paths.push(p); + } + + file_io + .delete_stream(futures::stream::iter(paths.clone())) + .await + .unwrap(); + + for p in &paths { + assert!(!file_io.exists(p).await.unwrap()); + } + } } From 56fac2bbd98aebb271f2c6016aa3e8e80d5010f7 Mon Sep 17 00:00:00 2001 From: Li Jiajia Date: Fri, 10 Jul 2026 04:56:21 -0400 Subject: [PATCH 04/10] fix(rest): route vended credentials by prefix and keep them in object cache --- crates/catalog/rest/src/catalog.rs | 5 ++- crates/iceberg/src/io/file_io.rs | 4 +- crates/iceberg/src/io/object_cache.rs | 11 ++++++ crates/iceberg/src/io/storage/config/mod.rs | 41 +++++++++++++++++++- crates/iceberg/src/table.rs | 42 +++++++++++++++++++++ crates/iceberg/src/transaction/mod.rs | 18 ++++++++- 6 files changed, 116 insertions(+), 5 deletions(-) diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index f9c7f2e2a8..41fabf9185 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -639,7 +639,8 @@ impl RestCatalog { let mut builder = FileIOBuilder::new(factory).with_props(props.clone()); // Vended credentials are scoped per location prefix: give each its own - // storage so reads/writes use the matching credentials. + // storage. Paths under no vended prefix fall back to the default `props` + // above, which carry no credentials. if let Some(creds) = storage_credentials { for cred in creds { let mut prefixed = props.clone(); @@ -2735,6 +2736,7 @@ mod tests { .build(), Some(Arc::new(LocalFsStorageFactory)), Runtime::current(), + None, ); let table = catalog @@ -2774,6 +2776,7 @@ mod tests { RestCatalogConfig::builder().uri(server.url()).build(), Some(Arc::new(LocalFsStorageFactory)), Runtime::current(), + None, ); catalog diff --git a/crates/iceberg/src/io/file_io.rs b/crates/iceberg/src/io/file_io.rs index 8666f6c24a..973098e2ff 100644 --- a/crates/iceberg/src/io/file_io.rs +++ b/crates/iceberg/src/io/file_io.rs @@ -115,6 +115,8 @@ impl FileIO { /// prefix storage if any, else the default. Built once, then cached. fn get_storage(&self, path: &str) -> Result> { // `prefixed` is sorted longest-first, so the first match is most specific. + // Selection is by longest matching string prefix, per the Iceberg REST + // spec's storage-credentials semantics (and Java's `S3FileIO`). for ps in self.prefixed.iter() { if path.starts_with(&ps.prefix) { return Self::get_or_build(&ps.storage, &self.factory, &ps.config); @@ -315,7 +317,7 @@ impl FileIOBuilder { }) .collect(); // Longest prefix first so routing picks the most specific match. - prefixed.sort_by(|a, b| b.prefix.len().cmp(&a.prefix.len())); + prefixed.sort_by_key(|item| std::cmp::Reverse(item.prefix.len())); FileIO { config: self.config, factory: self.factory, diff --git a/crates/iceberg/src/io/object_cache.rs b/crates/iceberg/src/io/object_cache.rs index 9d7815569b..29af0169f6 100644 --- a/crates/iceberg/src/io/object_cache.rs +++ b/crates/iceberg/src/io/object_cache.rs @@ -95,6 +95,17 @@ impl ObjectCache { } } + /// Returns a cache that uses `file_io` for future cache misses. + pub(crate) fn with_file_io(mut self, file_io: FileIO) -> Self { + self.file_io = file_io; + self + } + + #[cfg(test)] + pub(crate) fn file_io(&self) -> &FileIO { + &self.file_io + } + /// Retrieves an Arc [`Manifest`] from the cache /// or retrieves one from FileIO and parses it if not present pub(crate) async fn get_manifest(&self, manifest_file: &ManifestFile) -> Result> { diff --git a/crates/iceberg/src/io/storage/config/mod.rs b/crates/iceberg/src/io/storage/config/mod.rs index 2350aab6dd..d75f31ae1e 100644 --- a/crates/iceberg/src/io/storage/config/mod.rs +++ b/crates/iceberg/src/io/storage/config/mod.rs @@ -51,12 +51,23 @@ use serde::{Deserialize, Serialize}; /// which storage backend to use. The storage type is determined by the /// explicit factory selection. /// ``` -#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize, Default)] +#[derive(Clone, PartialEq, Eq, Serialize, Deserialize, Default)] pub struct StorageConfig { /// Configuration properties for the storage backend props: HashMap, } +impl std::fmt::Debug for StorageConfig { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + // Property values may hold vended credentials (e.g. `s3.secret-access-key`, + // `s3.session-token`). Debug is reachable through the `FileIO`/`Table` + // derives, so print only the keys and never the secret values. + f.debug_struct("StorageConfig") + .field("keys", &self.props.keys().collect::>()) + .finish_non_exhaustive() + } +} + impl StorageConfig { /// Create a new empty StorageConfig. pub fn new() -> Self { @@ -151,6 +162,34 @@ mod tests { assert!(config.props().is_empty()); } + #[test] + fn test_debug_redacts_credential_values() { + let config = StorageConfig::from_props(HashMap::from([ + ("s3.access-key-id".to_string(), "vended-key".to_string()), + ( + "s3.secret-access-key".to_string(), + "super-secret".to_string(), + ), + ("s3.session-token".to_string(), "vended-token".to_string()), + ])); + + let rendered = format!("{config:?}"); + // Secret values must never appear in Debug output (reachable via FileIO/Table). + assert!( + !rendered.contains("super-secret"), + "leaked secret: {rendered}" + ); + assert!( + !rendered.contains("vended-token"), + "leaked token: {rendered}" + ); + // Keys stay visible so routing/config is still diagnosable. + assert!( + rendered.contains("s3.secret-access-key"), + "keys hidden: {rendered}" + ); + } + #[test] fn test_storage_config_get() { let config = StorageConfig::new().with_prop("region", "us-east-1"); diff --git a/crates/iceberg/src/table.rs b/crates/iceberg/src/table.rs index 2885e33644..a273a5c819 100644 --- a/crates/iceberg/src/table.rs +++ b/crates/iceberg/src/table.rs @@ -222,6 +222,12 @@ impl Table { /// Sets the [`Table`] `FileIO` and returns an updated instance. pub(crate) fn with_file_io(mut self, file_io: FileIO) -> Self { + self.object_cache = Arc::new( + self.object_cache + .as_ref() + .clone() + .with_file_io(file_io.clone()), + ); self.file_io = file_io; self } @@ -416,7 +422,9 @@ mod tests { use super::*; use crate::encryption::SensitiveBytes; use crate::encryption::kms::MemoryKeyManagementClient; + use crate::io::{FileIOBuilder, MemoryStorageFactory}; use crate::spec::TableProperties; + use crate::test_utils::test_runtime; fn load_test_metadata(filename: &str) -> TableMetadata { let path = format!( @@ -501,6 +509,40 @@ mod tests { assert_eq!(table.identifier.name(), "table"); } + #[test] + fn test_with_file_io_updates_object_cache_file_io() { + let metadata = load_test_metadata("TableMetadataV2ValidMinimal.json"); + let original_file_io = FileIOBuilder::new(Arc::new(MemoryStorageFactory)) + .with_prop("marker", "original") + .build(); + let replacement_file_io = FileIOBuilder::new(Arc::new(MemoryStorageFactory)) + .with_prop("marker", "replacement") + .build(); + + let table = Table::builder() + .metadata(metadata) + .identifier(TableIdent::from_strs(["ns", "table"]).unwrap()) + .file_io(original_file_io) + .runtime(test_runtime()) + .build() + .unwrap() + .with_file_io(replacement_file_io); + + assert_eq!( + table.file_io().config().get("marker").map(String::as_str), + Some("replacement") + ); + assert_eq!( + table + .object_cache() + .file_io() + .config() + .get("marker") + .map(String::as_str), + Some("replacement") + ); + } + fn make_kms() -> Arc { let kms = MemoryKeyManagementClient::new(); kms.add_master_key("master-1").unwrap(); diff --git a/crates/iceberg/src/transaction/mod.rs b/crates/iceberg/src/transaction/mod.rs index 50f74a8178..bdce9b0b56 100644 --- a/crates/iceberg/src/transaction/mod.rs +++ b/crates/iceberg/src/transaction/mod.rs @@ -252,6 +252,12 @@ impl Transaction { )?; } + // A location change moves metadata/data to a new prefix that the refresh + // load's vended credentials do not cover, so it needs a post-commit reload. + let location_changed = existing_updates + .iter() + .any(|update| matches!(update, TableUpdate::SetLocation { .. })); + let table_commit = TableCommit::builder() .ident(self.table.identifier().to_owned()) .updates(existing_updates) @@ -259,8 +265,16 @@ impl Transaction { .build(); let committed = catalog.update_table(table_commit).await?; - // Reuse the credentialed FileIO we already hold; the commit response has none. - Ok(committed.with_file_io(self.table.file_io().clone())) + if location_changed { + // The new location has its own vended credentials; the reused FileIO is + // scoped to the old prefix, so reload the table to pick them up. + catalog.load_table(committed.identifier()).await + } else { + // The commit response carries no credentials. Reuse the FileIO from the + // refresh load above (not `self.table`, which is left untouched when the + // metadata is unchanged) so freshly vended credentials are not dropped. + Ok(committed.with_file_io(refreshed.file_io().clone())) + } } } From 016a2a1c1a8342a721d13aa870334cd5459076fd Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Tue, 11 Aug 2026 16:31:08 -0400 Subject: [PATCH 05/10] fix incorrect signature --- crates/catalog/rest/src/catalog.rs | 2 -- 1 file changed, 2 deletions(-) diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index 41fabf9185..d8b3ec6dbd 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -2736,7 +2736,6 @@ mod tests { .build(), Some(Arc::new(LocalFsStorageFactory)), Runtime::current(), - None, ); let table = catalog @@ -2776,7 +2775,6 @@ mod tests { RestCatalogConfig::builder().uri(server.url()).build(), Some(Arc::new(LocalFsStorageFactory)), Runtime::current(), - None, ); catalog From 881d40f2ca573ab5c2b10f87efd2e4300cf78cbf Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Tue, 18 Aug 2026 09:23:48 -0400 Subject: [PATCH 06/10] chore: reformat Cargo.toml reqwest entry per taplo `taplo fmt --check` in the lint job rejects the single-line `features` array added for `native-tls`; it exceeds the 80-column width, so taplo expands it. Co-Authored-By: Claude Opus 5 (1M context) --- Cargo.toml | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index 47f30e4d7d..c2b9e80869 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -131,7 +131,10 @@ pretty_assertions = "1.4" pyo3 = "0.28" rand = "0.9.3" regex = "1.11.3" -reqwest = { version = "0.12.12", default-features = false, features = ["json", "native-tls"] } +reqwest = { version = "0.12.12", default-features = false, features = [ + "json", + "native-tls", +] } roaring = { version = "0.11" } rstest = "0.26" serde = { version = "1.0.219", features = ["rc"] } From b89c3a21349a4373c886b7375a655b07486a4ef6 Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Tue, 18 Aug 2026 09:23:48 -0400 Subject: [PATCH 07/10] chore(ci): repin swatinem/rust-cache to v2.9.2 The `# v2` comment no longer matches the pinned commit now that the `v2` tag has moved, which zizmor flags as ref-version-mismatch. Repin to the v2.9.2 commit, matching apache/iceberg-rust main. Co-Authored-By: Claude Opus 5 (1M context) --- .github/workflows/bindings_python_ci.yml | 2 +- .github/workflows/ci.yml | 8 ++++---- .github/workflows/public-api.yml | 2 +- 3 files changed, 6 insertions(+), 6 deletions(-) diff --git a/.github/workflows/bindings_python_ci.yml b/.github/workflows/bindings_python_ci.yml index 377e4d71cb..a80b9349d2 100644 --- a/.github/workflows/bindings_python_ci.yml +++ b/.github/workflows/bindings_python_ci.yml @@ -88,7 +88,7 @@ jobs: uses: ./.github/actions/setup-builder - name: Cache Rust artifacts if: runner.os != 'Linux' - uses: swatinem/rust-cache@e18b497796c12c097a38f9edb9d0641fb99eee32 # v2 + uses: swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2.9.2 with: key: bindings-python save-if: ${{ github.event_name == 'push' && github.ref == 'refs/heads/main' }} diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index d3c1713484..e6c827a95e 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -122,7 +122,7 @@ jobs: uses: ./.github/actions/setup-builder - name: Cache Rust artifacts - uses: swatinem/rust-cache@e18b497796c12c097a38f9edb9d0641fb99eee32 # v2 + uses: swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2.9.2 with: save-if: ${{ github.event_name == 'push' && github.ref == 'refs/heads/main' }} @@ -148,7 +148,7 @@ jobs: uses: ./.github/actions/setup-builder - name: Cache Rust artifacts - uses: swatinem/rust-cache@e18b497796c12c097a38f9edb9d0641fb99eee32 # v2 + uses: swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2.9.2 with: save-if: ${{ github.event_name == 'push' && github.ref == 'refs/heads/main' }} @@ -188,7 +188,7 @@ jobs: uses: ./.github/actions/setup-builder - name: Cache Rust artifacts - uses: swatinem/rust-cache@e18b497796c12c097a38f9edb9d0641fb99eee32 # v2 + uses: swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2.9.2 with: save-if: ${{ github.event_name == 'push' && github.ref == 'refs/heads/main' }} @@ -222,7 +222,7 @@ jobs: repo-token: ${{ secrets.GITHUB_TOKEN }} - name: Cache Rust artifacts - uses: swatinem/rust-cache@e18b497796c12c097a38f9edb9d0641fb99eee32 # v2 + uses: swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2.9.2 with: key: ${{ matrix.test-suite.name }} save-if: ${{ github.event_name == 'push' && github.ref == 'refs/heads/main' }} diff --git a/.github/workflows/public-api.yml b/.github/workflows/public-api.yml index f5149d30d4..a8d1030560 100644 --- a/.github/workflows/public-api.yml +++ b/.github/workflows/public-api.yml @@ -44,7 +44,7 @@ jobs: uses: ./.github/actions/setup-builder - name: Cache Rust artifacts - uses: swatinem/rust-cache@e18b497796c12c097a38f9edb9d0641fb99eee32 # v2 + uses: swatinem/rust-cache@6323deb102c322ba6fcbdcafc7e3dddab59af2b6 # v2.9.2 with: save-if: ${{ github.event_name == 'push' && github.ref == 'refs/heads/main' }} From f26d5ca0504b66e43c10aa179df4226a5806d9eb Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Tue, 18 Aug 2026 14:59:11 -0400 Subject: [PATCH 08/10] chore: regenerate public-api snapshots The snapshots were never regenerated after the downstream commits that added `RequestAuthenticator`/`TokenProvider` to iceberg-catalog-rest and `DeltaWriter`/`PositionDeleteFileWriter` plus the `RollingFileWriterBuilder` schema parameter to iceberg, so `make check-public-api` fails. Regenerated with cargo-public-api 0.51.0 on nightly-2026-03-05, matching CI. Co-Authored-By: Claude Opus 5 (1M context) --- crates/catalog/rest/public-api.txt | 56 ++++++++++++++++++++++ crates/iceberg/public-api.txt | 74 +++++++++++++++++++++++++++++- 2 files changed, 128 insertions(+), 2 deletions(-) diff --git a/crates/catalog/rest/public-api.txt b/crates/catalog/rest/public-api.txt index 027df29b24..84e8f043d1 100644 --- a/crates/catalog/rest/public-api.txt +++ b/crates/catalog/rest/public-api.txt @@ -1,4 +1,21 @@ pub mod iceberg_catalog_rest +pub struct iceberg_catalog_rest::AuthenticatorConfig +pub iceberg_catalog_rest::AuthenticatorConfig::credential: core::option::Option<(core::option::Option, alloc::string::String)> +pub iceberg_catalog_rest::AuthenticatorConfig::custom_authenticator: core::option::Option> +pub iceberg_catalog_rest::AuthenticatorConfig::extra_oauth_params: std::collections::hash::map::HashMap +pub iceberg_catalog_rest::AuthenticatorConfig::token: core::option::Option +pub iceberg_catalog_rest::AuthenticatorConfig::token_endpoint: alloc::string::String +pub struct iceberg_catalog_rest::BearerTokenAuthenticator +impl iceberg_catalog_rest::BearerTokenAuthenticator +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::new(provider: alloc::sync::Arc) -> Self +impl core::clone::Clone for iceberg_catalog_rest::BearerTokenAuthenticator +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::clone(&self) -> iceberg_catalog_rest::BearerTokenAuthenticator +impl core::fmt::Debug for iceberg_catalog_rest::BearerTokenAuthenticator +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl iceberg_catalog_rest::RequestAuthenticator for iceberg_catalog_rest::BearerTokenAuthenticator +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::authenticate_request<'life0, 'life1, 'async_trait>(&'life0 self, req: &'life1 mut reqwest::async_impl::request::Request) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::invalidate_cache<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::regenerate_cache<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait pub struct iceberg_catalog_rest::CommitTableRequest pub iceberg_catalog_rest::CommitTableRequest::identifier: core::option::Option pub iceberg_catalog_rest::CommitTableRequest::requirements: alloc::vec::Vec @@ -152,6 +169,15 @@ impl serde_core::ser::Serialize for iceberg_catalog_rest::NamespaceResponse pub fn iceberg_catalog_rest::NamespaceResponse::serialize<__S>(&self, __serializer: __S) -> core::result::Result<<__S as serde_core::ser::Serializer>::Ok, <__S as serde_core::ser::Serializer>::Error> where __S: serde_core::ser::Serializer impl<'de> serde_core::de::Deserialize<'de> for iceberg_catalog_rest::NamespaceResponse pub fn iceberg_catalog_rest::NamespaceResponse::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> +pub struct iceberg_catalog_rest::OAuth2TokenProvider +impl iceberg_catalog_rest::OAuth2TokenProvider +pub fn iceberg_catalog_rest::OAuth2TokenProvider::new(client: reqwest::async_impl::client::Client, client_id: core::option::Option, client_secret: alloc::string::String, token_endpoint: alloc::string::String, extra_headers: http::header::map::HeaderMap, extra_oauth_params: std::collections::hash::map::HashMap, cached_token: core::option::Option) -> Self +impl core::fmt::Debug for iceberg_catalog_rest::OAuth2TokenProvider +pub fn iceberg_catalog_rest::OAuth2TokenProvider::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl iceberg_catalog_rest::TokenProvider for iceberg_catalog_rest::OAuth2TokenProvider +pub fn iceberg_catalog_rest::OAuth2TokenProvider::invalidate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::OAuth2TokenProvider::regenerate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::OAuth2TokenProvider::token<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait pub struct iceberg_catalog_rest::RegisterTableRequest pub iceberg_catalog_rest::RegisterTableRequest::metadata_location: alloc::string::String pub iceberg_catalog_rest::RegisterTableRequest::name: alloc::string::String @@ -207,6 +233,7 @@ pub fn iceberg_catalog_rest::RestCatalog::update_namespace<'life0, 'life1, 'asyn pub fn iceberg_catalog_rest::RestCatalog::update_table<'life0, 'async_trait>(&'life0 self, commit: iceberg::catalog::TableCommit) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait pub struct iceberg_catalog_rest::RestCatalogBuilder impl iceberg_catalog_rest::RestCatalogBuilder +pub fn iceberg_catalog_rest::RestCatalogBuilder::with_authenticator(self, authenticator: alloc::sync::Arc) -> Self pub fn iceberg_catalog_rest::RestCatalogBuilder::with_client(self, client: reqwest::async_impl::client::Client) -> Self impl core::default::Default for iceberg_catalog_rest::RestCatalogBuilder pub fn iceberg_catalog_rest::RestCatalogBuilder::default() -> Self @@ -217,6 +244,15 @@ pub type iceberg_catalog_rest::RestCatalogBuilder::C = iceberg_catalog_rest::Res pub fn iceberg_catalog_rest::RestCatalogBuilder::load(self, name: impl core::convert::Into, props: std::collections::hash::map::HashMap) -> impl core::future::future::Future> + core::marker::Send pub fn iceberg_catalog_rest::RestCatalogBuilder::with_runtime(self, runtime: iceberg::runtime::Runtime) -> Self pub fn iceberg_catalog_rest::RestCatalogBuilder::with_storage_factory(self, storage_factory: alloc::sync::Arc) -> Self +pub struct iceberg_catalog_rest::StaticTokenProvider +impl iceberg_catalog_rest::StaticTokenProvider +pub fn iceberg_catalog_rest::StaticTokenProvider::new(token: impl core::convert::Into) -> Self +impl core::fmt::Debug for iceberg_catalog_rest::StaticTokenProvider +pub fn iceberg_catalog_rest::StaticTokenProvider::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl iceberg_catalog_rest::TokenProvider for iceberg_catalog_rest::StaticTokenProvider +pub fn iceberg_catalog_rest::StaticTokenProvider::invalidate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::StaticTokenProvider::regenerate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::StaticTokenProvider::token<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait pub struct iceberg_catalog_rest::StorageCredential pub iceberg_catalog_rest::StorageCredential::config: std::collections::hash::map::HashMap pub iceberg_catalog_rest::StorageCredential::prefix: alloc::string::String @@ -266,3 +302,23 @@ pub fn iceberg_catalog_rest::UpdateNamespacePropertiesResponse::deserialize<__D> pub const iceberg_catalog_rest::REST_CATALOG_PROP_DISABLE_HEADER_REDACTION: &str pub const iceberg_catalog_rest::REST_CATALOG_PROP_URI: &str pub const iceberg_catalog_rest::REST_CATALOG_PROP_WAREHOUSE: &str +pub trait iceberg_catalog_rest::RequestAuthenticator: core::marker::Send + core::marker::Sync + core::fmt::Debug +pub fn iceberg_catalog_rest::RequestAuthenticator::authenticate_request<'life0, 'life1, 'async_trait>(&'life0 self, req: &'life1 mut reqwest::async_impl::request::Request) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_catalog_rest::RequestAuthenticator::invalidate_cache<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::RequestAuthenticator::regenerate_cache<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl iceberg_catalog_rest::RequestAuthenticator for iceberg_catalog_rest::BearerTokenAuthenticator +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::authenticate_request<'life0, 'life1, 'async_trait>(&'life0 self, req: &'life1 mut reqwest::async_impl::request::Request) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::invalidate_cache<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::BearerTokenAuthenticator::regenerate_cache<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub trait iceberg_catalog_rest::TokenProvider: core::marker::Send + core::marker::Sync + core::fmt::Debug +pub fn iceberg_catalog_rest::TokenProvider::invalidate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::TokenProvider::regenerate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::TokenProvider::token<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl iceberg_catalog_rest::TokenProvider for iceberg_catalog_rest::OAuth2TokenProvider +pub fn iceberg_catalog_rest::OAuth2TokenProvider::invalidate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::OAuth2TokenProvider::regenerate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::OAuth2TokenProvider::token<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl iceberg_catalog_rest::TokenProvider for iceberg_catalog_rest::StaticTokenProvider +pub fn iceberg_catalog_rest::StaticTokenProvider::invalidate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::StaticTokenProvider::regenerate<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg_catalog_rest::StaticTokenProvider::token<'life0, 'async_trait>(&'life0 self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait diff --git a/crates/iceberg/public-api.txt b/crates/iceberg/public-api.txt index 5959f10c74..91616c92b1 100644 --- a/crates/iceberg/public-api.txt +++ b/crates/iceberg/public-api.txt @@ -3083,6 +3083,7 @@ pub fn iceberg::transaction::Transaction::expire_snapshots(&self) -> iceberg::tr pub fn iceberg::transaction::Transaction::fast_append(&self) -> iceberg::transaction::append::FastAppendAction pub fn iceberg::transaction::Transaction::new(table: &iceberg::table::Table) -> Self pub fn iceberg::transaction::Transaction::replace_sort_order(&self) -> iceberg::transaction::sort_order::ReplaceSortOrderAction +pub fn iceberg::transaction::Transaction::row_delta(&self) -> iceberg::transaction::row_delta::RowDeltaAction pub fn iceberg::transaction::Transaction::update_location(&self) -> iceberg::transaction::update_location::UpdateLocationAction pub fn iceberg::transaction::Transaction::update_schema(&self) -> iceberg::transaction::update_schema::UpdateSchemaAction pub fn iceberg::transaction::Transaction::update_statistics(&self) -> iceberg::transaction::update_statistics::UpdateStatisticsAction @@ -3112,6 +3113,7 @@ pub struct iceberg::writer::base_writer::data_file_writer::DataFileWriter iceberg::writer::CurrentFileStatus for iceberg::writer::base_writer::data_file_writer::DataFileWriter where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_file_path(&self) -> alloc::string::String pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_row_num(&self) -> usize +pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_schema(&self) -> iceberg::spec::SchemaRef pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_written_size(&self) -> usize impl iceberg::writer::IcebergWriter for iceberg::writer::base_writer::data_file_writer::DataFileWriter where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::close<'life0, 'async_trait>(&'life0 mut self) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait @@ -3145,8 +3147,57 @@ pub struct iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteW impl iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig pub fn iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig::new(equality_ids: alloc::vec::Vec, original_schema: iceberg::spec::SchemaRef) -> iceberg::Result pub fn iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig::projected_arrow_schema_ref(&self) -> &arrow_schema::schema::SchemaRef +impl core::clone::Clone for iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig +pub fn iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig::clone(&self) -> iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig impl core::fmt::Debug for iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig pub fn iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteWriterConfig::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub mod iceberg::writer::base_writer::position_delete_writer +pub struct iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter +impl iceberg::writer::IcebergWriter for iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter::close<'life0, 'async_trait>(&'life0 mut self) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter::write<'life0, 'async_trait>(&'life0 mut self, batch: arrow_array::record_batch::RecordBatch) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl core::fmt::Debug for iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder +impl iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder::new(inner: iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder, config: iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig) -> Self +impl iceberg::writer::IcebergWriterBuilder for iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator +pub type iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder::R = iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder::build<'life0, 'async_trait>(&'life0 self, partition_key: core::option::Option) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl core::fmt::Debug for iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig +impl iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig::arrow_schema() -> arrow_schema::schema::SchemaRef +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig::new(partition_value: core::option::Option, partition_spec_id: i32, referenced_data_file: core::option::Option) -> Self +impl core::clone::Clone for iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig::clone(&self) -> iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig +impl core::fmt::Debug for iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteWriterConfig::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub mod iceberg::writer::combined_writer +pub mod iceberg::writer::combined_writer::delta_writer +pub struct iceberg::writer::combined_writer::delta_writer::DeltaWriter +pub iceberg::writer::combined_writer::delta_writer::DeltaWriter::data_writer: DW +pub iceberg::writer::combined_writer::delta_writer::DeltaWriter::eq_delete_writer: EDW +pub iceberg::writer::combined_writer::delta_writer::DeltaWriter::pos_delete_writer: PDW +pub iceberg::writer::combined_writer::delta_writer::DeltaWriter::seen_rows: std::collections::hash::map::HashMap +pub iceberg::writer::combined_writer::delta_writer::DeltaWriter::unique_cols: alloc::vec::Vec +impl iceberg::writer::IcebergWriter for iceberg::writer::combined_writer::delta_writer::DeltaWriter where DW: iceberg::writer::IcebergWriter + iceberg::writer::CurrentFileStatus, PDW: iceberg::writer::IcebergWriter, EDW: iceberg::writer::IcebergWriter +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriter::close<'life0, 'async_trait>(&'life0 mut self) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriter::write<'life0, 'async_trait>(&'life0 mut self, batch: arrow_array::record_batch::RecordBatch) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub struct iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder +impl iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::new(data_writer_builder: DWB, pos_delete_writer_builder: PDWB, eq_delete_writer_builder: EDWB, unique_cols: alloc::vec::Vec) -> Self +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::with_max_seen_rows(self, max_seen_rows: usize) -> Self +impl iceberg::writer::IcebergWriterBuilder for iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder where DWB: iceberg::writer::IcebergWriterBuilder, PDWB: iceberg::writer::IcebergWriterBuilder, EDWB: iceberg::writer::IcebergWriterBuilder, ::R: iceberg::writer::CurrentFileStatus +pub type iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::R = iceberg::writer::combined_writer::delta_writer::DeltaWriter<::R, ::R, ::R> +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::build<'life0, 'async_trait>(&'life0 self, partition_key: core::option::Option) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl core::clone::Clone for iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::clone(&self) -> iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder +impl core::fmt::Debug for iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct iceberg::writer::combined_writer::delta_writer::Position +pub const iceberg::writer::combined_writer::delta_writer::DEFAULT_MAX_SEEN_ROWS: usize pub mod iceberg::writer::file_writer pub mod iceberg::writer::file_writer::location_generator pub struct iceberg::writer::file_writer::location_generator::DefaultFileNameGenerator @@ -3186,12 +3237,13 @@ pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter: impl iceberg::writer::CurrentFileStatus for iceberg::writer::file_writer::rolling_writer::RollingFileWriter pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_file_path(&self) -> alloc::string::String pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_row_num(&self) -> usize +pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_schema(&self) -> iceberg::spec::SchemaRef pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_written_size(&self) -> usize pub struct iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder impl iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder::build(&self) -> iceberg::writer::file_writer::rolling_writer::RollingFileWriter -pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder::new(inner_builder: B, target_file_size: usize, file_io: iceberg::io::FileIO, location_generator: L, file_name_generator: F) -> Self -pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder::new_with_default_file_size(inner_builder: B, file_io: iceberg::io::FileIO, location_generator: L, file_name_generator: F) -> Self +pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder::new(inner_builder: B, schema: iceberg::spec::SchemaRef, target_file_size: usize, file_io: iceberg::io::FileIO, location_generator: L, file_name_generator: F) -> Self +pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder::new_with_default_file_size(inner_builder: B, schema: iceberg::spec::SchemaRef, file_io: iceberg::io::FileIO, location_generator: L, file_name_generator: F) -> Self impl core::clone::Clone for iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder::clone(&self) -> iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder impl core::fmt::Debug for iceberg::writer::file_writer::rolling_writer::RollingFileWriterBuilder @@ -3200,6 +3252,7 @@ pub struct iceberg::writer::file_writer::ParquetWriter impl iceberg::writer::CurrentFileStatus for iceberg::writer::file_writer::ParquetWriter pub fn iceberg::writer::file_writer::ParquetWriter::current_file_path(&self) -> alloc::string::String pub fn iceberg::writer::file_writer::ParquetWriter::current_row_num(&self) -> usize +pub fn iceberg::writer::file_writer::ParquetWriter::current_schema(&self) -> iceberg::spec::SchemaRef pub fn iceberg::writer::file_writer::ParquetWriter::current_written_size(&self) -> usize impl iceberg::writer::file_writer::FileWriter for iceberg::writer::file_writer::ParquetWriter pub async fn iceberg::writer::file_writer::ParquetWriter::close(self) -> iceberg::Result> @@ -3209,6 +3262,7 @@ impl iceberg::writer::file_writer::ParquetWriterBuilder pub fn iceberg::writer::file_writer::ParquetWriterBuilder::from_table_properties(table_props: &iceberg::spec::TableProperties, schema: iceberg::spec::SchemaRef) -> Self pub fn iceberg::writer::file_writer::ParquetWriterBuilder::new(props: parquet::file::properties::WriterProperties, schema: iceberg::spec::SchemaRef) -> Self pub fn iceberg::writer::file_writer::ParquetWriterBuilder::new_with_match_mode(props: parquet::file::properties::WriterProperties, schema: iceberg::spec::SchemaRef, match_mode: iceberg::arrow::FieldMatchMode) -> Self +pub fn iceberg::writer::file_writer::ParquetWriterBuilder::with_arrow_schema(self, arrow_schema: alloc::sync::Arc) -> iceberg::Result pub fn iceberg::writer::file_writer::ParquetWriterBuilder::with_match_mode(self, match_mode: iceberg::arrow::FieldMatchMode) -> Self impl core::clone::Clone for iceberg::writer::file_writer::ParquetWriterBuilder pub fn iceberg::writer::file_writer::ParquetWriterBuilder::clone(&self) -> iceberg::writer::file_writer::ParquetWriterBuilder @@ -3262,18 +3316,22 @@ pub fn iceberg::writer::partitioning::fanout_writer::FanoutWriter::writ pub trait iceberg::writer::CurrentFileStatus pub fn iceberg::writer::CurrentFileStatus::current_file_path(&self) -> alloc::string::String pub fn iceberg::writer::CurrentFileStatus::current_row_num(&self) -> usize +pub fn iceberg::writer::CurrentFileStatus::current_schema(&self) -> iceberg::spec::SchemaRef pub fn iceberg::writer::CurrentFileStatus::current_written_size(&self) -> usize impl iceberg::writer::CurrentFileStatus for iceberg::writer::file_writer::ParquetWriter pub fn iceberg::writer::file_writer::ParquetWriter::current_file_path(&self) -> alloc::string::String pub fn iceberg::writer::file_writer::ParquetWriter::current_row_num(&self) -> usize +pub fn iceberg::writer::file_writer::ParquetWriter::current_schema(&self) -> iceberg::spec::SchemaRef pub fn iceberg::writer::file_writer::ParquetWriter::current_written_size(&self) -> usize impl iceberg::writer::CurrentFileStatus for iceberg::writer::base_writer::data_file_writer::DataFileWriter where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_file_path(&self) -> alloc::string::String pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_row_num(&self) -> usize +pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_schema(&self) -> iceberg::spec::SchemaRef pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter::current_written_size(&self) -> usize impl iceberg::writer::CurrentFileStatus for iceberg::writer::file_writer::rolling_writer::RollingFileWriter pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_file_path(&self) -> alloc::string::String pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_row_num(&self) -> usize +pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_schema(&self) -> iceberg::spec::SchemaRef pub fn iceberg::writer::file_writer::rolling_writer::RollingFileWriter::current_written_size(&self) -> usize pub trait iceberg::writer::IcebergWriter: core::marker::Send + 'static pub fn iceberg::writer::IcebergWriter::close<'life0, 'async_trait>(&'life0 mut self) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait @@ -3284,6 +3342,12 @@ pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriter:: impl iceberg::writer::IcebergWriter for iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteFileWriter where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator pub fn iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteFileWriter::close<'life0, 'async_trait>(&'life0 mut self) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait pub fn iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteFileWriter::write<'life0, 'async_trait>(&'life0 mut self, batch: arrow_array::record_batch::RecordBatch) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl iceberg::writer::IcebergWriter for iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter::close<'life0, 'async_trait>(&'life0 mut self) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter::write<'life0, 'async_trait>(&'life0 mut self, batch: arrow_array::record_batch::RecordBatch) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl iceberg::writer::IcebergWriter for iceberg::writer::combined_writer::delta_writer::DeltaWriter where DW: iceberg::writer::IcebergWriter + iceberg::writer::CurrentFileStatus, PDW: iceberg::writer::IcebergWriter, EDW: iceberg::writer::IcebergWriter +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriter::close<'life0, 'async_trait>(&'life0 mut self) -> core::pin::Pin>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriter::write<'life0, 'async_trait>(&'life0 mut self, batch: arrow_array::record_batch::RecordBatch) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait pub trait iceberg::writer::IcebergWriterBuilder: core::marker::Send + core::marker::Sync + 'static pub type iceberg::writer::IcebergWriterBuilder::R: iceberg::writer::IcebergWriter pub fn iceberg::writer::IcebergWriterBuilder::build<'life0, 'async_trait>(&'life0 self, partition_key: core::option::Option) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait @@ -3293,6 +3357,12 @@ pub fn iceberg::writer::base_writer::data_file_writer::DataFileWriterBuilder iceberg::writer::IcebergWriterBuilder for iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteFileWriterBuilder where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator pub type iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteFileWriterBuilder::R = iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteFileWriter pub fn iceberg::writer::base_writer::equality_delete_writer::EqualityDeleteFileWriterBuilder::build<'life0, 'async_trait>(&'life0 self, partition_key: core::option::Option) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl iceberg::writer::IcebergWriterBuilder for iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder where B: iceberg::writer::file_writer::FileWriterBuilder, L: iceberg::writer::file_writer::location_generator::LocationGenerator, F: iceberg::writer::file_writer::location_generator::FileNameGenerator +pub type iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder::R = iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriter +pub fn iceberg::writer::base_writer::position_delete_writer::PositionDeleteFileWriterBuilder::build<'life0, 'async_trait>(&'life0 self, partition_key: core::option::Option) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait +impl iceberg::writer::IcebergWriterBuilder for iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder where DWB: iceberg::writer::IcebergWriterBuilder, PDWB: iceberg::writer::IcebergWriterBuilder, EDWB: iceberg::writer::IcebergWriterBuilder, ::R: iceberg::writer::CurrentFileStatus +pub type iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::R = iceberg::writer::combined_writer::delta_writer::DeltaWriter<::R, ::R, ::R> +pub fn iceberg::writer::combined_writer::delta_writer::DeltaWriterBuilder::build<'life0, 'async_trait>(&'life0 self, partition_key: core::option::Option) -> core::pin::Pin> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait pub macro iceberg::ensure_data_valid! #[non_exhaustive] pub enum iceberg::ErrorKind pub iceberg::ErrorKind::CatalogCommitConflicts From b5120491020f8dc08f026a55655d5b01d307cc5f Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Sun, 23 Aug 2026 18:49:14 -0400 Subject: [PATCH 09/10] Add oauth2 token refreshing on expiry --- crates/catalog/rest/public-api.txt | 2 +- crates/catalog/rest/src/catalog.rs | 2 -- crates/catalog/rest/src/token.rs | 28 +++++++++++++++++++++------- 3 files changed, 22 insertions(+), 10 deletions(-) diff --git a/crates/catalog/rest/public-api.txt b/crates/catalog/rest/public-api.txt index 84e8f043d1..3331c292e5 100644 --- a/crates/catalog/rest/public-api.txt +++ b/crates/catalog/rest/public-api.txt @@ -171,7 +171,7 @@ impl<'de> serde_core::de::Deserialize<'de> for iceberg_catalog_rest::NamespaceRe pub fn iceberg_catalog_rest::NamespaceResponse::deserialize<__D>(__deserializer: __D) -> core::result::Result::Error> where __D: serde_core::de::Deserializer<'de> pub struct iceberg_catalog_rest::OAuth2TokenProvider impl iceberg_catalog_rest::OAuth2TokenProvider -pub fn iceberg_catalog_rest::OAuth2TokenProvider::new(client: reqwest::async_impl::client::Client, client_id: core::option::Option, client_secret: alloc::string::String, token_endpoint: alloc::string::String, extra_headers: http::header::map::HeaderMap, extra_oauth_params: std::collections::hash::map::HashMap, cached_token: core::option::Option) -> Self +pub fn iceberg_catalog_rest::OAuth2TokenProvider::new(client: reqwest::async_impl::client::Client, client_id: core::option::Option, client_secret: alloc::string::String, token_endpoint: alloc::string::String, extra_headers: http::header::map::HeaderMap, extra_oauth_params: std::collections::hash::map::HashMap) -> Self impl core::fmt::Debug for iceberg_catalog_rest::OAuth2TokenProvider pub fn iceberg_catalog_rest::OAuth2TokenProvider::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result impl iceberg_catalog_rest::TokenProvider for iceberg_catalog_rest::OAuth2TokenProvider diff --git a/crates/catalog/rest/src/catalog.rs b/crates/catalog/rest/src/catalog.rs index d8b3ec6dbd..29e2678cd2 100644 --- a/crates/catalog/rest/src/catalog.rs +++ b/crates/catalog/rest/src/catalog.rs @@ -244,8 +244,6 @@ impl AuthenticatorConfig { self.token_endpoint.clone(), extra_headers, self.extra_oauth_params.clone(), - // If there's a preexisting token, seed the cache with that token. - self.token.clone(), )); return Some(Arc::new(BearerTokenAuthenticator::new(provider))); } diff --git a/crates/catalog/rest/src/token.rs b/crates/catalog/rest/src/token.rs index c1a8685728..e7dbf094d7 100644 --- a/crates/catalog/rest/src/token.rs +++ b/crates/catalog/rest/src/token.rs @@ -168,6 +168,10 @@ pub struct OAuth2TokenProvider { /// Most recently fetched token. cached_token: Mutex>, + + /// Expiration time of the cached token, if known. If `None`, we don't know when it expires. + /// If the token is expired, we will fetch a new one. + cached_token_expiration: Mutex>, } impl Debug for OAuth2TokenProvider { @@ -199,7 +203,6 @@ impl OAuth2TokenProvider { token_endpoint: String, extra_headers: HeaderMap, extra_oauth_params: HashMap, - cached_token: Option, ) -> Self { Self { client, @@ -208,12 +211,13 @@ impl OAuth2TokenProvider { token_endpoint, extra_headers, extra_oauth_params, - cached_token: Mutex::new(cached_token), + cached_token: Mutex::new(None), + cached_token_expiration: Mutex::new(None), } } /// Just fetch a token. Don't store it or do anything else. - async fn exchange_credential_for_token(&self) -> Result { + async fn exchange_credential_for_token(&self) -> Result<(String, Option)> { let mut params = HashMap::with_capacity(4); params.insert("grant_type", "client_credentials"); if let Some(client_id) = &self.client_id { @@ -273,7 +277,7 @@ impl OAuth2TokenProvider { })?; Err(Error::from(e)) }?; - Ok(auth_res.access_token) + Ok((auth_res.access_token, auth_res.expires_in)) } } @@ -282,25 +286,35 @@ impl TokenProvider for OAuth2TokenProvider { /// Fetch a new token if we don't already have one cached. async fn token(&self) -> Result { let mut cached = self.cached_token.lock().await; + let mut cached_expiration = self.cached_token_expiration.lock().await; - if let Some(token) = cached.clone() { + if let Some(token) = cached.clone() + && cached_expiration.is_none_or(|exp| std::time::Instant::now() < exp) + { return Ok(token); } - let token = self.exchange_credential_for_token().await?; + let (token, expires_in) = self.exchange_credential_for_token().await?; *cached = Some(token.clone()); + *cached_expiration = + expires_in.map(|secs| std::time::Instant::now() + std::time::Duration::from_secs(secs)); Ok(token) } async fn invalidate(&self) -> Result<()> { *self.cached_token.lock().await = None; + *self.cached_token_expiration.lock().await = None; Ok(()) } /// Try fetching and caching a new token. If that fails, keep the old token. async fn regenerate(&self) -> Result<()> { let mut cached = self.cached_token.lock().await; - *cached = Some(self.exchange_credential_for_token().await?); + let mut cached_expiration = self.cached_token_expiration.lock().await; + let (token, expires_in) = self.exchange_credential_for_token().await?; + *cached = Some(token); + *cached_expiration = + expires_in.map(|secs| std::time::Instant::now() + std::time::Duration::from_secs(secs)); Ok(()) } } From 7e2d814da13e2c75176b9cb6e5fe13faf9d07594 Mon Sep 17 00:00:00 2001 From: Patrick Butler Date: Thu, 27 Aug 2026 16:36:17 -0400 Subject: [PATCH 10/10] use checked add for token expiry --- crates/catalog/rest/src/token.rs | 42 +++++++++++++++++++++++++++----- 1 file changed, 36 insertions(+), 6 deletions(-) diff --git a/crates/catalog/rest/src/token.rs b/crates/catalog/rest/src/token.rs index e7dbf094d7..380ee9fe89 100644 --- a/crates/catalog/rest/src/token.rs +++ b/crates/catalog/rest/src/token.rs @@ -27,6 +27,7 @@ use std::collections::HashMap; use std::fmt::Debug; use std::sync::Arc; +use std::time::{Duration, Instant}; use async_trait::async_trait; use http::{Method, StatusCode}; @@ -171,7 +172,7 @@ pub struct OAuth2TokenProvider { /// Expiration time of the cached token, if known. If `None`, we don't know when it expires. /// If the token is expired, we will fetch a new one. - cached_token_expiration: Mutex>, + cached_token_expiration: Mutex>, } impl Debug for OAuth2TokenProvider { @@ -281,6 +282,12 @@ impl OAuth2TokenProvider { } } +/// How far ahead of expiry a cached token is replaced. +/// +/// A token handed out at the very edge of its lifetime can expire while the request carrying it +/// is still in flight, so it is retired early instead. +const REFRESH_MARGIN: Duration = Duration::from_secs(600); + #[async_trait] impl TokenProvider for OAuth2TokenProvider { /// Fetch a new token if we don't already have one cached. @@ -289,15 +296,17 @@ impl TokenProvider for OAuth2TokenProvider { let mut cached_expiration = self.cached_token_expiration.lock().await; if let Some(token) = cached.clone() - && cached_expiration.is_none_or(|exp| std::time::Instant::now() < exp) + && cached_expiration.is_none_or(|exp| Instant::now() + REFRESH_MARGIN < exp) { return Ok(token); } let (token, expires_in) = self.exchange_credential_for_token().await?; + // Resolve the expiry before touching the cache, so a failure here leaves the previous + // token and its expiry paired rather than storing a new token against a stale expiry. + let expiration = expiration_instant(expires_in)?; *cached = Some(token.clone()); - *cached_expiration = - expires_in.map(|secs| std::time::Instant::now() + std::time::Duration::from_secs(secs)); + *cached_expiration = expiration; Ok(token) } @@ -312,9 +321,30 @@ impl TokenProvider for OAuth2TokenProvider { let mut cached = self.cached_token.lock().await; let mut cached_expiration = self.cached_token_expiration.lock().await; let (token, expires_in) = self.exchange_credential_for_token().await?; + let expiration = expiration_instant(expires_in)?; *cached = Some(token); - *cached_expiration = - expires_in.map(|secs| std::time::Instant::now() + std::time::Duration::from_secs(secs)); + *cached_expiration = expiration; Ok(()) } } + +/// Converts an OAuth2 `expires_in` lifetime into the instant the token goes stale. +/// +/// `None` means the server reported no lifetime, which callers read as "never expires". +/// A lifetime so large that it overflows the monotonic clock is an error rather than `None`, +/// since silently treating it as non-expiring would keep a dead token forever. +fn expiration_instant(expires_in: Option) -> Result> { + expires_in + .map(|secs| { + Instant::now() + .checked_add(Duration::from_secs(secs)) + .ok_or_else(|| { + Error::new( + ErrorKind::Unexpected, + "OAuth2 token lifetime overflows the monotonic clock", + ) + .with_context("expires_in", secs.to_string()) + }) + }) + .transpose() +}