Conversation
…face Adds `CometS3ScopedCredentialProvider`, an opt-in `@Public` sub-interface of `CometS3CredentialProvider` that lets a vendor advertise the narrowest known-safe S3-key prefixes the STS session it is about to vend will authorize. The native layer (added in a follow-up commit) will use this hint to key the `object_store` cache one entry per distinct scope on a bucket instead of the single per-bucket store used today. `CometS3CredentialDispatcher` gains a `getPolicyLocationsFor(long, String, String, int)` static method: it looks up the registered provider by handle, checks `instanceof CometS3ScopedCredentialProvider`, dispatches, and normalizes a null return to an empty list. Base-interface-only providers keep today's behavior — the method returns `Collections.emptyList()` without touching the provider, which the native side reads as "no scope hint" and falls back to single-entry-per-bucket caching. The scope hint is advisory only. The correctness ground is a 403-retry wrapper inside Comet native (added later in this series): on 403 the wrapper invalidates the entry and re-fires the SPI with the actual failing path in context, so over-reporting is self-healing at the cost of one 403 per newly discovered scope boundary and under-reporting only costs extra SPI churn. `CometPublicApiSuite` gains the new type in its pinned `@Public` set; `MinioCometS3CredentialProvider` (test fixture) implements the sub-interface with an installable prefix list for the upcoming scope-aware IT scenarios; `CometS3ScopedCredentialProviderTest` covers dispatch to base vs. scoped provider, null normalization, exception propagation, and handle validation.
Second half of the scope-hint SPI landing (Layer 1 was #TBD). Introduces the
native-side plumbing that Layer 3 wires into the parquet object-store cache:
* jni-bridge: add `method_get_policy_locations_for` alongside the existing
`ensure_initialized` / `get_credentials_for_path` static-method IDs. Signature
matches the Java dispatcher entry:
`(JLjava/lang/String;Ljava/lang/String;I)Ljava/util/List;`
* CometS3CredentialBridge::fetch_policy_locations — thin JNI wrapper that reads
a `java.util.List<String>` back from the dispatcher. A null return, a base
(non-scoped) provider, or an empty list all normalize to an empty vec, which
the cache layer will interpret as "single-entry-per-bucket" semantics for
backward compatibility. Rewires the module-level doc-comment to describe the
new scope-hint contract (advisory prefixes, S3 remains authoritative via the
403-retry safety net layered above).
* RetryOn403ObjectStore (new `parquet::objectstore::retry`): correctness ground
for the advisory scope hint. On a single `Error::PermissionDenied` from the
wrapped store, invoke a caller-provided rebuild closure once, then retry the
same operation against the rebuilt store. A second 403 propagates without a
further retry. Only `PermissionDenied` (403) triggers the retry —
`Unauthenticated` (401) is treated as a permanent credential-config error.
Scope of interception is deliberately narrow: `put_opts`, `get_opts`,
`get_ranges`, `list_with_delimiter`, `copy_opts`, `rename_opts`. Stream
variants (`list`, `list_with_offset`, `delete_stream`) surface per-item errors
as-is; per-item retry would require materializing the stream and there is no
observed failure mode outside `get_opts`/`get_ranges` on the parquet path.
`put_multipart_opts` also skips retry — a partially-uploaded multipart cannot
be transparently restarted.
* Eight unit tests covering: pass-through success, rebuild-then-retry, second-
403 propagation, rebuild idempotency across ops, rebuild-error surface,
401-not-retried, Send+Sync bound, composes-with-Arc<Mutex<...>>.
No callers yet: Layer 3 (scope-aware cache in `parquet_support`) will construct
the rebuild closure and wire wrapped stores into the registry. Compiles clean
with the expected dead-code warnings on the new symbols.
Final layer of the scope-hint SPI landing. Wires
`CometS3ScopedCredentialProvider` (Layer 1) and the JNI bridge / retry
wrapper (Layer 2) into `parquet_support`'s object-store registry so that a
single bucket can host multiple scoped stores, and any misdiagnosed 403
transparently falls back to a widened-scope rebuild instead of aborting the
scan.
Native changes
--------------
* `parquet::objectstore::s3` — factor the builder into two entry points so
`parquet_support` can perform its own bridge lifecycle around the cache
lookup:
- `try_construct_bridge(url, configs) -> Result<Option<Arc<Bridge>>>`
returns `Ok(None)` when `comet.credential.provider.class` is absent
(no SPI in play), `Ok(Some(_))` when the dispatcher successfully
resolves, `Err(_)` on init failure.
- `create_store_with_bridge(url, configs, bridge, min_ttl)` takes the
pre-constructed bridge (or `None` to fall back to the AWS credential
chain) and returns the raw `AmazonS3` store. The old monolithic
`create_store` is deleted — the internal `test_create_store` now calls
`create_store_with_bridge(&url, &configs, None, Duration::from_secs(300))`
directly.
* `parquet::parquet_support` — replace the per-key
`HashMap<Key, Arc<dyn ObjectStore>>` with `HashMap<Key, Vec<ScopeEntry>>`
so one bucket can carry disjoint prefix-scoped stores. `ScopeEntry` pairs
a `Vec<String>` of prefix hints (empty ⇒ catchall) with the store
produced under those hints. Three private helpers make the intent
explicit:
- `path_covered(path, prefixes)` — `true` when `prefixes` is empty
(catchall) or any prefix is a byte-prefix of `path`.
- `find_matching_scope(entries, path)` — first entry whose prefixes
cover `path`; used on the read path.
- `invalidate_scope(cache, key, prefixes)` — drops the matching entry
before we push its catchall replacement; prunes the key when the
vector empties.
The rewritten `prepare_object_store_with_configs` calls
`try_construct_bridge` before the cache lookup so it can also drive the
scope-hint fetch (`bridge.fetch_policy_locations().unwrap_or_default()`).
On a cache miss the raw store is wrapped in `RetryOn403ObjectStore`
whenever a bridge is present; the rebuild closure re-runs
`try_construct_bridge` + `create_store_with_bridge`, then swaps the
cache entry for a catchall entry (empty prefixes), so subsequent reads on
the same bucket route through the widened store without another 403.
Design trade-off worth spelling out: on 403 we widen to catchall rather
than requesting a fresh scope for the failing path. The alternative
would require threading the failing `Path` into the rebuild closure,
which the `OnceCell` "rebuild at most once" semantics do not naturally
support. In exchange for the simpler correctness net we lose vendor-
scope granularity on that bucket for the remainder of the process.
Well-behaved providers that never overreport a scope are unaffected;
overreporting providers pay a widened-cache cost once per bucket.
* `ObjectStoreCacheKey` is `(String, u64, bool)` in the tree, not the
2-tuple the plan assumed — the third element is `hdfs_backend`. All
three existing seed-based tests (`check_isolated_stores`,
`native_s3_aliases_share_cache_and_registration_identity`,
`keeps_native_file_url_separate_from_explicit_hadoop_file_routing`) now
seed `Vec<ScopeEntry>` values under the unchanged 3-tuple keys.
* Three new unit tests exercise the private helpers:
- `path_covered_treats_empty_prefixes_as_catchall`
- `find_matching_scope_returns_first_covering_entry`
- `invalidate_scope_drops_matching_entries_and_prunes_empty_keys`
All 60 tests in `parquet::objectstore::*` and 16 tests in
`parquet::parquet_support::tests` pass locally.
Docs
----
* User guide: new "Scope hints via `CometS3ScopedCredentialProvider`"
section on the S3 credential providers page describing the opt-in
sub-interface, when to implement it, backward-compatibility, and the
advisory / 403-retry-authoritative contract.
* Contributor guide: new "Scope-aware `object_store` registry" section on
the SPI design page, spelling out the `Vec<ScopeEntry>` shape, why the
outer cache key stayed 3-tuple rather than absorbing prefixes, and the
widen-to-catchall design trade-off.
Scala IT
--------
`CometS3CredentialBridgeSuite` gains two Docker-tagged (Minio) scenarios
against the real bridge:
* `scoped provider: two reads inside the same scope share one object_store
entry` — installs a single scope prefix, reads two paths under it,
verifies the scope hint is fetched only once (cache reuses the
`ScopeEntry`) but credentials are refetched per read.
* `scoped provider: reads under disjoint scopes get separate object_store
entries` — installs scope A, reads under A, swaps to scope B, reads under
B; verifies a fresh scope-hint fetch was triggered (proving cache
fragmentation on disjoint prefixes rather than accidental sharing).
Minio does not enforce per-prefix denial out of the box, so the full
overreport-then-403-recover path is proven by
`native/core/src/parquet/objectstore/retry.rs`'s eight unit tests rather
than an end-to-end IT. Deferred to a follow-up if per-prefix Minio user
policies become available in the test base.
Replace the scope-hint design with CometS3LocationScopedCredentialProvider, an opt-in extension of CometS3CredentialProvider whose getPolicyLocations(bucket) lists every location in a bucket that has its own policy. For such a provider, create_store returns a LocationScopedObjectStore that is cached and registered once per bucket like any S3 store. It serves each request with the store of the longest location covering the request's path, matched one segment at a time after percent-decoding, with the bucket root as an implicit location. A location's store is built on first use from a template whose region was resolved when the bucket's store was created, so nothing blocks on the Tokio runtime inside an async read. A 403 fetches the locations again, once for all the reads routed from the same snapshot whether the fetch succeeds or fails, and retries once if the path now routes to a different location. The dispatcher returns null for a provider that does not implement the interface, without calling it, and create_store builds the same store for it as before. Otherwise the dispatcher copies the provider's list into a String[], so list code runs inside the checked JNI call. Removes RetryOn403ObjectStore, the scope-hint JNI entry point, and the multi-entry object store cache.
Scala 2.12 cannot choose between java.util.Set.of(E) and Set.of(E...) for a single argument, so the Spark 3.4 and 3.5 builds failed to compile the suite. Use Collections.singleton instead.
Each location's bridge called CometS3CredentialBridge::new again with empty catalog properties, so it only reached the bucket's provider registration while create_store also passed an empty map. Once create_store forwards the fs.s3a.* map (apache#6023), ensureInitialized would miss the existing key and create a second provider on whichever thread first read the location, often a Tokio worker with no context class loader. Location bridges now come from CometS3CredentialBridge::for_path, which reuses the bucket bridge's handle and bucket string and creates only the path string. They share the registration whatever create_store passes, and skip an ensureInitialized round trip per location. The MinIO suite asserts that one provider instance serves every location.
A LocationScopedObjectStore fetched the provider's locations again only after a 403. A provider with no policy for a path throws instead, and the bridge reported that as a plain Generic error. So a location added while the executors were up was never picked up when the provider vends no credential for the bucket root, and a location the provider dropped kept being used after it stopped vending it. Both stayed broken until the executors restarted. The bridge now gives its failures a CredentialProviderError source, which object_store passes through to the read unchanged. The store treats a read that fails with one like a 403: the same bounded refresh, and one retry if the path now routes to a different location. Other errors still return without a refresh. Closes apache#6221.
|
@snmvaughan could you help review? |
…resh-on-credential-error
sunchao
left a comment
There was a problem hiding this comment.
Summary
- Prior state and problem: Credential-provider exceptions bypassed location refresh, leaving long-lived executors unable to discover added or removed S3 policy locations.
- Design approach: Mark bridge failures with
CredentialProviderErrorand reuse the existing bounded refresh and retry path. - Correctness / compatibility analysis: Confirmed that
object_store0.13.2 preserves the typed error through S3 reads. Unchanged routes return the original error, unrelated failures do not refresh, and retries remain bounded. Compared Spark Parquet sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and experimental 4.2.0. The change affects Comet’s credential routing without changing Spark-facing APIs or Parquet semantics. - Key design decisions: Reusing snapshot generations and the existing refresh lock avoids additional persistent state. Successful reads gain no additional provider calls. Error-chain inspection and location refresh occur on failed reads.
- Implementation sketch: The bridge supplies the typed source, both
get_optsandget_rangesrecognize it, and focused tests and documentation describe the new trigger. - Behavioral changes worth calling out: Added locations can become readable when the root provider throws. Removed locations can reroute to their current parent. Reads retry only when the refreshed route changes.
- Suggested improvements: No introduced P1/P2 issues found within this review.
Reviewed the entire five-file diff from base e5d0b7575f1b7bfe4b19da7246f2d9733357e18f to head 848bdfce4a8a7199f9eb8373d8d71e1b45e76cc2. Confirmed non-draft status and read the supplied snapshot and live discussion. Routed skills: review-comet-pr and review-comet-ffi-pr.
Exact-head CI at 2026-09-25 17:10 UTC: 18 successful checks, six running and 13 skipped, with no failures. Native build, Rust tests, clippy and TPC-H passed. CI logs confirm all 21 location-scoped tests passed. Spark 4.1 Comet suites, TPC-DS and label-run preflight remained running. Spark SQL, Iceberg and macOS suites were skipped in the main PR run.
Validation: 25 isolated Rust tests passed using the exact routing module and extracted error type, including four probes through real AmazonS3 stores with mocked credential providers. Restoring the old 403-only trigger made four new regression tests fail. Local validation did not execute the JVM/MinIO integration or a full Comet build. Project files remain unchanged.
|
This looks good. |
…6223) (#6276) A LocationScopedObjectStore fetched the provider's locations again only after a 403. A provider with no policy for a path throws instead, and the bridge reported that as a plain Generic error, so a location added while the executors were up was never picked up, and a dropped location kept being used, until the executors restarted. The bridge now gives its failures a CredentialProviderError source, which object_store passes through to the read unchanged. The store treats a read that fails with one like a 403: the same bounded refresh, and one retry if the path now routes to a different location. Other errors still return without a refresh. (cherry picked from commit 68c4e9b)
Regenerate the 1.1.0 changelog with the generator from #6358, for the same range as #6282 (1.0.0..021c378), and format it with prettier 3.9.9. The generator used for #6282 credited each PR only to the author of its merge commit, so it left out Michael Taranov, whose commits from #6092 are in #6219, and credited the backports only to the person who opened them. The PR lines for #5310, #6219, #6223 and the four backports now name their co-authors, and the credits count each co-authored PR, which adds Michael Taranov and brings the contributor count to 41.
Which issue does this PR close?
Closes #6212.
Closes #6221.
Rationale for this change
#6031 added
CometS3LocationScopedCredentialProvider. Its location-scoped store fetches the provider's locations again only after a 403, but a provider with no policy for a path throws instead. So a location added while the executors are up is never picked up when the provider vends no credential for the bucket root, and a location the provider drops keeps being used after it stops vending it. Reads under those paths fail until the executors restart. #6212 and #6221 describe the same problem in more detail.What changes are included in this PR?
This is option (a) from #6212.
CometS3CredentialBridge::get_credentialgives its failures aCredentialProviderErrorsource. The message is unchanged.object_storepasses a credential provider's error through to the read unchanged, so the location-scoped store can recognize it.LocationScopedObjectStoretreats a read that fails with that source like a 403. It uses the same bounded refresh, where the reads routed from one snapshot share one attempt, and the same single retry when the path now routes to a different location. Other errors still return without a refresh.Nothing in the
@PublicAPI changes. Comet just callsgetPolicyLocationsin one more situation, and the guide already asks providers to make it safe to call at any time.How are these changes tested?
There are new tests in
location_scoped.rs. The existing fake-provider harness can now make a location's credential fail the way the bridge reports a provider exception.get_optsandget_ranges.The four positive tests fail if the trigger goes back to 403s only, and the last one fails if every error triggers a refresh. The
location_scoped,objectstoreandparquet_supportRust tests pass, and clippy with-D warningsand fmt are clean. I did not run the MinIOCometS3CredentialBridgeSuite, which needs Docker and is not part of PR CI.