Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
323 changes: 106 additions & 217 deletions Cargo.lock

Large diffs are not rendered by default.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ mockito = "1"
motore-macros = "0.4.3"
murmur3 = "0.5.2"
once_cell = "1.20"
opendal = "0.58"
opendal = "0.59"
ordered-float = "4"
parquet = "59.2"
pilota = "0.11.10"
Expand Down
2 changes: 2 additions & 0 deletions crates/catalog/rest/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -35,13 +35,15 @@ chrono = { workspace = true }
http = { workspace = true }
iceberg = { workspace = true }
itertools = { workspace = true }
rand = { workspace = true }
Comment thread
zakariya-s marked this conversation as resolved.
reqwest = { workspace = true }
serde = { workspace = true }
serde_derive = { workspace = true }
serde_json = { workspace = true }
tokio = { workspace = true }
tracing = { workspace = true }
typed-builder = { workspace = true }
typetag = { workspace = true }
uuid = { workspace = true, features = ["v4"] }

[dev-dependencies]
Expand Down
19 changes: 19 additions & 0 deletions crates/catalog/rest/public-api.txt
Original file line number Diff line number Diff line change
Expand Up @@ -181,6 +181,20 @@ impl serde_core::ser::Serialize for iceberg_catalog_rest::ListTablesResponse
pub fn iceberg_catalog_rest::ListTablesResponse::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::ListTablesResponse
pub fn iceberg_catalog_rest::ListTablesResponse::deserialize<__D>(__deserializer: __D) -> core::result::Result<Self, <__D as serde_core::de::Deserializer>::Error> where __D: serde_core::de::Deserializer<'de>
pub struct iceberg_catalog_rest::LoadCredentialsResponse
pub iceberg_catalog_rest::LoadCredentialsResponse::storage_credentials: alloc::vec::Vec<iceberg_catalog_rest::StorageCredential>
impl core::clone::Clone for iceberg_catalog_rest::LoadCredentialsResponse
pub fn iceberg_catalog_rest::LoadCredentialsResponse::clone(&self) -> iceberg_catalog_rest::LoadCredentialsResponse
impl core::cmp::Eq for iceberg_catalog_rest::LoadCredentialsResponse
impl core::cmp::PartialEq for iceberg_catalog_rest::LoadCredentialsResponse
pub fn iceberg_catalog_rest::LoadCredentialsResponse::eq(&self, other: &iceberg_catalog_rest::LoadCredentialsResponse) -> bool
impl core::fmt::Debug for iceberg_catalog_rest::LoadCredentialsResponse
pub fn iceberg_catalog_rest::LoadCredentialsResponse::fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result
impl core::marker::StructuralPartialEq for iceberg_catalog_rest::LoadCredentialsResponse
impl serde_core::ser::Serialize for iceberg_catalog_rest::LoadCredentialsResponse
pub fn iceberg_catalog_rest::LoadCredentialsResponse::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::LoadCredentialsResponse
pub fn iceberg_catalog_rest::LoadCredentialsResponse::deserialize<__D>(__deserializer: __D) -> core::result::Result<Self, <__D as serde_core::de::Deserializer>::Error> where __D: serde_core::de::Deserializer<'de>
pub struct iceberg_catalog_rest::LoadTableResult
pub iceberg_catalog_rest::LoadTableResult::config: std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>
pub iceberg_catalog_rest::LoadTableResult::metadata: iceberg::spec::table_metadata::TableMetadata
Expand Down Expand Up @@ -223,6 +237,7 @@ pub fn iceberg_catalog_rest::NoopAuthManager::fmt(&self, f: &mut core::fmt::Form
impl iceberg_catalog_rest::AuthManager for iceberg_catalog_rest::NoopAuthManager
pub fn iceberg_catalog_rest::NoopAuthManager::catalog_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::NoopAuthManager::init_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::boxed::Box<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::NoopAuthManager::table_session<'life0, 'life1, 'life2, 'life3, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _table: &'life2 iceberg::catalog::TableIdent, _props: &'life3 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>, parent: alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait
pub struct iceberg_catalog_rest::OAuth2Manager
impl iceberg_catalog_rest::OAuth2Manager
pub fn iceberg_catalog_rest::OAuth2Manager::new(token_endpoint: impl core::convert::Into<alloc::string::String>) -> Self
Expand All @@ -235,6 +250,7 @@ pub fn iceberg_catalog_rest::OAuth2Manager::fmt(&self, f: &mut core::fmt::Format
impl iceberg_catalog_rest::AuthManager for iceberg_catalog_rest::OAuth2Manager
pub fn iceberg_catalog_rest::OAuth2Manager::catalog_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, client: &'life1 iceberg_catalog_rest::HttpClient, props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::OAuth2Manager::init_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, client: &'life1 iceberg_catalog_rest::HttpClient, props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::boxed::Box<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::OAuth2Manager::table_session<'life0, 'life1, 'life2, 'life3, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _table: &'life2 iceberg::catalog::TableIdent, props: &'life3 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>, parent: alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: '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
Expand Down Expand Up @@ -386,11 +402,14 @@ pub const iceberg_catalog_rest::REST_CATALOG_PROP_WAREHOUSE: &str
pub trait iceberg_catalog_rest::AuthManager: core::fmt::Debug + core::marker::Send + core::marker::Sync
pub fn iceberg_catalog_rest::AuthManager::catalog_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, client: &'life1 iceberg_catalog_rest::HttpClient, props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::AuthManager::init_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, client: &'life1 iceberg_catalog_rest::HttpClient, props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::boxed::Box<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::AuthManager::table_session<'life0, 'life1, 'life2, 'life3, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _table: &'life2 iceberg::catalog::TableIdent, _props: &'life3 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>, parent: alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait
impl iceberg_catalog_rest::AuthManager for iceberg_catalog_rest::NoopAuthManager
pub fn iceberg_catalog_rest::NoopAuthManager::catalog_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::NoopAuthManager::init_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::boxed::Box<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::NoopAuthManager::table_session<'life0, 'life1, 'life2, 'life3, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _table: &'life2 iceberg::catalog::TableIdent, _props: &'life3 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>, parent: alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait
impl iceberg_catalog_rest::AuthManager for iceberg_catalog_rest::OAuth2Manager
pub fn iceberg_catalog_rest::OAuth2Manager::catalog_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, client: &'life1 iceberg_catalog_rest::HttpClient, props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::OAuth2Manager::init_session<'life0, 'life1, 'life2, 'async_trait>(&'life0 self, client: &'life1 iceberg_catalog_rest::HttpClient, props: &'life2 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::boxed::Box<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait
pub fn iceberg_catalog_rest::OAuth2Manager::table_session<'life0, 'life1, 'life2, 'life3, 'async_trait>(&'life0 self, _client: &'life1 iceberg_catalog_rest::HttpClient, _table: &'life2 iceberg::catalog::TableIdent, props: &'life3 std::collections::hash::map::HashMap<alloc::string::String, alloc::string::String>, parent: alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<alloc::sync::Arc<dyn iceberg_catalog_rest::AuthSession>>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait, 'life3: 'async_trait
pub trait iceberg_catalog_rest::AuthSession: core::fmt::Debug + core::marker::Send + core::marker::Sync
pub fn iceberg_catalog_rest::AuthSession::authenticate<'life0, 'life1, 'async_trait>(&'life0 self, request: &'life1 mut iceberg_catalog_rest::HttpRequest) -> core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = iceberg::error::Result<()>> + core::marker::Send + 'async_trait)>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait
51 changes: 47 additions & 4 deletions crates/catalog/rest/src/auth/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -25,9 +25,10 @@ use std::fmt::Debug;
use std::sync::Arc;

use async_trait::async_trait;
use iceberg::Result;
use iceberg::{Error, ErrorKind, Result, TableIdent};
pub use oauth2::OAuth2Manager;

use crate::catalog::{REST_CATALOG_PROP_AUTH_TYPE, RestCatalogConfig};
use crate::client::HttpClient;
use crate::request::HttpRequest;

Expand All @@ -36,6 +37,32 @@ pub const AUTH_TYPE_NONE: &str = "none";
/// `rest.auth.type` value selecting OAuth2 token authentication.
pub const AUTH_TYPE_OAUTH2: &str = "oauth2";

/// Builds the auth manager selected by the `rest.auth.type` configuration,
/// like Java's `AuthManagers.loadAuthManager`.
pub(crate) fn load_auth_manager(config: &RestCatalogConfig) -> Result<Arc<dyn AuthManager>> {
let auth_type = config.auth_type();
// Java parity (`AuthManagers`): make the inference visible so users
// configure the type explicitly.
if auth_type == AUTH_TYPE_OAUTH2 && !config.has_explicit_auth_type() {
tracing::warn!(
"Inferring {REST_CATALOG_PROP_AUTH_TYPE}={AUTH_TYPE_OAUTH2} from the configured \
OAuth properties; set it explicitly to avoid this warning"
);
}
match auth_type.as_str() {
AUTH_TYPE_NONE => Ok(Arc::new(NoopAuthManager)),
AUTH_TYPE_OAUTH2 => Ok(Arc::new(OAuth2Manager::from_config(config)?)),
other => Err(Error::new(
ErrorKind::DataInvalid,
format!(
"unknown '{REST_CATALOG_PROP_AUTH_TYPE}': {other}; use \
`RestSessionCatalogBuilder::with_auth_manager` or \
`RestCatalogBuilder::with_auth_manager` to inject a custom auth manager"
),
)),
}
}

/// Creates the [`AuthSession`]s used to authenticate REST catalog requests.
///
/// A manager is exclusively scoped to one catalog and must not be reused by
Expand All @@ -46,9 +73,9 @@ pub const AUTH_TYPE_OAUTH2: &str = "oauth2";
/// Catalog initialization calls [`AuthManager::catalog_session`] exactly once;
/// later sessions may rely on the state established by that call.
///
/// Both methods are handed the catalog's [`HttpClient`], which an
/// implementation may reuse for its own requests (e.g. a token exchange) so
/// that they share the catalog's connection pool and configuration.
/// Session-construction methods are handed the catalog's [`HttpClient`], which
/// an implementation may reuse for its own requests (e.g. a token exchange)
/// so that they share the catalog's connection pool and configuration.
#[async_trait]
pub trait AuthManager: Debug + Send + Sync {
/// Session used for the initial `/v1/config` handshake, given the
Expand All @@ -73,6 +100,22 @@ pub trait AuthManager: Debug + Send + Sync {
client: &HttpClient,
props: &HashMap<String, String>,
) -> Result<Arc<dyn AuthSession>>;

/// Returns a session for requests associated with `table`.
///
/// `props` are the unmerged properties returned by the table endpoint.
/// The default preserves the catalog session; managers should return a
/// child session only when the table properties contain an authentication
/// override.
async fn table_session(
&self,
_client: &HttpClient,
_table: &TableIdent,
_props: &HashMap<String, String>,
parent: Arc<dyn AuthSession>,
) -> Result<Arc<dyn AuthSession>> {
Ok(parent)
}
}

/// Authenticates outgoing REST catalog requests.
Expand Down
Loading
Loading