Skip to content

feat: per-location credentials for the S3 credential SPI - #6031

Merged
andygrove merged 7 commits into
apache:mainfrom
snmvaughan:feature/comet-s3-scoped-credential-provider
Sep 25, 2026
Merged

andygrove merged 7 commits into
apache:mainfrom
snmvaughan:feature/comet-s3-scoped-credential-provider

Conversation

@snmvaughan

@snmvaughan snmvaughan commented Sep 18, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

Closes #6207, which proposes the new @Public API for agreement, as the versioning policy asks.

Rationale for this change

On the Parquet path, object_store::CredentialProvider::get_credential() receives no request path, so one S3 store presents one credential, and Comet caches one store per (scheme://bucket, config_hash, backend). A CometS3CredentialProvider therefore gets one credential per bucket, requested with the path of the first file Comet reads. Vendors whose policies differ by location within a bucket (one STS session for warehouse/sales, another for warehouse/finance) get 403s on every location but the first.

What changes are included in this PR?

  • CometS3LocationScopedCredentialProvider (new @Public opt-in extension of CometS3CredentialProvider) with one method, List<String> getPolicyLocations(String bucket), returning every location in the bucket that has its own policy. Each request uses the credential of the longest location covering its path, matched one segment at a time after percent-decoding; the bucket root is an implicit location. Comet requests a location's credential with the location itself as the path. Locations apply to native Parquet reads; the Iceberg path is unchanged.
  • CometS3CredentialDispatcher.getPolicyLocations returns null for providers that do not implement the interface, without calling them, and otherwise a String[] copy, so provider list code runs inside the checked JNI call. A null list or location throws.
  • LocationScopedObjectStore (native): registered and cached once per bucket like any S3 store. It routes each request to its location's store and builds each location's store on first use from a template whose region was resolved up front, so nothing calls block_on inside an async read. On a 403 it fetches the locations again, once for all the reads routed from the same snapshot whether the fetch succeeds or fails, and retries each read once if its path now routes elsewhere.
  • Base providers are unaffected: create_store builds the same store as before with the same calls in the same order; the only addition is one dispatcher call on a cache miss that does not reach the provider.
  • Removes RetryOn403ObjectStore, the scope-hint JNI entry point, and the multi-entry object store cache from the earlier revision.
  • Docs: user guide section, design notes, and the versioning policy's public API list.

How are these changes tested?

  • Rust unit tests (parquet::objectstore::location_scoped):
    • Routing: longest covering location, segment boundaries (sales vs sales_eu), percent-decoding (Spark's %3A escapes), duplicate and root locations, and invalid locations (empty and .. segments, URIs, control characters, invalid UTF-8).
    • Reads through one store: three locations in rotating order (A/B/A), with lazy store construction.
    • 403 handling: unchanged locations; a location added after the snapshot; a retry that is also denied; concurrent 403s on get_opts and get_ranges sharing one refresh; reads on separate threads waiting for a refresh already in progress; concurrent reads sharing one failed refresh; a later read asking again; a failed store build reported as itself.
    • The multi-threaded test fails when the refresh mutex is removed, and the failed-refresh test fails when failures are not shared.
  • parquet::objectstore::s3: building a store from a template inside the Tokio runtime with no endpoint or region configured.
  • JUnit (CometS3LocationScopedCredentialProviderTest): base providers return null without being called; empty, null, null-element, non-String, lazy-list and provider exceptions; unknown handles.
  • CometPublicApiSuite pins the new @Public type.
  • CometS3CredentialBridgeSuite (MinIO, manual suite): a dedicated bucket with a per-bucket provider override reads files from four locations, including the bucket root, in one partition and asserts the credential path each read requested; a second test asserts that scans after a warm-up scan reuse the bucket's locations without asking again. The existing base-provider test now also asserts the base provider still receives a file path. The suite needs Docker and has not been run on this revision yet.

cargo fmt, cargo clippy --all-targets --workspace -- -D warnings, cargo test -p datafusion-comet --lib (438 passed), the JUnit tests and CometPublicApiSuite above, spotless, scalastyle, and Prettier pass locally.

@github-actions github-actions Bot added enhancement New feature or request area:scan Parquet scan / data reading area:ffi Arrow FFI / JNI boundary labels Sep 18, 2026
…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.
@snmvaughan
snmvaughan force-pushed the feature/comet-s3-scoped-credential-provider branch from 410c433 to 0d79257 Compare September 18, 2026 23:23

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Correctness

Reviewed 8a86d804 against base 5705a58ac.

The existing bucket-level cache can bind all reads to the first file's credential scope. This change adds an opt-in Java scope hint, stores several scope entries per existing configuration/backend key, chooses the longest covering prefix, and rebuilds a store after a 403. Keeping configuration and backend identity in the outer key preserves the existing separation of explicitly configured credentials and endpoints.

Four P2 findings are reproduced at the assigned head: recovery redirects the original scope to the replacement store, documented prefixes do not match native URL paths, an overlapping 403 can skip its retry, and region discovery can re-enter Tokio during async recovery. The first issue also affects a single partition containing files from several scopes because the planner prepares one store using its first file.

I compared the maintained Spark 3.5 and 4.0 Parquet reader paths at 5947fd6e and 03f28fc4. They initialize readers for the individual file split and its Hadoop context. Comet must continue to read each authorized file and surface genuine authorization failures. This PR changes credential routing and errors, not SQL type conversion, null handling, ANSI behavior, or arithmetic boundaries. Maintained 3.4 and 4.1 source branches were unavailable, so those versions remain unqualified.

Validation

The offline Rust fixture copied the production cache-matching and retry methods verbatim, used synthetic stores, and linked cached once_cell 1.21.4 and Tokio 1.53.1. It reproduced all four conditions, with controls showing that the original A store and the fresh concurrent-request store still authorize the failed paths. The nested-runtime probe uses a ready synthetic region future. No AWS, credentials, JNI execution, full Comet build, Spark suite, or Minio suite was involved. The author's test claims are separate from this evidence.

At the final cutoff, only the label check had succeeded. Comet CI and CodeQL were awaiting approval with zero product jobs executed. The new same-scope integration test also expects no additional policy callback, while production explicitly calls the callback before every cache lookup. Its URI-shaped hints and that assertion should be reconciled with the intended contract when adding end-to-end regression coverage.

Performance

The prefix mismatch prevents cache reuse for conforming nonempty hints. The component fixture measured zero hits in 1,000 same-scope lookups, versus 1,000 hits for the leading-slash format used by the Rust tests. This is a functional hit-count comparison, not an end-to-end latency measurement.

The exact lookup helper's median time for 20,000 lookups per sample, across three samples, was about 2 ns with one entry, 126 ns with 64 entries, and 2,378 ns with 1,024 entries for a hit. Misses had similar costs. These local synthetic timings show the linear scan cost only. They exclude JNI, credential vending, network I/O, and store construction. Scope hints and bridge/global-reference construction also run before cache hits. A focused native/JVM same-scope hit and concurrent-miss benchmark would establish whether that repeated work is justified.

Design

The opt-in sub-interface is a useful compatibility boundary, and longest-prefix selection can let disjoint scopes coexist. The store's mutable backing identity currently undermines that model: an entry retains its original scope while its wrapper starts using credentials for another scope. A path-aware selector over stable scoped stores would keep the routing rule in one place and allow one retry per operation without a process-lifetime retry latch.

Empty hints and policy-fetch errors become catchall entries. They need the same multi-scope recovery guarantees as nonempty hints. The implementation also wraps base-interface providers, despite the user guide saying their retry path stays inactive. The documentation should describe the final behavior consistently. The retained vector is process-wide and grows with distinct scope lists, so the old bucket-count-only size argument no longer describes its bounds. I found no newly demonstrated cross-user or cross-catalog credential disclosure and do not make that claim.

Abstraction & complexity

The cache helper, SPI bridge, and retry decorator each have a clear purpose, but scope selection and replacement ownership are split across the cache and wrapper. Consolidating path-based selection would remove the contradictory lifetime rules highlighted above. The appended rebuilt entry is also a raw store, so recovery behavior differs depending on how an entry was created. Please exercise the final routing design through actual scoped read sequences and concurrent failures, in addition to the current fake-store unit tests.

Comment on lines +104 to +107
self.rebuilt
.get()
.cloned()
.unwrap_or_else(|| Arc::clone(&self.inner))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Correctness

[P2] Could the rebuilt store be selected by request path without permanently replacing the backing store for every operation on this wrapper? The planner constructs one store from the first file and uses it for all files in a Parquet partition. With an A-only initial store and a B-only rebuilt store, reading A1, B1, then A2 returns success, success, then 403. The original store still authorizes A2, but current() now routes it to B and already_rebuilt() suppresses recovery. Appending a new scope entry does not repair the original cache entry or scans already holding this wrapper. I reproduced that sequence using these exact methods and synthetic scoped stores. Please retain path-specific routing and bound retries per operation, with A/B/A and three-disjoint-scope coverage.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Agreed, and I think this constrains where the fix can live. A path-aware selector inside the cache wouldn't be consulted at read time: prepare_object_store_with_configs registers one store per ObjectStoreUrl, and that URL carries scheme, config hash, backend and host but no scope. DefaultObjectStoreRegistry::register_store is a plain insert, so two NativeScan operators on the same bucket under different scopes overwrite each other, and FileScanConfig::open resolves from that single URL at execute time. That's also why the one-file-per-partition case you mention can't be fixed in the cache -- a file group holds many paths but only one URL.

Would it work to make the registered store itself the selector? One wrapper per (bucket, config_hash, backend) holding the scoped stores and dispatching on location inside get_opts and get_ranges would put routing in one place, give one retry per operation with no process-lifetime latch, and handle a partition that spans scopes. The alternative is splitting file groups by scope and registering a URL per scope, which looks considerably more invasive.

Either way, the comment above register_object_store currently documents the invariant that Vec<ScopeEntry> breaks, so that needs updating too.

Comment on lines +789 to +792
prefixes
.iter()
.filter(|p| path.starts_with(p.as_str()))
.map(|p| p.len())

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Correctness

[P2] Could both sides be canonicalized to the bucket-relative representation required by CometS3ScopedCredentialProvider before this comparison? requested_path is taken directly from url.path(), so it starts with /, while the Java contract explicitly requires prefixes without a leading slash. A conforming hint such as warehouse/table therefore never covers /warehouse/table/file. With only that scoped entry present, the exact lookup helper produced zero hits in 1,000 requests, causing a new store and HTTP pool to be constructed on each preparation. The absolute S3 URI used by the new Minio test also never matches. The comparison should enforce the documented path-component boundary too, so table cannot match table_extra. Please exercise the Java contract's prefix format through the native cache.

Comment on lines +177 to +181
let store = self.current();
match store.get_opts(location, options.clone()).await {
Err(e) if is_forbidden(&e) && !self.already_rebuilt() => {
debug!("RetryOn403: 403 on get({location}); rebuilding store");
let rebuilt = self.rebuild_once(Some(location))?;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Correctness

[P2] Could this decision distinguish a 403 from the original store from a 403 returned by the rebuilt store? Two requests can capture the original store before either completes. If the first finishes, rebuilds successfully, and retries, the second then observes already_rebuilt() == true and returns its original 403 without trying the fresh store. A deterministic interleaving of these exact methods reproduced one successful request and one failed request, even though a later call for the failed path succeeds on the rebuilt store. OnceCell prevents duplicate construction but does not make this guard safe. Please let each operation retry once against the published replacement and cover the same race in get_ranges.

Comment on lines +991 to +996
let (rebuilt_raw, _path) = objectstore::s3::create_store_with_bridge(
&url_for_rebuild,
&configs_for_rebuild,
rebuilt_bridge,
Duration::from_secs(300),
)?;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Correctness

[P2] Could recovery reuse the resolved region or perform the rebuild without synchronously entering Tokio? This closure is invoked inside the async get_opts/get_ranges operation. For an ordinary AWS bucket with neither fs.s3a.endpoint nor fs.s3a.endpoint.region configured, create_store_with_bridge calls get_runtime().block_on(resolve_bucket_region(bucket)). Tokio rejects that nested runtime entry with a panic, including when the region cache would return immediately. I verified the latter with the locked Tokio 1.53.1 and a ready cached-region future, without making an AWS request. Thus the first recoverable 403 can abort the scan instead of retrying. Please cover recovery without an explicit endpoint or region.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for working on this. sunchao's four findings all look right to me, so I won't repeat them. A few additional things below.

One process note first: Preflight is currently failing on the markdown formatting check, and that gates everything downstream. PR Build, the Spark SQL matrix and the Iceberg suites are all skipped, so nothing has actually been built or tested on this change yet. npx prettier "**/*.md" --write fixes both of the changed docs.

let entries = cache.entry(cache_key_for_rebuild.clone()).or_default();
// Deduplicate: a concurrent rebuild may have already inserted an
// entry with the same prefixes.
if !entries.iter().any(|e| e.prefixes == fresh_prefixes) {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

This dedup check means the append that the whole recovery story depends on usually doesn't happen. In the overreporting case the vendor returns the same over-broad prefix on the rebuild as it did on the initial fetch, so fresh_prefixes equals the existing entry's prefixes, the check fires, and nothing gets pushed. The wrapper is left silently rebound to the new store with no entry recording it, which is the state sunchao describes on the retry wrapper. Should the dedup key on something that tells the two sessions apart rather than on the prefix list?

/// non-string elements as end-of-list / skip respectively. Only called from
/// [`CometS3CredentialBridge::fetch_policy_locations`] where the JVM dispatcher normalizes null
/// returns to `Collections.emptyList()`, but the extra null-check keeps the helper self-contained.
fn read_java_string_list(

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Every other JNI call in this file goes through jni_static_call!, which runs check_exception for you. This helper calls call_method_unchecked directly, so if a vendor hands back a lazily evaluated List that throws from size() or get(i), the exception stays pending and the next JNI call on that thread misbehaves. Could it use the checked path instead?

Separately, line 448 casts each element straight through JString::from_raw with no type check, while the doc comment above says non-string elements are skipped. Worth making the code match the comment, since a vendor returning a List<Object> would currently be undefined behavior rather than a skipped element.

"org.apache.comet.cloud.s3.CometS3CredentialProvider",
"org.apache.comet.cloud.s3.CometS3Credentials")
"org.apache.comet.cloud.s3.CometS3Credentials",
"org.apache.comet.cloud.s3.CometS3ScopedCredentialProvider")

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

docs/source/about/versioning_policy.md enumerates the public API as well, and it still lists only the four pre-existing types. The comment at the top of this suite asks for the doc and this list to change together, so could CometS3ScopedCredentialProvider be added there too? It would also help to say what adding a method to the sub-interface later would mean for compatibility, since the existing text covers that question for CometS3CredentialProvider.

}
```

Comet uses the hint to key the native `object_store` registry as `(bucket, config_hash, backend, scope_prefixes)` instead of the base three-tuple, so two prefixes with disjoint scopes each get their own store. **The hint is advisory.** S3 itself is authoritative: if the vendor overreports a scope and S3 returns 403 anyway, the native cache transparently rebuilds the store under a widened (catchall) scope, retries the operation once, and continues. A second 403 propagates to Spark as a real error.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Two things on this paragraph. The registry isn't keyed as (bucket, config_hash, backend, scope_prefixes) -- the design doc has a whole section explaining why that shape was rejected, and the code keeps the three-tuple with the prefixes hanging off the entry. And the rebuild doesn't widen to a catchall; it re-fires the SPI and uses whatever prefixes come back.

Together with the point sunchao raised about base-interface providers being wrapped despite this page saying their retry path stays inactive, that's three claims on this page that need reconciling with the final behavior. The backward compatibility promise rests on the third one, so I'd like that one to be accurate in particular.

// plumbing: the scope hint reaches Rust and the cache preserves scope granularity.
// ---------------------------------------------------------------------

test("scoped provider: two reads inside the same scope share one object_store entry") {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Beyond the URI-shaped hint and the callback assertion sunchao flagged, there's a third thing working against this test. The native object_store cache is a process-wide static and MinioCometS3CredentialProvider defaults to an empty hint, so the earlier tests in this suite have already seeded a catchall entry under the same bucket key. find_matching_scope matches that at effective length 0, so the scoped entry never gets exercised even once the prefix format is fixed.

Would using a bucket the other tests don't touch be enough here, or is there a way to reset the native cache between tests?

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Follow-up at 8a86d804 against 5705a58ac: the reviewed source pair is unchanged, and the four P2 findings in the previous review remain.

One qualification to the registry-overwrite discussion: the current Parquet path captures the selected store in its reader factory, and each reader clones that retained store. Locked DataFusion 55.1.0 uses this supplied factory instead of its default factory. Replacing the registry entry therefore does not by itself retarget an already-created Comet Parquet reader. The earlier A/B/A finding remains because the retained retry wrapper changes its own backing store.

I corroborated the existing JNI exception comment: an exception from a lazy list's size() or get() remains pending on an already-attached JVM thread while the caller falls back to an empty scope. The JNI local frame releases references but does not clear that exception. The prefix-list deduplication, API documentation, user documentation, and cache-isolation comments are also supported by the source. The scoped tests need an isolated cache key and assertions about the selected store or scoped access, since policy callback counts alone do not establish routing.

Preflight now fails Prettier on both changed S3 credential documentation files. All product build, Spark, and benchmark jobs were skipped. CodeQL succeeded only for Actions. The logged CI checkout has the same tree as the reviewed head.

This follow-up used source and locked-dependency checks. I did not rerun the earlier synthetic Rust probes or execute JNI, Spark, or Minio tests. Maintained Spark 3.4 and 4.1 source branches remain unavailable.

@andygrove andygrove added this to the 1.1.0 milestone Sep 21, 2026
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.
@snmvaughan

Copy link
Copy Markdown
Contributor Author

Thanks both. The findings were right, and most of them came from one design choice: the store used for a request was decided by a mutable, process-lifetime latch inside the retry wrapper, and scopes were discovered one session at a time from 403s. I've replaced that design rather than patching it, in cb1bdec.

New design. The opt-in interface is now CometS3LocationScopedCredentialProvider.getPolicyLocations(String bucket): the provider lists every location in the bucket that has its own policy, instead of reporting what one session covers. As andygrove suggested, the registered store is the selector: LocationScopedObjectStore is cached and registered once per (bucket, config_hash, backend), holds a store per location, and routes each request by location to the longest covering location. The bucket root is an implicit location.

How each finding is addressed:

  • A/B/A and partitions spanning scopes (sunchao): routing is per request with no mutable backing store. There's a test that reads three locations in rotating order through one store.
  • Prefix format and segment boundary (sunchao): locations and request paths are both canonicalized with Path::from_url_path and matched segment by segment, so table doesn't cover table_extra. Tests cover the contract's format, slash variants and %-escapes.
  • Concurrent 403s (sunchao): a 403 re-fetches the locations, then retries once if the route changed. The retry budget is per request, with no latch. Every refresh attempt starts a new snapshot generation, so requests routed from the same snapshot share one attempt, including a failed one. Tests cover two in-flight get_opts and get_ranges requests failing together, reads on separate threads waiting for a refresh in progress, and a shared failed refresh.
  • Nested runtime (sunchao): the region is resolved when the bucket's store is created. Location stores are built from that template and never block_on, which a test checks inside the runtime.
  • Dedup (andygrove): gone. Nothing is appended on recovery; recovery only re-reads the provider's list.
  • JNI (andygrove): the dispatcher returns a String[] copy. List code runs inside the checked jni_static_call!, non-String elements fail with ArrayStoreException in Java, and native frees each element's local ref.
  • Registry invariant (andygrove, sunchao): one store per key again, so the register_object_store comment holds as written.
  • Base providers (sunchao, andygrove): unaffected. The dispatcher returns null without calling them, and create_store builds the same store as before. The user guide now says exactly that.
  • Versioning policy (andygrove): the new type is listed, along with what adding a method or changing how paths match locations would mean. The policy also asks for the addition to be agreed in an issue first; I'll open one and link it from the PR.
  • Docs (andygrove): the user guide and design notes are rewritten against the new behavior, and Prettier passes.
  • Integration tests (andygrove, sunchao): a dedicated bucket with a per-bucket provider override gives the scoped tests their own cache key. They assert which credential path each read requested, with four locations read in one partition, rather than counting callbacks; the reuse test runs a warm-up scan first so it does not depend on test order. The suite needs Docker, and I haven't run it on this revision yet.

The PR description is updated. Could a maintainer approve the workflow runs for this revision? Comet CI and CodeQL are waiting for approval. Happy to walk through any of it.

@snmvaughan snmvaughan changed the title feat: scope-aware object_store cache for the S3 credential SPI feat: per-location credentials for the S3 credential SPI Sep 24, 2026
@comphead
comphead self-requested a review September 24, 2026 18:10
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.
@comphead
comphead self-requested a review September 24, 2026 19:25
factory_bucket.as_str(),
credential_path,
AccessMode::Read,
&HashMap::new(),

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Each location's bridge goes through CometS3CredentialBridge::new again, so it only lands on the bucket's registration because this &HashMap::new() matches what create_store passes. #6023 changes create_store to forward the fs.s3a.* map. In a trial merge of the two, s3.rs only conflicts in create_store, so this line survives unchanged. After that, ensureInitialized would miss the existing key and create a second provider with an empty map. It would do that on whatever thread first reads the location, which is often a Tokio worker with no context class loader, so a vendor jar that comes in through --jars can't be loaded when Comet is on extraClassPath.

Could the location bridges be derived from bridge instead, reusing its handle and bucket string and only creating a new path string? They would then share the registration by construction, whatever create_store passes, and we'd skip an ensureInitialized round trip per location.

This doesn't need to hold up this PR. We can file a follow-on issue for it, as long as the fix lands before or together with #6023.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks, deriving from the bucket's bridge closes this. The location bridges now share the handle by construction. I also checked that the bucket bridge is always Read in create_store, so the inherited mode matches what the factory passed before.

}
}

fn is_forbidden(err: &Error) -> bool {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

The refresh only fires on a 403 from S3. A provider with no policy for the path it's asked about will usually throw instead, and get_credential turns that into a Generic error, which skips the refresh. I tried two cases with a fake location store that returns that error, and both fetched the locations zero times across three reads. In the first, the provider has no bucket-wide policy and throws for /, then adds warehouse/finance mid-job. Reads under it route to / and keep failing, so the new location is never picked up. In the second, the provider folds warehouse/finance into warehouse and throws for the old location. Later reads under it keep failing, while a freshly built store serves the same path from /warehouse. The store lives as long as the executor, so both stay broken until a restart.

The user guide says a location added while a job runs is picked up, and suggests keying the provider's cache on the location, which leads naturally to that throw. Should a credential failure from a location's store trigger the same refresh as a 403? Or should the contract require getCredentialsForPath to answer for / and for locations the provider has dropped, and the guide say so?

This one can also be a follow-on issue rather than blocking the PR.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks @snmvaughan

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.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: A bucket’s cached S3 store requested credentials using its first file’s path, preventing reads across independently authorized locations.
  • Design approach: The opt-in CometS3LocationScopedCredentialProvider supplies policy locations. One bucket-level wrapper routes each read to the longest matching location after percent-decoding and segment comparison.
  • Correctness / compatibility analysis: The rewrite addresses the earlier A/B/A routing, prefix matching, concurrent-403, nested-runtime and JNI concerns. I compared Spark’s SparkPath and V1/V2 Parquet readers across 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. The per-file routing and decoding match those paths. No additional introduced P1/P2 issues found within this review; the existing credential-error refresh concern remains reproducible.
  • Key design decisions: Location bridges reuse the bucket’s provider registration, and region resolution precedes asynchronous reads. Keeping routing inside the registered store removes the previous mutable backing-store identity and simplifies ownership.
  • Implementation sketch: Java validates and copies locations into String[], JNI retrieves them through the checked call, and native code combines a location index, lazy stores and serialized refresh generations. The reuse test served six reads using three store constructions and no policy refreshes. This validates reuse, not end-to-end performance.
  • Behavioral changes worth calling out: Scoped providers receive a location path or /. A 403 refreshes locations and permits one retry when the route changes. Base providers retain their existing behavior, and Iceberg does not use location routing.
  • Suggested improvements: Address the existing credential-error refresh concern. At location_scoped.rs:130, only PermissionDenied triggers refresh. A provider exception becomes Generic, so an added location with no root credential, or a removed location whose provider now throws, remains unreadable through the cached store. The bounded reproduction exercised both get_opts and get_ranges: three failures, zero location refreshes, and a successful read through a fresh store. Recognize credential-fetch failures explicitly and apply the same bounded refresh and changed-route retry. This is existing feedback, so no duplicate finding is added.

Reviewed the entire 16-file diff from base 5705a58ac2ef7e6b2674c87f3b495f789ebfcc8d to head 13e6344631fef807166f3b5b94efcfc0fab31f41. The PR remained non-draft. Read AGENTS.md and the discussion snapshot, excluding Copilot. Routed skill: review-comet-pr; no sibling skill applied.

Exact-head CI: the label check succeeded. Comet CI and CodeQL remain action_required, each with zero jobs. There is no product-build or test verdict.

Validation: all 16 tests from the unchanged routing module passed in an isolated crate using the locked direct dependency versions, including object_store 0.13.2 and Tokio 1.53.1. Two additional credential-error probes confirmed the existing concern without network requests. All nine dispatcher JUnit tests passed using Java 21 with source/target 17. Java 17 release checking was unavailable. Full Comet builds, integrated JNI execution, Spark SQL, MinIO and end-to-end benchmarks were not run. Project files remain unchanged.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: One cached S3 store used the first file’s credential scope for the whole bucket, preventing reads across independently authorized locations.
  • Design approach: The opt-in CometS3LocationScopedCredentialProvider supplies policy locations. A bucket-level wrapper routes each read to the longest matching location.
  • Correctness / compatibility analysis: The current implementation addresses the earlier routing, concurrent-403, nested-runtime and JNI findings. I compared URI handling and per-file Parquet reads against verified Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No additional introduced P1/P2 issues found within this review. The existing credential-error refresh defect remains reproducible.
  • Key design decisions: Location bridges share the bucket’s provider registration, and region resolution precedes asynchronous reads. Keeping routing inside one wrapper simplifies ownership and removes the previous mutable backing-store identity.
  • Implementation sketch: Java validates and copies locations into String[], checked JNI retrieves them, and native code maintains a location index, lazy stores and serialized refresh generations. The reuse test served six reads with three store constructions and no policy refreshes. This establishes reuse, not end-to-end performance.
  • Behavioral changes worth calling out: Scoped providers receive a location path or /. A 403 refreshes locations and allows one retry when the route changes. Base providers retain their existing behavior. Iceberg does not use location routing.
  • Suggested improvements: Address the existing credential-error refresh concern. At native/core/src/parquet/objectstore/location_scoped.rs:130, only PermissionDenied triggers refresh. Provider exceptions become Generic, so adding a location when no root credential exists, or removing a location whose provider subsequently throws, leaves authorized files unreadable through the cached store. Both get_opts and get_ranges reproduced three failures with zero refreshes, while a fresh store succeeded. Recognize credential-fetch failures explicitly and apply bounded refresh and changed-route retry. This remains an existing P2 blocker and is not duplicated as a new finding.

Reviewed the entire 16-file diff from base 5705a58ac2ef7e6b2674c87f3b495f789ebfcc8d to head 13e6344631fef807166f3b5b94efcfc0fab31f41. The PR remains non-draft. Read AGENTS.md and all supplied discussion and reviews, excluding Copilot. Routed skill: review-comet-pr. No sibling skill applied.

Exact-head CI: the label check passed. Comet CI and CodeQL remain action_required, each with zero jobs. There is no product-build or test verdict.

Validation: reran all 16 routing tests using a byte-identical production module in an isolated crate with matching locked direct dependency versions. Both additional credential-error probes confirmed the existing defect without network requests. Recompiled the dispatcher tests with Java 17 and --release 17; all nine passed. Full Comet builds, integrated JNI execution, Spark SQL, MinIO and end-to-end benchmarks were not run. Project files remain unchanged.

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: One cached S3 store used the first file’s credential scope for the bucket, preventing reads across independently authorized locations.
  • Design approach: The opt-in CometS3LocationScopedCredentialProvider supplies policy locations. A bucket-level wrapper routes each request to its longest covering location.
  • Correctness / compatibility analysis: The rewrite addresses the earlier routing, concurrent-403, nested-runtime and JNI findings. I compared URI handling and per-file Parquet reads against verified Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0. No additional introduced P1/P2 issues found within this review. The existing credential-error refresh defect remains reproducible.
  • Key design decisions: Location bridges share the bucket’s provider registration, and region resolution precedes asynchronous reads. Keeping routing inside one wrapper simplifies ownership and removes mutable backing-store identity.
  • Implementation sketch: Java validates and copies locations into String[], checked JNI retrieves them, and native code combines a location index, lazy stores and serialized refresh generations. The reuse test serves six reads with three store constructions and no policy refreshes. This establishes reuse, not end-to-end performance.
  • Behavioral changes worth calling out: Scoped providers receive a location path or /. A 403 refreshes locations and permits one retry when the route changes. Base providers retain their behavior. Iceberg does not use location routing.
  • Suggested improvements: Address the existing credential-error refresh concern. At native/core/src/parquet/objectstore/location_scoped.rs:130, only PermissionDenied triggers refresh. Provider exceptions become Generic, so adding a location without a root credential, or removing a location whose provider subsequently throws, leaves authorized files unreadable through the cached store. Both cases reproduced three failures with zero refreshes through each of get_opts and get_ranges, while a fresh store succeeded. Recognize credential-fetch failures explicitly and apply bounded refresh with changed-route retry. This remains an existing P2 concern and is not duplicated as a new finding.

Reviewed the entire 16-file diff from base 5705a58ac2ef7e6b2674c87f3b495f789ebfcc8d to head 13e6344631fef807166f3b5b94efcfc0fab31f41. The PR remains non-draft. Read AGENTS.md and supplied reviews, comments and threads, excluding Copilot. Routed skill: review-comet-pr. No sibling skill applied.

Exact-head CI: Comet CI and Required Checks now pass, including native tests, all four Linux Spark 4.1 Comet suites, and TPC-H/TPC-DS result verification. CodeQL also passes. Jobs tested merge commit d3211cdcc9e83a26941ae5c8e3dcad077bb12628 with newer main. Fourteen of the sixteen changed files match the pinned head, including the routing, bridge, dispatcher and tests. Spark SQL, Iceberg, macOS and benchmark jobs were skipped.

Validation: reran all 16 routing tests using a byte-identical production module in an isolated crate with matching locked direct dependency versions. Both credential-error probes confirmed the existing defect without network requests. Recompiled the dispatcher tests using Java 17 and --release 17; all nine passed. No local full Comet build, integrated JNI/MinIO run, Spark SQL suite or end-to-end benchmark was performed. Project files remain unchanged.

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks @snmvaughan. Deriving the location bridges from the bucket's bridge closes the interaction with #6023. I checked 13e634463 locally, merged with main. The location_scoped, objectstore and parquet_support tests pass, and clippy and fmt are clean.

The refresh concern on location_scoped.rs:130 still reproduces, and I filed #6221 for it. I don't think it needs to hold up this PR. It is not a regression, because only providers that implement the new interface take the new path, and nobody does yet. It fails loudly, with no wrong results and no fallback to a broader credential, and a one-off provider error doesn't stick. What persists is stale routing after a vendor's locations change while the executors are up, until they restart.

I would like the fix to land in the same release that ships the interface, though. Otherwise vendors building against it would have to learn the workaround, and the guide's "picked up" claim wouldn't hold for a provider without a bucket-wide credential. Would you be able to pick up #6221?

@andygrove
andygrove added this pull request to the merge queue Sep 25, 2026
@andygrove

Copy link
Copy Markdown
Member

Thanks @snmvaughan. Deriving the location bridges from the bucket's bridge closes the interaction with #6023. I checked 13e634463 locally, merged with main. The location_scoped, objectstore and parquet_support tests pass, and clippy and fmt are clean.

The refresh concern on location_scoped.rs:130 still reproduces, and I filed #6221 for it. I don't think it needs to hold up this PR. It is not a regression, because only providers that implement the new interface take the new path, and nobody does yet. It fails loudly, with no wrong results and no fallback to a broader credential, and a one-off provider error doesn't stick. What persists is stale routing after a vendor's locations change while the executors are up, until they restart.

I would like the fix to land in the same release that ships the interface, though. Otherwise vendors building against it would have to learn the workaround, and the guide's "picked up" claim wouldn't hold for a provider without a bucket-wide credential. Would you be able to pick up #6221?

I'm going to create a PR now to fix the issue

@sunchao sunchao left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Summary

  • Prior state and problem: One cached S3 store used the first file’s credential scope for the bucket, preventing reads across independently authorized locations.
  • Design approach: CometS3LocationScopedCredentialProvider supplies policy locations. A bucket-level wrapper routes each request to its longest covering location after percent-decoding and segment comparison.
  • Correctness / compatibility analysis: The rewrite addresses the earlier routing, concurrent-403, nested-runtime and JNI findings. I verified relevant Spark sources for 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 and compared URI handling and per-file Parquet reads. No additional introduced P1/P2 issues found within this review. The previously reported credential-error refresh defect remains reproducible.
  • Key design decisions: Location bridges share the bucket’s provider registration, and region resolution happens before asynchronous reads. Keeping routing inside one wrapper simplifies ownership and removes the earlier mutable backing-store identity.
  • Implementation sketch: Java validates and copies locations into String[], checked JNI retrieves them, and native code maintains the location index, lazy stores and serialized refresh generations. The reuse test serves six reads with three store constructions and no policy refreshes. This demonstrates reuse, without establishing end-to-end performance.
  • Behavioral changes worth calling out: Scoped providers receive a location path or /. A 403 refreshes locations and allows one retry when the route changes. Base providers retain their behavior, and Iceberg does not use location routing.
  • Suggested improvements: Complete the existing credential-error refresh fix tracked in #6221. At native/core/src/parquet/objectstore/location_scoped.rs:130, only PermissionDenied triggers refresh. Provider exceptions become Generic, so adding a location without a root credential, or removing a location whose provider subsequently throws, leaves authorized files unreadable through the cached store. Both cases reproduced three failures and zero refreshes through each of get_opts and get_ranges, while a fresh store succeeded. Recognize credential-fetch failures explicitly and apply bounded refresh with changed-route retry. Maintainers accepted follow-up handling, but this P2 concern remains unresolved at the reviewed head and is not duplicated as a new finding.

Reviewed the entire 16-file diff from base 5705a58ac2ef7e6b2674c87f3b495f789ebfcc8d to head 13e6344631fef807166f3b5b94efcfc0fab31f41, including all prerequisite commits. The PR remains non-draft. Read AGENTS.md and the supplied reviews, issue comments, inline comments and threads, excluding Copilot. Routed skill: review-comet-pr. No sibling skill applied.

Exact-head CI: 23 checks passed and 14 were skipped. Comet CI, Required Checks and CodeQL passed, including native tests, all four Linux Spark 4.1 Comet suites and TPC-H/TPC-DS verification. CI tested merge commit d3211cdcc9e83a26941ae5c8e3dcad077bb12628 against newer main. Fourteen of the sixteen changed files match the pinned head, including the routing, bridge, dispatcher and tests. Spark SQL, Iceberg, macOS and benchmark jobs were skipped.

Validation: reran all 16 routing tests using a byte-identical production module in an isolated crate with matching locked direct dependency versions. Both additional credential-error probes confirmed the existing defect without network requests. Compiled the dispatcher tests with Java 17 and --release 17; all nine passed. No local full Comet build, integrated JNI/MinIO run, Spark SQL suite or end-to-end benchmark was performed. Project files remain unchanged.

Merged via the queue into apache:main with commit e5d0b75 Sep 25, 2026
37 checks passed
@snmvaughan
snmvaughan deleted the feature/comet-s3-scoped-credential-provider branch September 25, 2026 16:26
@andygrove andygrove added the backport-1.0 Candidate for backporting to 1.0 release branch label Sep 25, 2026
andygrove added a commit that referenced this pull request Sep 25, 2026
* feat: add CometS3ScopedCredentialProvider opt-in scope-hint sub-interface

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.

* feat: add scope-hint JNI bridge + RetryOn403ObjectStore correctness net

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.

* feat: scope-aware object_store cache with rebuild-on-403 wiring

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.

* feat: route S3 credentials per policy location

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.

* fix: compile CometS3CredentialBridgeSuite on Scala 2.12

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.

* fix: derive location-scoped bridges from the bucket's bridge

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 (#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.

(cherry picked from commit e5d0b75)

Adapted for branch-1.0:
- parquet/objectstore/s3.rs: kept branch-1.0's region-lookup error
  message ("Failed to resolve region: {e}") in S3StoreTemplate::new.
  main's longer message, which suggests fs.s3a.endpoint settings, comes
  from #5314 (S3 compliant filesystems), which is not on branch-1.0.
  That message was the only conflict; the rest of the file matches main.

Co-authored-by: Steve Vaughan <email@stevevaughan.me>
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:ffi Arrow FFI / JNI boundary area:scan Parquet scan / data reading backport-1.0 Candidate for backporting to 1.0 release branch enhancement New feature or request

Projects

None yet

Development

Successfully merging this pull request may close these issues.

S3 credential SPI: per-location credentials within a bucket

4 participants