Skip to content
Merged
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
Original file line number Diff line number Diff line change
Expand Up @@ -82,13 +82,13 @@ A Comet-side cache would have to either expose a tuning knob (TTL, max size, evi
- It is cached and registered like any other S3 store, one per key, so later scans on the bucket share it and the cache and registry keep one store per identity. As for any store, scans that miss the cache at the same moment each build one, and the last one cached wins.
- Each request is served by the store of the longest location covering its path, matched one segment at a time after percent-decoding, with the bucket root as an implicit location. Routing is per request, so a partition whose files span several locations reads each file with its own location's credential.
- A location's store is an `AmazonS3` whose bridge is bound to the location itself, so the vendor sees one stable path per location. The bridge is derived from the bucket's bridge, sharing its provider registration, so creating it calls no `ensureInitialized` and loads no classes on the Tokio worker that usually creates it. It is built on first use, usually inside an async read, from an `S3StoreTemplate` that resolved the region when the bucket's store was created, so building never blocks on the Tokio runtime.
- The locations are a snapshot. A 403 is either a real denial or a location added or removed since the snapshot, so the store fetches the locations again and retries the request once if its path now routes elsewhere. Every attempt, successful or not, starts a new snapshot generation, so requests routed from the same generation share one attempt and a failed attempt fails them all instead of each calling the provider in turn. The retry budget is per request, with no state that outlives it.
- The locations are a snapshot. A 403 is either a real denial or a location added or removed since the snapshot, and so is a failure to get a location's credential, because a provider with no policy for a path throws rather than vending a credential S3 would reject. On either, the store fetches the locations again and retries the request once if its path now routes elsewhere. The bridge gives its failures a `CredentialProviderError` source, which `object_store` passes through to the read unchanged, so the store can tell them from other errors. Every attempt, successful or not, starts a new snapshot generation, so requests routed from the same generation share one attempt and a failed attempt fails them all instead of each calling the provider in turn. The retry budget is per request, with no state that outlives it.

The locations come from asking for the bucket's whole list rather than which prefixes one session covers. Asking per session leaves Comet to discover the other scopes from 403s, which needs mutable per-store state and cannot route a partition that spans scopes. With the whole list up front, routing is a function of the path.

This is consistent with [Why no Comet-side cache](#why-no-comet-side-cache): the store caches no credentials, and every request still calls `getCredentialsForPath`. What it keeps is the provider's location list and a store for each location that has been read, so it grows with the vendor's policy list, not with the paths read. The list is only fetched again after a 403, so a change that causes none is not seen until the executor builds a new store.
This is consistent with [Why no Comet-side cache](#why-no-comet-side-cache): the store caches no credentials, and every request still calls `getCredentialsForPath`. What it keeps is the provider's location list and a store for each location that has been read, so it grows with the vendor's policy list, not with the paths read. The list is only fetched again after a 403 or a credential failure, so a change that causes neither is not seen until the executor builds a new store.

The dispatcher returns `null` for a provider that does not implement the interface without calling it, and Comet builds the same plain store as before. The Iceberg path does not use locations. Operations other than reads route by path without the 403 retry, since Comet only reads through these stores.
The dispatcher returns `null` for a provider that does not implement the interface without calling it, and Comet builds the same plain store as before. The Iceberg path does not use locations. Operations other than reads route by path without the retry, since Comet only reads through these stores.

## Path-specific behavior

Expand Down
2 changes: 1 addition & 1 deletion docs/source/user-guide/latest/s3-credential-providers.md
Original file line number Diff line number Diff line change
Expand Up @@ -224,7 +224,7 @@ Comet serves each request with the credential of the longest location that cover

Comet requests a location's credential by calling `getCredentialsForPath` with the location as the path, as you returned it but with a leading slash. Every request under a location shares that credential, so it must authorize every path the location is the longest match for, and your cache can key on the location. Locations apply to Comet's native Parquet reads only; Iceberg reads call `getCredentialsForPath` as they do for any provider.

**When Comet asks.** Comet calls `getPolicyLocations` when it creates the store for a bucket on an executor and keeps the answer for later reads of that bucket with the same S3 configuration. Reads that start at the same moment may each create a store and call it. If a read then fails with 403, Comet asks again, once for all the reads that failed on the same answer, and retries each read once if its path now falls under a different location, so a location added while a job runs is picked up. A location added or removed without causing a 403 is not seen until the executor creates a new store. Make `getPolicyLocations` thread-safe and independent of where it runs; it may be called on the driver or on executors.
**When Comet asks.** Comet calls `getPolicyLocations` when it creates the store for a bucket on an executor and keeps the answer for later reads of that bucket with the same S3 configuration. Reads that start at the same moment may each create a store and call it. If a read then fails with 403, or because `getCredentialsForPath` threw for the location Comet sent it to, Comet asks again, once for all the reads that failed on the same answer, and retries each read once if its path now falls under a different location. So a location added while a job runs is picked up even when you vend no credential for the bucket root, and a location you drop stops being used once its credential fails. A location added or removed without a read failing on it is not seen until the executor creates a new store. Make `getPolicyLocations` thread-safe and independent of where it runs; it may be called on the driver or on executors.

**Failures.** If `getPolicyLocations` throws or returns `null`, or returns a location that is `null` or invalid, the read fails. A location is invalid if, once decoded, it is not valid UTF-8 or has a segment that is empty, `.`, `..`, or contains a control character, so a URI such as `s3://bucket/a` is invalid too. Comet does not fall back to a broader credential.

Expand Down
17 changes: 16 additions & 1 deletion native/core/src/cloud/s3/credential_bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -339,14 +339,29 @@ struct RawCredentials {
expiration_epoch_millis: i64,
}

/// The bridge could not get a credential from the provider, as opposed to S3 rejecting one. It is
/// the source of the error `get_credential` returns, which `object_store` passes through to the
/// read unchanged. A `LocationScopedObjectStore` treats it like a 403, because a provider that has
/// no policy for a location throws, and that can mean the location changed since its snapshot.
#[derive(Debug)]
pub(crate) struct CredentialProviderError(pub(crate) String);

impl fmt::Display for CredentialProviderError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.0)
}
}

impl std::error::Error for CredentialProviderError {}

#[async_trait]
impl CredentialProvider for CometS3CredentialBridge {
type Credential = AwsCredential;

async fn get_credential(&self) -> object_store::Result<Arc<AwsCredential>> {
let raw = self.fetch_raw().map_err(|e| object_store::Error::Generic {
store: "S3",
source: e.to_string().into(),
source: Box::new(CredentialProviderError(e.to_string())),
})?;
Ok(Arc::new(AwsCredential {
key_id: raw.access_key_id,
Expand Down
Loading
Loading