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.
410c433 to
0d79257
Compare
sunchao
left a comment
There was a problem hiding this comment.
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.
| self.rebuilt | ||
| .get() | ||
| .cloned() | ||
| .unwrap_or_else(|| Arc::clone(&self.inner)) |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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.
| prefixes | ||
| .iter() | ||
| .filter(|p| path.starts_with(p.as_str())) | ||
| .map(|p| p.len()) |
There was a problem hiding this comment.
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.
| 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))?; |
There was a problem hiding this comment.
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.
| let (rebuilt_raw, _path) = objectstore::s3::create_store_with_bridge( | ||
| &url_for_rebuild, | ||
| &configs_for_rebuild, | ||
| rebuilt_bridge, | ||
| Duration::from_secs(300), | ||
| )?; |
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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) { |
There was a problem hiding this comment.
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( |
There was a problem hiding this comment.
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") |
There was a problem hiding this comment.
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. |
There was a problem hiding this comment.
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") { |
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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.
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.
|
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 How each finding is addressed:
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. |
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.
| factory_bucket.as_str(), | ||
| credential_path, | ||
| AccessMode::Read, | ||
| &HashMap::new(), |
There was a problem hiding this comment.
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.
There was a problem hiding this comment.
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 { |
There was a problem hiding this comment.
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.
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
left a comment
There was a problem hiding this comment.
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
CometS3LocationScopedCredentialProvidersupplies 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
SparkPathand 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, onlyPermissionDeniedtriggers refresh. A provider exception becomesGeneric, 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 bothget_optsandget_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
left a comment
There was a problem hiding this comment.
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
CometS3LocationScopedCredentialProvidersupplies 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, onlyPermissionDeniedtriggers refresh. Provider exceptions becomeGeneric, 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. Bothget_optsandget_rangesreproduced 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
left a comment
There was a problem hiding this comment.
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
CometS3LocationScopedCredentialProvidersupplies 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, onlyPermissionDeniedtriggers refresh. Provider exceptions becomeGeneric, 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 ofget_optsandget_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
left a comment
There was a problem hiding this comment.
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
left a comment
There was a problem hiding this comment.
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:
CometS3LocationScopedCredentialProvidersupplies 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, onlyPermissionDeniedtriggers refresh. Provider exceptions becomeGeneric, 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 ofget_optsandget_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.
* 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>
Which issue does this PR close?
Closes #6207, which proposes the new
@PublicAPI 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). ACometS3CredentialProvidertherefore 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 forwarehouse/sales, another forwarehouse/finance) get 403s on every location but the first.What changes are included in this PR?
CometS3LocationScopedCredentialProvider(new@Publicopt-in extension ofCometS3CredentialProvider) 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.getPolicyLocationsreturnsnullfor providers that do not implement the interface, without calling them, and otherwise aString[]copy, so provider list code runs inside the checked JNI call. Anulllist 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 callsblock_oninside 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.create_storebuilds 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.RetryOn403ObjectStore, the scope-hint JNI entry point, and the multi-entry object store cache from the earlier revision.How are these changes tested?
parquet::objectstore::location_scoped):salesvssales_eu), percent-decoding (Spark's%3Aescapes), duplicate and root locations, and invalid locations (empty and..segments, URIs, control characters, invalid UTF-8).get_optsandget_rangessharing 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.parquet::objectstore::s3: building a store from a template inside the Tokio runtime with no endpoint or region configured.CometS3LocationScopedCredentialProviderTest): base providers returnnullwithout being called; empty,null,null-element, non-String, lazy-list and provider exceptions; unknown handles.CometPublicApiSuitepins the new@Publictype.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 andCometPublicApiSuiteabove, spotless, scalastyle, and Prettier pass locally.