snmvaughan opened a new pull request, #6478:
URL: https://github.com/apache/datafusion-comet/pull/6478
## Which issue does this PR close?
No issue was filed for this separately. It extends #6031 (for #6207), which
gave native Parquet reads one credential per policy location, to the native
Iceberg path.
## Rationale for this change
#6031 gave native Parquet reads one credential per policy location and left
the Iceberg path unchanged. That path still gives each table one credential:
- `load_file_io` builds each `FileIO` with a credential bridge bound to one
reference path, which is the metadata location for scans and the data location
for writes.
- `iceberg-storage-opendal` attaches that same loader to the operator it
builds for every file.
So every data and delete file a scan reads, and every file a write produces,
is signed with the reference path's credential. That breaks for tables whose
files span locations with different policies:
- data outside the table location;
- a narrower policy nested under the table location;
- files in another bucket.
Reads of those files get 403s. Worse, if a broader credential covers a path
that a narrower policy restricts, the read succeeds even though the provider
would refuse it.
## What changes are included in this PR?
- **`LocationScopedS3StorageFactory` and `LocationScopedS3Storage`**
(native, new file `execution/operators/iceberg_location_scoped.rs`). These
replace `OpenDalStorageFactory::S3` when the configured provider implements
`CometS3LocationScopedCredentialProvider`.
- **Routing.** Every iceberg-rust `Storage` method receives the path it
operates on. The storage routes each call to the longest location covering that
path, in that path's bucket, and hands it to an OpenDAL S3 storage whose
credential loader is a bridge bound to that location. The key is the one
`iceberg-storage-opendal` asks S3 for: the rest of the path after
`scheme://bucket/`, normalized as opendal normalizes every path (whitespace
trimmed, leading and empty segments dropped). Iceberg paths aren't URIs, so the
key isn't percent-decoded: a partition escape such as the `%3A` in
`ts=2024-01-01T00%3A00` is part of it.
- **Sharing.** Bucket snapshots and location storages are built on first
use and shared by every `FileIO` of the provider registration and access mode,
keyed like the `FileIO` cache but without the metadata path. A new commit or
another table in the catalog therefore reuses them instead of calling
`getPolicyLocations` again, and a refresh through one `FileIO` serves the
others. Concurrent builds for one registration wait for the first. The registry
holds weak references, so the state lives as long as a cached `FileIO` uses it.
- **Retry on 403 or a signing failure.** `new_input` and `new_output`
return files bound to this storage, and `reader` returns a `FileRead` that
routes each range read at the current snapshot and reopens the file when that
routes it elsewhere. So any request that gets a 403, or that can't be signed
because the provider didn't produce its location's credential, fetches the
bucket's locations again, once per snapshot generation, and retries once if its
path now routes somewhere else, as on the Parquet path. A provider exception
never reaches S3: reqsign's chain turns it into "no credential", and
reqsign-core 3.3's signer then fails with `CredentialInvalid` without sending
the request, so that error is the second trigger. Reads, `write` and the
deletes all retry this way. `delete_stream` deletes a failed batch again
through its new routes. A streaming `writer` refreshes but doesn't retry, since
it can't send again what it already handed to a failed upload; Spark's task
retry then w
rites through the new route.
- **Shared routing** (new file `cloud/s3/policy_locations.rs`).
`LocationIndex` and the generation-shared refresh move out of
`location_scoped.rs` into `PolicyLocations`, which both stores now use. The
Parquet store's behavior is unchanged.
- **Wiring** (`iceberg_common.rs`). `build_s3_credential_loader` becomes
`build_s3_access`. It asks the reference bridge for its locations, then returns
either a single loader (base providers, IRSA, or the default chain) or the
location-scoped factory.
- If a location-scoped provider's `getPolicyLocations` throws, the read or
write fails, as on the Parquet path.
- For an S3-compatible alias scheme, `BlobHostPromotingS3Storage` wraps
the new storage.
- **`CometS3CredentialBridge::for_location`.** Builds a bridge for another
bucket and path that shares the reference bridge's provider registration, so
files in another bucket don't trigger a second `ensureInitialized`. On this
path the dispatch key is the catalog name, or the reference bucket when there
is none, and the provider is always given the bucket it's asked about, so one
registration serves every bucket.
- **`opendal` is no longer an optional dependency** of `native/core`. The
storage recognizes opendal's `PermissionDenied` and reqsign's
`CredentialInvalid` in the iceberg error's source chain.
`iceberg-storage-opendal` already builds the crate and `Cargo.lock` is
unchanged. `hdfs-opendal` now enables `opendal/services-hdfs`.
- **Docs.** Javadoc for `CometS3LocationScopedCredentialProvider` and the
dispatcher, the user guide, and a new section in the design notes. The error
message fidelity caveat now names the error a provider exception becomes
(`failed to load signing credential`) instead of an anonymous request. The docs
also record a limitation that predates this PR: a table whose metadata location
is hostless (`blob:///bucket/...`) never uses the configured provider, because
Comet finds no bucket to build its bridge from.
Providers that implement only the base interface are unaffected. The only
addition on their path is one dispatcher call when a `FileIO` is built, and
that call never reaches the provider.
## How are these changes tested?
- **New Rust unit tests** in
`execution::operators::iceberg_location_scoped`, using fake inner storages.
They cover:
- routing each file to the longest covering location, respecting
path-segment boundaries and falling back to the bucket root;
- routing files in a second bucket by that bucket's own locations;
- every `Storage` method routing by its own path, including how
`delete_stream` groups paths;
- input files and readers going through the routing;
- a 403 on a location added after the snapshot, both for a direct read and
through a reader;
- unchanged locations, a failed refresh, and non-403 errors;
- routing by the key S3 receives, with a percent escape, a `#`, and an
empty segment that opendal drops;
- a location the provider stops vending after the snapshot, whose signing
failure triggers the refresh;
- a reader that fails after another request refreshed the locations still
fetching them again, and concurrent range reads sharing one refresh whether or
not it changes their route;
- `write` retrying on a new location, a failed streaming writer refreshing
for the task's next attempt, and `delete_stream` deleting a failed batch again
through its new routes;
- a real OpenDAL S3 storage whose loader throws, which pins that the
failure is a signing error and not a 403, and that the storage's opendal is the
one Comet downcasts to;
- `FileIO`s of one registration sharing a refresh and building each
location's storage once, and a new bucket being fetched only once;
- the registry reusing a live value, dropping the slots of values that are
gone, and creating a value once when builds race;
- invalid locations being rejected, and a 403 or a signing failure being
recognized through iceberg's error.
- **Existing Rust tests.** The `LocationIndex` tests moved unchanged to
`cloud::s3::policy_locations`. The `parquet::objectstore::location_scoped`
tests pass unchanged on top of the shared refresh.
- **`CometS3CredentialBridgeSuite`** (MinIO, run manually). It creates an
Iceberg table in the scoped bucket, partitioned so its `eu` files fall under a
narrower location, reads it with the native scan, and checks which credential
path each file requested. The suite needs Docker and has not been run on this
revision yet.
These pass locally: `cargo fmt`, `cargo clippy --all-targets --workspace --
-D warnings`, `cargo test -p datafusion-comet --lib` (605 passed), spotless,
and Prettier. `./mvnw test-compile` also passes (JDK 17), so the `spark` module
and its tests compile. As a check on the new tests, disabling the 403 check
makes the six retry and 403 tests fail, routing by the percent-decoded or
unnormalized key fails `routes_by_the_key_s3_receives`, ignoring signing
failures fails the three signing-failure tests, refreshing against the reader's
opening snapshot fails `a_reader_refreshes_after_a_refresh_it_did_not_see`,
routing a reader again only after it fails fails
`concurrent_range_reads_share_one_refresh_that_changes_nothing`, dropping the
write, writer or `delete_stream` refresh fails its test, and handing out a
fresh registry slot on every call fails the three registry tests.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]