Skip to content

feat(rest): Support refreshing vended storage credentials - #2932

Open
zakariya-s wants to merge 15 commits into
apache:mainfrom
zakariya-s:feat/rest-vended-credential-refresh
Open

zakariya-s wants to merge 15 commits into
apache:mainfrom
zakariya-s:feat/rest-vended-credential-refresh

Conversation

@zakariya-s

@zakariya-s zakariya-s commented Jul 30, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

What changes are included in this PR?

This PR adds support for refreshing short-lived storage credentials vended by REST catalogs:

  • Adds a backend-independent StorageCredentialProvider interface to FileIO
  • Implements a REST credential provider for S3, GCS, and ADLS refresh endpoints
  • Uses table-scoped tokens and header.* properties for refresh requests
  • Caches credentials independently by cloud and selects credentials using the longest matching storage prefix
  • Adds jittered failure backoff while a cached credential remains valid
  • Adapts refreshed credentials into the OpenDAL/reqsign S3, GCS, and ADLS credential providers
  • Redacts credential-bearing configuration from Debug output

Azure credential refresh is not included because the current OpenDAL Azure backend does not expose the credential-provider and expiry hooks required for safe refresh.

ADLS refresh is now supported in OpenDAL 0.59.1 via the custom credential-provider.

Are these changes tested?

Yes

Comment thread crates/catalog/rest/src/catalog.rs Outdated
/// Disable header redaction in error logs (defaults to false for security)
pub const REST_CATALOG_PROP_DISABLE_HEADER_REDACTION: &str = "disable-header-redaction";
/// Identifier for a server-side scan plan associated with credential requests.
pub const REST_CATALOG_PROP_SCAN_PLAN_ID: &str = "rest.scan.plan-id";

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Server-side planning isn't supported yet so this will always be empty anyway. I don't think it's a problem to keep it until it eventually is supported

