diff --git a/docs/source/contributor-guide/s3-credential-provider-design.md b/docs/source/contributor-guide/s3-credential-provider-design.md index 8fa449dd555..e67168055fc 100644 --- a/docs/source/contributor-guide/s3-credential-provider-design.md +++ b/docs/source/contributor-guide/s3-credential-provider-design.md @@ -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 diff --git a/docs/source/user-guide/latest/s3-credential-providers.md b/docs/source/user-guide/latest/s3-credential-providers.md index df3d76005bb..04e6a2c9b85 100644 --- a/docs/source/user-guide/latest/s3-credential-providers.md +++ b/docs/source/user-guide/latest/s3-credential-providers.md @@ -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. diff --git a/native/core/src/cloud/s3/credential_bridge.rs b/native/core/src/cloud/s3/credential_bridge.rs index 415dcf2ae47..0993f551409 100644 --- a/native/core/src/cloud/s3/credential_bridge.rs +++ b/native/core/src/cloud/s3/credential_bridge.rs @@ -339,6 +339,21 @@ 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; @@ -346,7 +361,7 @@ impl CredentialProvider for CometS3CredentialBridge { async fn get_credential(&self) -> object_store::Result> { 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, diff --git a/native/core/src/parquet/objectstore/location_scoped.rs b/native/core/src/parquet/objectstore/location_scoped.rs index 7f368615e42..3144d5e6bc1 100644 --- a/native/core/src/parquet/objectstore/location_scoped.rs +++ b/native/core/src/parquet/objectstore/location_scoped.rs @@ -27,9 +27,11 @@ //! kept for the life of this store. //! //! The locations are a snapshot, so a 403 can mean a location was added or removed after it was -//! taken. A read (`get_opts` or `get_ranges`) that gets a 403 fetches the locations again, unless -//! another read already tried since this one was routed, and retries once if its path now routes to -//! a different location; otherwise the 403 is returned. A failed fetch fails every read that shared +//! taken. So can a failure to get a location's credential from the provider, because a provider +//! with no policy for a path throws rather than returning a credential S3 would reject. A read +//! (`get_opts` or `get_ranges`) that gets either error fetches the locations again, unless another +//! read already tried since this one was routed, and retries once if its path now routes to a +//! different location; otherwise the error is returned. A failed fetch fails every read that shared //! it. Other operations route by path without retrying, because Comet only reads through this store. use std::collections::HashMap; @@ -48,6 +50,8 @@ use object_store::{ }; use tokio::sync::Mutex; +use crate::cloud::s3::credential_bridge::CredentialProviderError; + const STORE: &str = "LocationScopedS3"; /// The credential path used for paths that no returned location covers. @@ -127,8 +131,21 @@ fn credential_path(location: &str) -> String { } } -fn is_forbidden(err: &Error) -> bool { - matches!(err, Error::PermissionDenied { .. }) +/// Whether `err` can mean the locations changed since the snapshot: S3 denied the request, or the +/// provider could not produce the credential of the location the request was routed to. +fn may_mean_stale_locations(err: &Error) -> bool { + matches!(err, Error::PermissionDenied { .. }) || is_credential_failure(err) +} + +fn is_credential_failure(err: &Error) -> bool { + let mut source = std::error::Error::source(err); + while let Some(e) = source { + if e.is::() { + return true; + } + source = e.source(); + } + false } /// The store chosen for one request, and the snapshot it was chosen from. @@ -143,7 +160,7 @@ struct Inner { source: LocationSource, factory: LocationStoreFactory, index: RwLock>, - /// Serializes refreshes, so the 403s from one snapshot share one attempt. + /// Serializes refreshes, so the failed reads from one snapshot share one attempt. refresh_lock: Mutex<()>, /// Location stores by credential path. stores: RwLock>>, @@ -183,8 +200,8 @@ impl Inner { )) } - /// Returns the route to retry on after `failed` returned the 403 `err` for `path`, or the error - /// the request should return. + /// Returns the route to retry on after `failed` returned `err` for `path`, or the error the + /// request should return. `err` is a 403 or a failure to get the location's credential. async fn retry_route(&self, path: &Path, failed: &Route, err: Error) -> Result { if let Err(refresh) = self.refresh(failed).await { return Err(Error::Generic { @@ -236,8 +253,8 @@ pub struct LocationScopedObjectStore { impl LocationScopedObjectStore { /// `locations` is the provider's first answer for `bucket`; `source` fetches it again after a - /// 403, and `factory` builds a location's store on first use. Fails if a location is not a - /// valid path. + /// 403 or a credential failure, and `factory` builds a location's store on first use. Fails if + /// a location is not a valid path. pub(crate) fn new( bucket: String, locations: Vec, @@ -297,7 +314,7 @@ impl ObjectStore for LocationScopedObjectStore { async fn get_opts(&self, location: &Path, options: GetOptions) -> Result { let route = self.inner.route(location)?; match route.store.get_opts(location, options.clone()).await { - Err(e) if is_forbidden(&e) => { + Err(e) if may_mean_stale_locations(&e) => { let retry = self.inner.retry_route(location, &route, e).await?; retry.store.get_opts(location, options).await } @@ -308,7 +325,7 @@ impl ObjectStore for LocationScopedObjectStore { async fn get_ranges(&self, location: &Path, ranges: &[Range]) -> Result> { let route = self.inner.route(location)?; match route.store.get_ranges(location, ranges).await { - Err(e) if is_forbidden(&e) => { + Err(e) if may_mean_stale_locations(&e) => { let retry = self.inner.retry_route(location, &route, e).await?; retry.store.get_ranges(location, ranges).await } @@ -374,13 +391,24 @@ mod tests { use std::time::Duration; use tokio::sync::Barrier; + /// How every read through one credential fails, set by credential path. + #[derive(Clone, Copy, Debug)] + enum Failure { + /// The provider throws, as it does for a path it has no policy for. + NoCredential, + /// Something unrelated to the credential fails, such as the connection. + Unreachable, + } + /// What S3 allows for one credential: reads under `allowed` succeed and every other read gets a /// 403. A success is reported as `NotFound` carrying the credential path, which shows which - /// location served the request without producing data. + /// location served the request without producing data. A credential path in `failures` fails + /// every read instead, checked on each read because the bridge asks the provider on each one. #[derive(Debug)] struct CredentialView { credential_path: String, allowed: Vec, + failures: Arc>>, gets: AtomicUsize, /// When set, each read waits here first, so a test can hold requests in flight together. gate: Option>, @@ -416,6 +444,31 @@ mod tests { if let Some(gate) = &self.gate { gate.wait().await; } + let failure = self + .failures + .lock() + .unwrap() + .get(&self.credential_path) + .copied(); + match failure { + // What `CometS3CredentialBridge::get_credential` returns when the provider throws. + Some(Failure::NoCredential) => { + return Err(Error::Generic { + store: "S3", + source: Box::new(CredentialProviderError(format!( + "no policy for {}", + self.credential_path + ))), + }) + } + Some(Failure::Unreachable) => { + return Err(Error::Generic { + store: "S3", + source: "connection refused".into(), + }) + } + None => {} + } let source = self.credential_path.clone().into(); let path = location.to_string(); if self.allowed.iter().any(|p| location.prefix_matches(p)) { @@ -458,6 +511,8 @@ mod tests { gate: Option<(String, Arc)>, /// Credential path whose store cannot be built. fail_build: Option, + /// How reads through a credential fail, shared with every store built. + failures: Arc>>, views: StdMutex>>, builds: AtomicUsize, } @@ -478,6 +533,7 @@ mod tests { refreshes: AtomicUsize::new(0), gate: None, fail_build: None, + failures: Arc::new(StdMutex::new(HashMap::new())), views: StdMutex::new(HashMap::new()), builds: AtomicUsize::new(0), } @@ -502,6 +558,14 @@ mod tests { *self.locations.lock().unwrap() = locations.iter().map(|l| l.to_string()).collect(); } + /// Makes every later read through `credential_path` fail with `failure`. + fn fail(&self, credential_path: &str, failure: Failure) { + self.failures + .lock() + .unwrap() + .insert(credential_path.to_string(), failure); + } + fn store(self: &Arc) -> LocationScopedObjectStore { let provider = Arc::clone(self); let source: LocationSource = Arc::new(move || { @@ -538,6 +602,7 @@ mod tests { .get(credential_path) .cloned() .unwrap_or_default(), + failures: Arc::clone(&provider.failures), gets: AtomicUsize::new(0), gate, }); @@ -876,4 +941,95 @@ mod tests { assert!(message.contains("bridge init failed"), "{message}"); assert!(!message.contains("policy locations"), "{message}"); } + + /// A provider with no bucket-wide policy throws when asked for the bucket root's credential. A + /// location added after the snapshot routes to the root until the locations are fetched again, + /// so that failure has to refresh them the way a 403 does. + async fn picks_up_a_location_when_the_root_credential_fails(use_ranges: bool) { + let provider = Arc::new(Provider::new( + &["warehouse/sales"], + &[ + ("/warehouse/sales", &["warehouse/sales"]), + ("/warehouse/finance", &["warehouse/finance"]), + ], + )); + provider.fail("/", Failure::NoCredential); + let store = provider.store(); + provider.set_locations(&["warehouse/sales", "warehouse/finance"]); + + for path in ["warehouse/finance/1", "warehouse/finance/2"] { + let served = if use_ranges { + served_by(get_ranges(&store, path).await) + } else { + served_by(get(&store, path).await) + }; + assert_eq!(served, "/warehouse/finance"); + } + assert_eq!(provider.refreshes.load(Ordering::SeqCst), 1); + assert_eq!(provider.gets("/"), 1); + } + + #[tokio::test] + async fn get_opts_picks_up_a_location_when_the_root_credential_fails() { + picks_up_a_location_when_the_root_credential_fails(false).await; + } + + #[tokio::test] + async fn get_ranges_picks_up_a_location_when_the_root_credential_fails() { + picks_up_a_location_when_the_root_credential_fails(true).await; + } + + /// A provider that folds a location into its parent stops vending the old location and throws + /// when asked for it. That failure refreshes the locations, and the parent serves the read. + #[tokio::test] + async fn moves_off_a_dropped_location_whose_credential_fails() { + let provider = Arc::new(Provider::new( + &["warehouse", "warehouse/finance"], + &[ + ("/warehouse", &["warehouse"]), + ("/warehouse/finance", &["warehouse/finance"]), + ], + )); + let store = provider.store(); + assert_eq!( + served_by(get(&store, "warehouse/finance/1").await), + "/warehouse/finance" + ); + + provider.set_locations(&["warehouse"]); + provider.fail("/warehouse/finance", Failure::NoCredential); + assert_eq!( + served_by(get(&store, "warehouse/finance/2").await), + "/warehouse" + ); + assert_eq!(provider.refreshes.load(Ordering::SeqCst), 1); + } + + /// A credential failure on a location the refreshed list still routes to is returned as it is, + /// the same as a 403 would be. + #[tokio::test] + async fn returns_a_credential_failure_when_the_locations_are_unchanged() { + let provider = Arc::new(Provider::new(&["a"], &[("/a", &["a"])])); + provider.fail("/a", Failure::NoCredential); + let store = provider.store(); + + let err = get(&store, "a/1").await.unwrap_err(); + assert!(is_credential_failure(&err), "got {err}"); + assert!(err.to_string().contains("no policy for /a"), "got {err}"); + assert_eq!(provider.refreshes.load(Ordering::SeqCst), 1); + assert_eq!(provider.gets("/a"), 1, "not retried on the same location"); + } + + /// Only a 403 or a credential failure can mean the locations changed, so any other error is + /// returned without fetching them again. + #[tokio::test] + async fn does_not_refresh_on_other_errors() { + let provider = Arc::new(Provider::new(&["a"], &[("/a", &["a"])])); + provider.fail("/a", Failure::Unreachable); + let store = provider.store(); + + let err = get(&store, "a/1").await.unwrap_err(); + assert!(err.to_string().contains("connection refused"), "got {err}"); + assert_eq!(provider.refreshes.load(Ordering::SeqCst), 0); + } } diff --git a/native/core/src/parquet/objectstore/s3.rs b/native/core/src/parquet/objectstore/s3.rs index 620bbdbdda5..3c6d3cc59f6 100644 --- a/native/core/src/parquet/objectstore/s3.rs +++ b/native/core/src/parquet/objectstore/s3.rs @@ -202,9 +202,10 @@ impl S3StoreTemplate { } /// Builds the store for a `CometS3LocationScopedCredentialProvider`. `bridge` was created on this -/// thread, which registered the provider. It fetches the locations again after a 403, and each -/// location's bridge is derived from it on first use, often on a Tokio worker, so every location -/// shares the bucket's provider registration without another `ensureInitialized` call. +/// thread, which registered the provider. It fetches the locations again after a 403 or a failure +/// to get a location's credential, and each location's bridge is derived from it on first use, +/// often on a Tokio worker, so every location shares the bucket's provider registration without +/// another `ensureInitialized` call. fn location_scoped_store( template: S3StoreTemplate, bucket: &str,