mbutrovich commented on code in PR #2932:
URL: https://github.com/apache/iceberg-rust/pull/2932#discussion_r4136989403


##########
crates/catalog/rest/src/catalog.rs:
##########
@@ -879,9 +873,29 @@ impl RestSessionCatalog {
                 )
             })?;
 
-        let file_io = FileIOBuilder::new(factory).with_props(props).build();
+        // If the catalog vends refreshable credentials for this table's 
storage,
+        // attach a provider so the backend re-fetches them before they expire.
+        // Only catalog authentication resolved from the properties can be
+        // rebuilt after FileIO serialization.
+        let credential_provider = build_vended_credential_provider(
+            &client.http_client,
+            client.auth_manager.as_ref(),
+            RestVendedCredentialProviderFactory::new(
+                &client.config.uri,
+                table.clone(),
+                table_config.unwrap_or_default(),
+            ),
+            &props,
+            self.auth_manager.is_none(),
+        )

Review Comment:
   `load_file_io` seeds the provider only from the flat `config` properties. At 
the head commit nothing reads `LoadTableResult::storage_credentials` (`grep -rn 
storage_credentials crates` only finds `types.rs`). The spec says "Clients must 
first check whether the respective credentials exist in the 
`storage-credentials` field before checking the `config` for credentials" 
([`LoadTableResult`](https://github.com/apache/iceberg/blob/48330b8dacab6662242d252b39c8444190979bb2/open-api/rest-catalog-open-api.yaml#L4249)).
 The provider already caches credentials per prefix, so should the initial 
`storage-credentials` seed that cache? As I read it, for a catalog that vends 
only through `storage-credentials`, the first file access has to fetch from the 
refresh endpoint. With refresh disabled, the table gets no credentials at all.
   
   #2651 also routes each path to the credential with the longest matching 
prefix, built from `storage-credentials`, but it uses a different mechanism 
(`FileIOBuilder::with_prefixed_props`). How do you see the two PRs fitting 
together? It would help to settle on one prefix-routing path before either 
lands.



##########
crates/iceberg/src/io/storage/mod.rs:
##########
@@ -139,4 +140,358 @@ pub trait StorageFactory: Debug + Send + Sync {
     /// A `Result` containing an `Arc<dyn Storage>` on success, or an error
     /// if the storage could not be created.
     fn build(&self, config: &StorageConfig) -> Result<Arc<dyn Storage>>;
+
+    /// Build a new Storage instance, optionally supplying a credential 
provider
+    /// that the backend can call to obtain and refresh short-lived 
credentials.
+    ///
+    /// Backends that cannot use the provider ignore it and use the credentials
+    /// in `config`, as they would without one. The default does exactly that.
+    #[allow(unused_variables)]
+    fn build_with_credentials(
+        &self,
+        config: &StorageConfig,
+        credential_provider: Option<Arc<dyn StorageCredentialProvider>>,
+    ) -> Result<Arc<dyn Storage>> {
+        self.build(config)
+    }
+}
+
+/// Supplies fresh, backend-specific storage credentials on demand.
+///
+/// A catalog that vends temporary credentials implements this trait so that
+/// storage backends can re-fetch credentials as they approach expiry instead
+/// of failing once the initial token's TTL runs out.
+///
+/// # Caching
+///
+/// [`load_credential`](Self::load_credential) may be called very frequently.
+/// Implementations must cache internally and only re-fetch when the current
+/// credential is at or near expiry; otherwise every object-store request could
+/// trigger a call back to the catalog.
+#[async_trait]
+pub trait StorageCredentialProvider: Debug + Send + Sync {
+    /// Return whether this provider has refresh configuration for `path`.
+    ///
+    /// Backends use this before replacing their normal credential chain. The
+    /// default is `true` for single-backend providers; multi-backend providers
+    /// should return `false` for schemes they do not configure.
+    fn supports_path(&self, _path: &str) -> bool {
+        true
+    }
+
+    /// Load a fresh credential for the storage location identified by `path`.
+    ///
+    /// `path` is the absolute location being accessed (e.g.
+    /// `s3://bucket/warehouse/db/table/...`). Providers that vend distinct
+    /// credentials per location prefix use it to select the most specific
+    /// match. When the selected credential has a declared
+    /// [`StorageCredential::prefix`], it must 
[cover](StorageCredential::covers) `path`.
+    async fn load_credential(&self, path: &str) -> Result<StorageCredential>;
+
+    /// Return a factory that rebuilds an equivalent provider in another 
process.
+    ///
+    /// [`FileIO::serialize_all`](crate::io::FileIO::serialize_all) serializes 
this
+    /// factory in place of the provider. The default reports that the provider
+    /// cannot be serialized.
+    fn factory(&self) -> Result<Arc<dyn StorageCredentialProviderFactory>> {
+        Err(Error::new(
+            ErrorKind::FeatureUnsupported,
+            "storage credential provider cannot be serialized",
+        ))
+    }
+}
+
+/// Serializable recipe that rebuilds a [`StorageCredentialProvider`] after
+/// [`FileIO`](crate::io::FileIO) deserialization.
+///
+/// Factories are serialized through [`typetag`](https://docs.rs/typetag), so
+/// implementations must use `#[typetag::serde]`, and the receiving binary must
+/// link the concrete implementation.
+#[typetag::serde(tag = "type")]
+pub trait StorageCredentialProviderFactory: Debug + Send + Sync {
+    /// Build a provider for a `FileIO` with the given storage configuration.
+    fn build(&self, config: &StorageConfig) -> Result<Arc<dyn 
StorageCredentialProvider>>;
+}
+
+/// A vended storage credential together with its scope and expiry.
+#[derive(Clone, Debug)]
+pub struct StorageCredential {
+    /// Storage-location prefix this credential is scoped to. `None` 
represents a
+    /// credential without a declared scope, sourced from flat storage 
properties.
+    prefix: Option<String>,
+    /// The backend-specific credential material.
+    kind: StorageCredentialKind,
+    /// When the credential expires, if known. `None` means non-expiring and
+    /// backends treat such a credential as always valid and never refresh it.
+    expires_at: Option<SystemTime>,
+}
+
+impl StorageCredential {
+    /// Create a storage credential with no declared scope or expiration.
+    pub fn new(kind: StorageCredentialKind) -> Self {
+        Self {
+            prefix: None,
+            kind,
+            expires_at: None,
+        }
+    }
+
+    /// Set the storage-location prefix this credential is scoped to.
+    pub fn with_prefix(mut self, prefix: impl Into<String>) -> Self {
+        self.prefix = Some(prefix.into());
+        self
+    }
+
+    /// Set when this credential expires.
+    pub fn with_expiration(mut self, expires_at: SystemTime) -> Self {
+        self.expires_at = Some(expires_at);
+        self
+    }
+
+    /// Return the storage-location prefix this credential is scoped to.
+    pub fn prefix(&self) -> Option<&str> {
+        self.prefix.as_deref()
+    }
+
+    /// Return whether this credential applies to `location`.
+    ///
+    /// A credential without a prefix covers every location. Otherwise the
+    /// prefix must match whole path segments of `location`, and scheme
+    /// aliases (`s3a`/`s3n` for `s3`, `gcs` for `gs`, and the plain-text
+    /// Azure schemes for their TLS variants) are treated as equal. A prefix
+    /// that is only a scheme, such as `s3`, covers every location with that
+    /// scheme.
+    pub fn covers(&self, location: &str) -> bool {
+        self.prefix
+            .as_deref()
+            .is_none_or(|prefix| storage_prefix_covers(prefix, location))
+    }
+
+    /// Return the backend-specific credential material.
+    pub fn kind(&self) -> &StorageCredentialKind {
+        &self.kind
+    }
+
+    /// Consume this credential and return its backend-specific material.
+    pub fn into_kind(self) -> StorageCredentialKind {
+        self.kind
+    }
+
+    /// Return when this credential expires.
+    pub fn expires_at(&self) -> Option<SystemTime> {
+        self.expires_at
+    }
+}
+
+/// Return whether the storage-location `prefix` covers `location`, with the
+/// matching rules of [`StorageCredential::covers`].
+pub fn storage_prefix_covers(prefix: &str, location: &str) -> bool {
+    let Some((location_scheme, location_rest)) = location.split_once("://") 
else {
+        return false;
+    };
+    let Some((prefix_scheme, prefix_rest)) = prefix.split_once("://") else {
+        return !prefix.is_empty() && canonical_scheme(prefix) == 
canonical_scheme(location_scheme);
+    };
+
+    canonical_scheme(prefix_scheme) == canonical_scheme(location_scheme)
+        && location_rest
+            .strip_prefix(prefix_rest)
+            .is_some_and(|remainder| {
+                prefix_rest.is_empty()
+                    || prefix_rest.ends_with('/')
+                    || remainder.is_empty()
+                    || remainder.starts_with('/')
+            })
+}

Review Comment:
   `storage_prefix_covers` only matches whole path segments, so a credential 
for `s3://bucket/tab` does not cover `s3://bucket/table/file.parquet`. The spec 
defines `prefix` as "a storage location prefix where the credential is 
relevant" and says clients should select the longest one 
([`StorageCredential`](https://github.com/apache/iceberg/blob/48330b8dacab6662242d252b39c8444190979bb2/open-api/rest-catalog-open-api.yaml#L3843-L3846)).
 It doesn't say whether a prefix has to end on a segment boundary. Java settles 
that with a plain string `startsWith` in 
[`S3FileIO.clientForStoragePath`](https://github.com/apache/iceberg/blob/48330b8dacab6662242d252b39c8444190979bb2/aws/src/main/java/org/apache/iceberg/aws/s3/S3FileIO.java#L387-L398),
 and it doesn't treat `s3a` and `s3` as the same scheme. With this code, a 
catalog that vends a prefix Java would accept gets a "does not cover" error 
here. Should this follow Java's matching? If the stricter rule is intentional, 
could the doc comment say 
 why?



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -399,17 +483,77 @@ impl OpenDalStorage {
     ///
     /// For most backends the URL host (bucket name) is sufficient. For HF the 
host
     /// encodes the repo type, not the repo identity, so a more specific key 
is used.
-    fn batch_key_for_path(&self, path: &str) -> String {
+    #[allow(unreachable_patterns)]
+    fn batch_key_for_path(&self, path: &str) -> Result<String> {
         match self {
             #[cfg(feature = "opendal-hf")]
-            OpenDalStorage::Hf { .. } => hf_batch_key(path),
-            _ => url::Url::parse(path)
+            OpenDalStorage::Hf { .. } => Ok(hf_batch_key(path)),
+            #[cfg(feature = "opendal-azdls")]
+            OpenDalStorage::Azdls { .. } => azdls_batch_key(path),
+            _ => Ok(url::Url::parse(path)
                 .ok()
                 .and_then(|u| u.host_str().map(|s| s.to_string()))
-                .unwrap_or_default(),
+                .unwrap_or_default()),
         }
     }
 
+    /// Return the dynamic credential provider that serves `path`, if any.
+    fn credential_provider_for_path(
+        &self,
+        path: &str,
+    ) -> Option<&Arc<dyn StorageCredentialProvider>> {
+        let provider: &Arc<dyn StorageCredentialProvider> = (match self {
+            #[cfg(feature = "opendal-s3")]
+            OpenDalStorage::S3 {
+                customized_credential_load: None,
+                credential_provider: Some(provider),
+                ..
+            } => Some(provider),
+            #[cfg(feature = "opendal-gcs")]
+            OpenDalStorage::Gcs {
+                credential_provider: Some(provider),
+                ..
+            } => Some(provider),
+            #[cfg(feature = "opendal-azdls")]
+            OpenDalStorage::Azdls {
+                credential_provider: Some(provider),
+                ..
+            } => Some(provider),
+            _ => None,
+        })?;
+        provider.supports_path(path).then_some(provider)
+    }
+
+    /// Returns a key that keeps bulk deletes within one operator and 
credential
+    /// scope. Loading the credential is normally a cache hit and avoids 
rebuilding
+    /// an operator for every path while preventing a batch from crossing 
prefixes.
+    async fn delete_batch_key_for_path(&self, path: &str) -> 
Result<DeleteBatchKey> {
+        let credential_location = match 
self.credential_provider_for_path(path) {
+            Some(provider) => {
+                let credential = provider.load_credential(path).await?;
+                if !credential.covers(path) {
+                    return Err(Error::new(
+                        ErrorKind::DataInvalid,
+                        format!(
+                            "vended credential prefix {:?} does not cover 
storage location {path:?}",
+                            credential.prefix()
+                        ),
+                    ));
+                }
+                Some(match credential.prefix() {
+                    Some(prefix) => prefix.to_string(),
+                    None => utils::storage_root(path)?,
+                })

Review Comment:
   When the credential has no prefix, the batch looks up credentials for the 
storage root (`s3://bucket/`) for as long as its `Deleter` lives. What happens 
with a catalog that returns the initial credential in `config`, with no prefix, 
and returns prefix-scoped entries from the refresh endpoint? I tried this 
against `RestVendedCredentialProvider` at the head commit. I seeded an unscoped 
credential that expires in 2 seconds and a refresh endpoint that returns a 
credential for `s3://bucket/table`. After the seed expired, 
`load_credential("s3://bucket/table/data/f.parquet")` returned the refreshed 
credential, but `load_credential("s3://bucket/")` failed with `no unexpired 
vended credential matches storage location: s3://bucket/`. The deleters are 
closed at the end of `delete_stream`, so a stream that runs past the seed's 
expiry would fail on flush even though every path in it is covered. Before 
expiry, each root lookup inside the refresh window also fetches, finds no entry 
for `s3://buc
 ket/`, and calls `record_failure`, which backs off refresh for the whole AWS 
cache.
   
   One option is to flush the root-keyed deleter as soon as a path in the 
stream resolves to a prefixed credential, because that means the unscoped 
credential is being replaced while it is still valid. Could you also add a test 
for an unscoped seed followed by a scoped refresh?



##########
crates/catalog/rest/src/catalog.rs:
##########
@@ -1470,7 +1478,7 @@ impl SessionCatalog for RestSessionCatalog {
         };
 
         let file_io = self
-            .load_file_io(Some(&response.metadata_location), None)
+            .load_file_io(commit.identifier(), 
Some(&response.metadata_location), None)

Review Comment:
   After a commit, the table returned here gets a `FileIO` built from the 
catalog properties only (`table_config` is `None`). So it has neither the 
vended credentials nor a refresh provider. `main` already behaves this way, but 
#2931 lists long-running writes as a motivation. A writer that keeps using the 
table returned by a commit loses refresh at that point. Should this reuse the 
`FileIO` from the table being committed? If that belongs in separate work, 
could you open a tracking issue and link it from #2931?



##########
crates/iceberg/src/io/file_io.rs:
##########
@@ -123,14 +132,28 @@ impl FileIO {
     ///
     /// All storage configuration properties are included in the serialized 
representation. These
     /// properties may contain credentials or other sensitive values, so the 
returned bytes must be
-    /// protected in transit and at rest by the application embedding this 
crate.
+    /// protected in transit and at rest by the application embedding this 
crate. A serialized
+    /// credential provider may likewise carry catalog authentication and 
vended credentials.
     ///
     /// Storage factories are serialized through 
[`typetag`](https://docs.rs/typetag). Third-party
     /// factories must use `#[typetag::serde]` on their [`StorageFactory`] 
implementation.
+    ///
+    /// A credential provider is serialized as the
+    /// [`StorageCredentialProviderFactory`] returned by
+    /// [`StorageCredentialProvider::factory`], and rebuilt on 
deserialization. Serialization fails
+    /// when the provider cannot be rebuilt in another process; use
+    /// [`FileIO::without_credential_provider`] to serialize without it.
     pub fn serialize_all(&self) -> Result<Vec<u8>> {
+        let credential_provider = self
+            .credential_provider
+            .as_ref()
+            .map(|provider| provider.factory())
+            .transpose()?;

Review Comment:
   With an injected `AuthManager`, a table whose config has a refresh endpoint 
now gets a `FileIO` that `serialize_all()` rejects with `FeatureUnsupported`. 
On `main` the same `FileIO` serializes, and 
`test_injected_auth_manager_file_io_serializes_without_provider_only` checks 
the new behavior. A caller that ships `FileIO` to workers would start failing 
after upgrading and would need a code change (`without_credential_provider`) to 
recover. Should `serialize_all` drop a provider that can't be rebuilt and 
serialize the static credentials, as `main` does, perhaps with a 
`tracing::warn!`? If failing is the behavior we want, could the PR description 
call it out as a behavior change?



##########
crates/storage/opendal/src/lib.rs:
##########
@@ -218,6 +233,7 @@ fn default_memory_operator() -> Operator {
 
 /// OpenDAL-based storage implementation.
 #[derive(Clone, Debug, Serialize, Deserialize)]
+#[non_exhaustive]

Review Comment:
   The new public fields on the `OpenDalStorage` variants 
(`credential_provider`, `sas_tokens`) and `#[non_exhaustive]` on the enum and 
its variants break downstream code that constructs or exhaustively matches 
`OpenDalStorage`. The repo marks breaking PRs with `!` in the title, as in 
#2838 (`feat!(rest): ...`). Could the title and description say this is 
breaking, or, once split, the PR that carries this change?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to