Comment thread crates/catalog/rest/src/catalog.rs Outdated
Comment on lines +1077 to +1084
let config = response
.config
.into_iter()
.chain(self.user_config.props.clone())
.collect();
let file_io = self
.load_file_io(Some(metadata_location), Some(config))
.await?;

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Comment thread crates/catalog/rest/src/client.rs
Comment thread crates/catalog/rest/Cargo.toml
Comment thread crates/iceberg/src/io/storage/mod.rs
Comment on lines +51 to +53
if let Some(no_auth) = m.remove(GCS_NO_AUTH)
&& is_truthy(no_auth.to_lowercase().as_str())
{

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Drive-by fix since this looked pretty bad. AWS did this correctly, but GCS would disable this even if gcs.no-auth was set to true

Comment on lines +340 to +345
let url = url::Url::parse(path).map_err(|e| {
Error::new(
ErrorKind::DataInvalid,
format!("Invalid gcs url: {path}: {e}"),
)
})?;

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

S3 had this validation before above but GCS didn't, so another drive-by fix

Comment thread crates/storage/opendal/src/lib.rs Outdated
/// `reqsign` [`Timestamp`](reqsign_core::time::Timestamp) used on backend
/// credential types (e.g. `AwsCredential::expires_in`, `google::Token::expires_at`).
#[cfg(any(feature = "opendal-s3", feature = "opendal-gcs"))]
pub(crate) fn system_time_to_timestamp(

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Didn't want to rely on directly on reqsign's Timestamp, but I'm happy to hear different opinions

/// It contains the location schemes it backs, the property keys it is configured
/// with, and how to parse its credential. The generic provider stays free of
/// any per-cloud knowledge.
struct CloudRefresh {

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This could also be made into a trait I guess, and we could split aws and gcp support into different modules, but I thought it was small enough to keep it as a struct and make consts for the different cloud providers

Comment thread crates/catalog/rest/src/credential.rs
Comment thread crates/catalog/rest/src/credential.rs
@mbutrovich
mbutrovich self-requested a review July 31, 2026 20:45
@zakariya-s

Copy link
Copy Markdown
Contributor Author

Hi @mbutrovich! Would it be possible to get a first-round review of this PR when you have time please?

@mbutrovich

Copy link
Copy Markdown
Collaborator

Hi @mbutrovich! Would it be possible to get a first-round review of this PR when you have time please?

Yep, it's in my queue! Just slammed with review requests :(

Thanks for your patience!

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

First pass, thanks for tackling this @zakariya-s. I have specific feedback, and also some ideas how we might break this up for other reviewers since this is a lot to review in one pass. A split along the layers already present in the diff looks fairly clean and, other than one ordering constraint, each piece is independently testable rather than a bare stub:

  1. The Debug-redaction hardening (RestCatalogConfig, HttpClient, StorageConfig, LoadTableResult, StorageCredential, OpenDalResolvingStorage, plus the is_sensitive_header broadening) and the two small unrelated fixes riding along (the register_table config merge, the GCS no-auth truthy parsing). None of this depends on credential refresh existing, and it already has its own tests.
  2. The StorageCredentialProvider trait and credential types in the iceberg crate, plus StorageFactory::build_with_credentials with its safe default (errors if a provider is supplied and the factory has not opted in). Covered by finding 5. No behavior change for existing backends.
  3. OpenDAL S3 consuming the trait (the adapter, the anonymous-access guard, the delete_stream change). Testable on its own with a hand-rolled provider, same as the tests already in this diff, without needing anything from the REST side.
  4. OpenDAL GCS consuming the trait, same shape as 3, independent of it.

Finding 1 (lost delete batching) sits in the shared uses_dynamic_credentials/delete_stream code in lib.rs rather than in s3.rs or gcs.rs specifically, so it isn't purely a 3-or-4 problem: whichever of the two lands first introduces that shared mechanism, and the other reuses it as-is, so the fix only needs to happen once.

  1. The REST catalog fetch/cache/jitter/backoff logic (RestVendedCredentialProvider, CloudRefresh, resolve_endpoint, the table-scoped client for_table). Covered by findings 2, 3, 4, and the root cause of 6 (the fresh-HttpClient-per-table-scoped-provider behavior lives in client.rs/credential.rs, both part of this PR). This can be written and tested against the trait directly with mockito, same as the tests already in this diff, without touching OpenDAL at all.
  2. Wiring RestCatalog::load_file_io to actually attach the provider (catalog.rs:561-565). Covered by finding 7 (the inline path at that call site) and the rest of finding 6 (this is the call site that turns "every load_file_io call re-runs the OAuth handshake" from a latent property of for_table into something that fires on every load_table/create_table/register_table).

This piece should land last on purpose, not just for tidiness: client.refresh-credentials-endpoint and gcs.oauth2.refresh-credentials-enabled are property keys the Java client and the REST spec already define, so a production catalog that's Java-interoperable may already be sending them today regardless of which client is asking. If the wiring lands before both the S3 and GCS sides can consume a provider, every existing rust client hitting such a catalog would start failing table loads on S3 or GCS the moment that PR merges, since the default build_with_credentials errors whenever a provider is supplied to a factory that hasn't opted in. So 3 and 4 both need to land before 6, even though 3, 4, and 5 are otherwise independent of each other and can be reviewed in any order.

#2931 is currently a single feature request rather than a tracking issue for a multi-PR stack. Worth turning it into one, or opening a separate tracking issue with a checklist for the pieces above, so reviewers can see the whole plan and where a given PR sits in it before reviewing any single piece.

Comment thread crates/storage/opendal/src/lib.rs Outdated
Comment on lines +423 to +431
fn uses_dynamic_credentials(&self, path: &str) -> bool {
match self {
#[cfg(feature = "opendal-s3")]
OpenDalStorage::S3 {
credential_provider: Some(provider),
..
} => provider.supports_path(path),
#[cfg(feature = "opendal-gcs")]
OpenDalStorage::Gcs {

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

crates/storage/opendal/src/lib.rs:423-431 (uses_dynamic_credentials) and :610-628 (delete_stream).

When a path is served by a credential provider, delete_stream skips the shared per-bucket Deleter and instead calls create_operator + a single op.delete(relative_path) per path, sequentially, inside the stream loop. The non-dynamic branch batches deletes through OpenDAL's Deleter (which can use bulk delete APIs); the dynamic branch does neither batching nor concurrency, and rebuilds the operator from scratch for every single file.

For expire_snapshots/purge on a table with vended-credential refresh enabled, this turns what would be a handful of batched multi-object delete calls into one HTTP round trip per file, plus an operator-construction cost per file. The code comment explains why deletes can't share a Deleter across different credential-prefix scopes (correctness: batch_key_for_path only groups by bucket, not by credential scope), but the fix taken forfeits batching entirely rather than partially, i.e. grouping deletes by (bucket, matched credential prefix) instead of just bucket would preserve batched delete within each credential-scope group. Was that considered?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Yeah I left this in the initial commit for sake of simplicity, we now batch by DeleteBatchKey which includes the credential scope so batching behaviour is now restored

Comment thread crates/catalog/rest/src/credential.rs Outdated
Comment on lines +265 to +300
let refreshed = self.fetch(configured).await.and_then(|entries| {
let credential = longest_prefix_match(&entries, path)
.filter(|entry| entry.is_unexpired(SystemTime::now()))
.map(|entry| entry.credential.clone())
.ok_or_else(|| {
Error::new(
ErrorKind::Unexpected,
format!("no unexpired vended credential matches storage location: {path}"),
)
})?;
Ok((entries, credential))
});

match refreshed {
Ok((entries, credential)) => {
let mut cache = configured.cache.lock().await;
cache.entries = entries;
cache.consecutive_failures = 0;
cache.retry_not_before = None;
Ok(credential)
}
Err(fetch_error) => {
let mut cache = configured.cache.lock().await;
cache.consecutive_failures = cache.consecutive_failures.saturating_add(1);
cache.retry_not_before =
Instant::now().checked_add(failure_backoff(cache.consecutive_failures));

// Graceful degradation: while the cached credential remains
// usable, serve it and retry after jittered backoff. Expired
// credentials are never served.
fallback
.filter(|entry| entry.is_unexpired(SystemTime::now()))
.map(|entry| entry.credential)
.ok_or(fetch_error)
}
}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

crates/catalog/rest/src/credential.rs:265-300 (refresh_credential), contrast with S3FileIO.refreshStorageCredentials()/GCSFileIO.refreshStorageCredentials() in Java.

Java's actual multi-prefix refresh (used by both S3FileIO and GCSFileIO, not the single-credential VendedCredentialsProvider/OAuth2RefreshCredentialsHandler used as an SDK credentials provider) is unconditional: on each scheduled refresh it fetches the credentials endpoint once, keeps every entry matching the cloud's root prefix, and replaces storageCredentials wholesale — no per-path filtering happens at refresh time at all.

The Rust refresh_credential instead does per-path filtering inline: it fetches all entries for a cloud, then immediately narrows to longest_prefix_match(&entries, path).filter(unexpired) for this specific call's path. If that narrowing yields nothing (no entry covers this path, or the covering entry happens to already be expired), the whole outcome is treated as Err and:

  • the freshly-fetched entries are never written to cache.entries (only the Ok((entries, credential)) branch at line 279-284 updates the cache) — so if the response contained entries for other prefixes (as the existing longest_prefix_match_ignores_freshness test exercises), they're thrown away even though a subsequent call for a different, valid path would have to fetch them all over again;
  • cache.consecutive_failures is incremented and retry_not_before backoff is armed (lines 288-290) even though the catalog responded successfully — it just didn't vend anything for this path.

Given Java's model of "cache everything the endpoint returns, unconditionally," was per-path filtering at refresh time (rather than only at cache-read time, where it already happens in cache_decision/longest_prefix_match) intentional here?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Good point that I missed, now all results are cached (minus the malformed or expired ones) and then it's filtered

Comment thread crates/catalog/rest/src/credential.rs Outdated
Comment on lines +238 to +249
parsed
.storage_credentials
.into_iter()
.filter(|sc| configured.cloud.matches_location(&sc.prefix))
.map(|sc| {
(configured.cloud.parse_credential)(&sc.config, Some(sc.prefix)).map(
|credential| {
CachedEntry::new(credential, configured.cloud.jitter_prefetch)
},
)
})
.collect()

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

collect() into Result<Vec<CachedEntry>> short-circuits on the first parse_credential error. If a server vends N credentials for one cloud and one entry is missing a required field, all N become unusable (and, per finding 2, this also counts as a "failure" for backoff purposes) rather than just the one bad entry. Java's equivalent (VendedCredentialsProvider.refreshCredential) only ever expects a single S3-prefixed entry and asserts on it directly, so there's no directly analogous "partial batch" behavior to compare against — but given this PR's own design supports N entries per cloud, is one bad entry meant to invalidate all the others?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

fetch() now returns a ParsedCredentials which splits between the successful creds and errors

Comment thread crates/catalog/rest/src/credential.rs Outdated
Comment on lines +441 to +443
let enabled = props
.get(cloud.enabled_key)
.is_none_or(|value| value.parse().unwrap_or(false));

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

str::parse::<bool>() only accepts the exact strings "true"/"false". Java's PropertyUtil.propertyAsBoolean uses Boolean.parseBoolean, which is case-insensitive for "true" ("True", "TRUE" all parse as true). A config value of client.refresh-credentials-enabled: "True" would enable refresh in Java but silently disable it in Rust (parse error → unwrap_or(false)). Low severity (fails closed either way), but worth a case-insensitive comparison to match Java's actual accepted input space.

Comment on lines +198 to +240
pub struct StorageCredential {
/// Storage-location prefix this credential is scoped to. `None` represents a
/// credential without a declared scope, sourced from flat storage properties.
pub prefix: Option<String>,
/// The backend-specific credential material.
pub 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.
pub expires_at: Option<SystemTime>,
}

/// Backend-specific credential material.
#[derive(Clone, Debug)]
pub enum StorageCredentialKind {
/// Amazon S3 credentials.
S3(S3Credential),
/// Google Cloud Storage credentials.
Gcs(GcsCredential),
}

/// Temporary Amazon S3 credentials.
#[derive(Clone)]
pub struct S3Credential {
/// AWS access key ID.
pub access_key_id: String,
/// AWS secret access key.
pub secret_access_key: String,
/// AWS session token, set for temporary (STS/vended) credentials.
pub session_token: Option<String>,
}

impl Debug for S3Credential {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("S3Credential").finish_non_exhaustive()
}
}

/// Temporary Google Cloud Storage credentials (an OAuth2 access token).
#[derive(Clone)]
pub struct GcsCredential {
/// OAuth2 bearer token used to access GCS.
pub token: String,
}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

These are new public types (in public-api.txt) that third-party StorageCredentialProvider implementors must construct by hand. Every field is pub, with no constructor. That's inconsistent with StorageConfig in the same module (crates/iceberg/src/io/storage/config/mod.rs:55-58), which keeps props private and exposes with_prop/from_props instead. Was a constructor considered, or is direct struct-literal construction intentional here?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry this wasn't intentional, they're no longer public

Comment thread crates/catalog/rest/src/catalog.rs Outdated
Comment thread crates/storage/opendal/src/lib.rs Outdated
Comment thread crates/catalog/rest/src/credential.rs
Comment thread crates/catalog/rest/src/client.rs
Comment thread crates/iceberg/src/io/storage/mod.rs
@zakariya-s

Copy link
Copy Markdown
Contributor Author

Thanks @mbutrovich for the first review! I've made some changes to address your comments and pushed them to this PR for posterity. I will follow your suggestion and split this PR and use a tracker issue.

@zakariya-s

Copy link
Copy Markdown
Contributor Author

First PR available here: #2976

Adapt table-scoped credential refresh to AuthManager and AuthSession.
@zakariya-s

Copy link
Copy Markdown
Contributor Author

Hey @CTTY, here's the full PR to see the broader picture. Happy if you wish to review the full PR if it's easier compared to splitting it.

@zakariya-s

Copy link
Copy Markdown
Contributor Author

apache/opendal#8030 has been included as part of opendal-core 0.59.0, so Azure support is now possible. Unfortunately, 0.59.0 doesn't compile due to a packaging issue so we will have to wait for a patch to upgrade the dependency.

@zakariya-s

Copy link
Copy Markdown
Contributor Author

@CTTY this now includes Azure support too, but as I mentioned before OpenDAL's latest release is broken so checks are failing. Hopefully 0.59.1 should release soon.

Update OpenDAL to 0.59.1 now that the fixed release is available.
@zakariya-s

Copy link
Copy Markdown
Contributor Author

0.59.1 was released, so Azure support should now be complete. This PR should be fully ready for review

@zakariya-s

Copy link
Copy Markdown
Contributor Author

Hey @mbutrovich @CTTY could I get a review here please? Thanks

Rebuild the refresh provider after FileIO serialization and fix several
correctness issues found in review.

- Serialize the credential provider as a StorageCredentialProviderFactory
  that rebuilds it from the FileIO properties and connects to the catalog
  lazily, like Java's VendedCredentialsProvider. Catalogs with an injected
  AuthManager refuse to serialize; FileIO::without_credential_provider
  drops the provider instead.
- Match credential prefixes on whole path segments and treat scheme
  aliases (s3a/s3n, gcs, abfs/wasb) as equal, via
  StorageCredential::covers.
- Merge refreshed credentials prefix by prefix, keeping an unexpired
  cached credential when the response has no valid replacement.
- Look up bulk-delete credentials by the batch scope, so a refresh during
  a delete can no longer abort it, and share one reqsign adapter across
  S3, GCS and ADLS.
- Cap the refresh buffer at half the credential lifetime, so short-lived
  credentials are not re-fetched on every operation.
- Ignore credential providers in storage factories that cannot use them
  instead of failing every operation, matching main and Java.
- Accept ADLS SAS tokens keyed by host (adls.sas-token.<host>) as Java and
  Polaris send them, and accept explicit http endpoints for abfs/wasb.
- Move auth manager selection to auth::load_auth_manager and auth property
  helpers onto RestCatalogConfig.
- Mark StorageCredentialKind and OpenDalStorage #[non_exhaustive].
@mbutrovich
mbutrovich self-requested a review September 29, 2026 16:41

@mbutrovich mbutrovich left a comment

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks for the updates @zakariya-s, and for turning #2931 into a tracking issue.

Since my last review the PR has grown to about 4,800 added lines across 26 files. It now also includes ADLS, FileIO serialization of providers, a new AuthManager::table_session method, and the OpenDAL 0.58 to 0.59 bump. I don't think a change this size can get through review here as one PR, even with it split into commits. #2976 has also drifted from this branch. As far as I can tell it has no StorageCredentialProviderFactory and no AzdlsCredential. Could we keep #2932 as a draft reference and land the work as a stack of smaller PRs, with the #2931 task list updated to match? Here is one way to cut it along the layers already in the diff:

  1. Debug redaction and the unrelated fixes: StorageConfig, LoadTableResult, StorageCredential, TokenResponse, the register_table config merge, and gcs.no-auth parsing. None of these depend on refresh.
  2. The core API in #2976, updated to match this branch: StorageCredentialProvider, StorageCredential, StorageCredentialKind with S3 and GCS, StorageFactory::build_with_credentials, FileIOBuilder::with_credential_provider, and the property constants. StorageCredentialKind is #[non_exhaustive], so the ADLS variant can be added later without a break.
  3. Provider serialization: StorageCredentialProvider::factory, StorageCredentialProviderFactory, and FileIO::without_credential_provider. This changes what FileIO::serialize_all does, so it deserves its own discussion.
  4. The OpenDAL S3 adapter and the shared pieces it introduces (VendedCredentialSource, DeleteBatchKey), tested with a hand-written provider like the tests already in this diff.
  5. The OpenDAL GCS adapter.
  6. The OpenDAL 0.59 bump as its own PR, then the ADLS storage pieces: AzdlsCredential, AzdlsSasTokens, and the ADLS adapter.
  7. REST auth: AuthManager::table_session, HttpClient::for_table, and the move of auth-manager selection into load_auth_manager. This adds a method to a public trait and doesn't touch storage.
  8. The REST provider in credential.rs for S3 only: fetch, cache, prefix selection, and backoff, tested against the trait with mockito. credential.rs alone is 2,452 lines (about 1,030 of code and 1,420 of tests), so starting with one cloud keeps this piece reviewable.
  9. GCS and ADLS in the REST provider, including the account-keyed seed strategy that ADLS needs.
  10. The wiring in RestSessionCatalog::load_file_io.

Everything depends on 2. Pieces 5 and 6 reuse the shared code from 4, 8 needs 2 and 7, and the wiring goes last. Some neighbors could be combined if a reviewer prefers, but each piece can be tested on its own. The comments below apply to whichever piece they end up in.

Comment on lines +880 to +890
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(),
)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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). 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.


let file_io = self
.load_file_io(Some(&response.metadata_location), None)
.load_file_io(commit.identifier(), Some(&response.metadata_location), None)

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Comment on lines 146 to +151
pub fn serialize_all(&self) -> Result<Vec<u8>> {
let credential_provider = self
.credential_provider
.as_ref()
.map(|provider| provider.factory())
.transpose()?;

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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?

Comment on lines +288 to +305
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('/')
})
}

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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). 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, 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?

Comment on lines +543 to +546
Some(match credential.prefix() {
Some(prefix) => prefix.to_string(),
None => utils::storage_root(path)?,
})

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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://bucket/, 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?


/// OpenDAL-based storage implementation.
#[derive(Clone, Debug, Serialize, Deserialize)]
#[non_exhaustive]

Copy link
Copy Markdown
Collaborator

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Support refreshing vended storage credentials for REST catalog tables

2 participants