mbutrovich commented on code in PR #2976:
URL: https://github.com/apache/iceberg-rust/pull/2976#discussion_r4201140750
##########
crates/iceberg/src/io/storage/mod.rs:
##########
@@ -139,4 +140,179 @@ 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_credential_provider(
+ &self,
+ config: &StorageConfig,
+ credential_provider: Option<Arc<dyn StorageCredentialProvider>>,
+ ) -> Result<Arc<dyn Storage>> {
+ self.build(config)
+ }
Review Comment:
`FileIO` builds every storage through this method, so a factory that
implements only `build` drops the provider and logs nothing. A `RestCatalog`
accepts a user's factory through
[`CatalogBuilder::with_storage_factory`](https://github.com/apache/iceberg-rust/blob/f2409e9f5eaff1c643985013d927e17fdf16e6ec/crates/iceberg/src/catalog/mod.rs#L157),
for example a wrapper around `OpenDalStorageFactory` that adds metrics. As I
read it, once the catalog wiring lands, a wrapper like that would keep using
the initial vended credentials until they expire, and the first symptom would
be an authorization error from the object store. Could the default log a
warning once when it drops a provider, the way `serialize_all` does with
[`warn_once_without_provider`](https://github.com/apache/iceberg-rust/blob/e5e50f39960c320072b1dd852748d668ac766b11/crates/iceberg/src/io/file_io.rs#L101-L111)?
##########
crates/storage/opendal/src/lib.rs:
##########
@@ -1098,11 +1354,353 @@ mod tests {
}
}
+ /// Vends a credential for the whole `s3` scheme until it has served
+ /// `table-c`, then credentials scoped to each table. The scheme-wide
+ /// credential expires once the `table-c` batch is signed.
+ #[cfg(feature = "opendal-s3")]
+ #[derive(Debug, Default)]
+ struct ChangingScopeProvider {
+ scoped: AtomicBool,
+ scheme_wide_expired: AtomicBool,
+ }
+
+ #[cfg(feature = "opendal-s3")]
+ #[async_trait]
+ impl StorageCredentialProvider for ChangingScopeProvider {
+ fn supports_path(&self, _path: &str) -> bool {
+ true
+ }
+
+ async fn load_credential(&self, path: &str) ->
Result<StorageCredential> {
+ if path.starts_with("s3://bucket/table-c") {
+ self.scoped.store(true, Ordering::SeqCst);
+ }
+ // Signing the scoped batch means the scheme-wide credential is
gone.
+ if path == "s3://bucket/table-c" {
+ self.scheme_wide_expired.store(true, Ordering::SeqCst);
+ }
+ match path.strip_prefix("s3://bucket/") {
+ Some(rest) if self.scoped.load(Ordering::SeqCst) &&
!rest.is_empty() => {
+ let table = rest.split('/').next().unwrap_or_default();
+ Ok(s3_credential(format!("s3://bucket/{table}"),
"SCOPED_AK"))
+ }
+ _ if self.scheme_wide_expired.load(Ordering::SeqCst) =>
Err(Error::new(
+ ErrorKind::Unexpected,
+ "the scheme-wide credential expired",
+ )),
+ _ => Ok(s3_credential("s3", "SCHEME_AK")),
+ }
+ }
+ }
+
+ #[cfg(feature = "opendal-s3")]
+ #[tokio::test]
+ async fn test_delete_stream_signs_each_batch_with_its_vended_credential() {
+ let mut server = mockito::Server::new_async().await;
+ let mut delete = |path: &str, access_key: &str| {
+ server
+ .mock("DELETE", path)
+ .match_header(
+ "authorization",
+
mockito::Matcher::Regex(format!("Credential={access_key}/")),
+ )
+ .expect(1)
+ .with_status(204)
+ };
+ // The scheme-wide batch is flushed, with its own credential, before
+ // the scoped credential that replaces it is used.
+ let scheme_wide = delete("/bucket/table-b/f.parquet", "SCHEME_AK")
+ .create_async()
+ .await;
+ let scoped = delete("/bucket/table-c/f.parquet", "SCOPED_AK")
+ .create_async()
+ .await;
+
+ let mut config = S3Config::default();
+ config.endpoint = Some(server.url());
+ config.region = Some("us-east-1".to_string());
+ config.disable_config_load = true;
+ config.disable_ec2_metadata = true;
+ let storage = OpenDalStorage::S3 {
+ config: Arc::new(config),
+ customized_credential_load: None,
+ credential_provider:
Some(Arc::new(ChangingScopeProvider::default())),
+ client_config: OpenDalClientConfig::default(),
+ };
+ storage
+ .delete_stream(
+ futures::stream::iter([
+ "s3://bucket/table-b/f.parquet".to_string(),
+ "s3://bucket/table-c/f.parquet".to_string(),
+ ])
+ .boxed(),
+ )
+ .await
+ .unwrap();
+
+ scheme_wide.assert_async().await;
+ scoped.assert_async().await;
+ }
+
+ /// Vends `token` as a GCS credential, or fails when it is `None`.
+ #[cfg(feature = "opendal-gcs")]
+ #[derive(Debug)]
+ struct GcsTokenProvider(Option<&'static str>);
+
+ #[cfg(feature = "opendal-gcs")]
+ #[async_trait]
+ impl StorageCredentialProvider for GcsTokenProvider {
+ fn supports_path(&self, _path: &str) -> bool {
+ true
+ }
+
+ async fn load_credential(&self, _path: &str) ->
Result<StorageCredential> {
+ let token = self
+ .0
+ .ok_or_else(|| Error::new(ErrorKind::Unexpected, "refresh
failed"))?;
+ Ok(StorageCredential::new(
+ "gs",
+ HashMap::from([(GCS_TOKEN.to_string(), token.to_string())]),
+ ))
+ }
+ }
+
+ #[cfg(feature = "opendal-gcs")]
+ fn gcs_storage(server: &mockito::Server, provider: GcsTokenProvider) ->
OpenDalStorage {
+ let config = gcs_config_parse(HashMap::from([
+ (GCS_SERVICE_HOST.to_string(), server.url()),
+ (GCS_TOKEN.to_string(), "static-token".to_string()),
+ ]))
+ .unwrap();
+ OpenDalStorage::Gcs {
+ config: Arc::new(config),
+ credential_provider: Some(Arc::new(provider)),
+ client_config: OpenDalClientConfig::default(),
+ }
+ }
+
+ #[cfg(feature = "opendal-gcs")]
+ #[tokio::test]
+ async fn test_gcs_signs_with_the_vended_token() {
+ let mut server = mockito::Server::new_async().await;
+ let mock = server
+ .mock("GET", mockito::Matcher::Any)
+ .match_header("authorization", "Bearer vended-token")
+ .expect_at_least(1)
+ .with_status(404)
+ .create_async()
+ .await;
+
+ let storage = gcs_storage(&server,
GcsTokenProvider(Some("vended-token")));
+ assert!(!storage.exists("gs://bucket/file.parquet").await.unwrap());
+ mock.assert_async().await;
+ }
Review Comment:
Could you add a test where one operator signs a request, the credential
nears expiry, and the next request goes out with a newly loaded credential? The
tests here check the first signed request, and the adapter tests in `s3.rs` and
`gcs.rs` check that the expiry is parsed. None of them covers the reload
itself, which is the behavior this PR adds. For example, a provider could
return a different token on each call with an expiry 60 seconds out, and the
test could call `exists` twice on one operator from `create_operator`, with
mocks that expect the first token and then the second. reqsign (the
request-signing library OpenDAL uses) reloads inside its 120-second window, so
the two requests carry different tokens. #2931 plans an OpenDAL 0.59 bump next,
and a test like this would catch any change in how the newer reqsign caches
credentials.
##########
crates/iceberg/src/io/storage/mod.rs:
##########
@@ -139,4 +140,179 @@ 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_credential_provider(
+ &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.
+///
+/// # Debug output
+///
+/// [`FileIO`](crate::io::FileIO) and storage implementations print the
+/// provider in their `Debug` output, so the provider's `Debug` implementation
+/// must not expose credentials or other secrets.
+///
+/// # 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.
Review Comment:
How close to expiry counts as "near" here? The S3 and GCS backends pass the
credential to reqsign, the request-signing library OpenDAL uses. reqsign treats
a cached credential as stale 120 seconds before it expires
([AWS](https://docs.rs/reqsign-aws-core/3.2.0/src/reqsign_aws_core/credential.rs.html#49-51),
[GCS](https://docs.rs/reqsign-google/3.0.4/src/reqsign_google/credential.rs.html#248-250))
and calls the provider again on every request after that. It also rejects a
reloaded credential that expires within 10 seconds
([`validate_refreshed_credential`](https://docs.rs/reqsign-core/3.3.2/src/reqsign_core/signer.rs.html#159-176),
with the 10-second headroom for
[AWS](https://docs.rs/reqsign-aws-v4/3.3.1/src/reqsign_aws_v4/sign_request.rs.html#32)
and
[GCS](https://docs.rs/reqsign-google/3.0.4/src/reqsign_google/sign_request.rs.html#35)).
I checked this at the head commit. I built one GCS operator with
`create_operator`, used a provider that counts its calls and sets
`gcs.oauth2.token-expires-at`, and called `exists` three times on that operator:
- With a token that expires in an hour, the provider was called once.
- With a token that expires in 60 seconds, the provider was called three
times.
- With a token that expires in 5 seconds, every request failed with
`refreshed signing credential expires before the requested operation deadline`.
Could this section tell implementers to return a credential with more than
two minutes left, so they refresh earlier than that? A provider that refreshes
one minute before expiry gets a call on every request for a minute. A provider
that refreshes at expiry fails every request in the last 10 seconds.
--
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]