snmvaughan commented on issue #6462:
URL:
https://github.com/apache/datafusion-comet/issues/6462#issuecomment-5914201613
# Design sketch: location-scoped credentials on the Iceberg path
Code references are against `apache/datafusion-comet` `main` at `9603ad128`
and `apache/iceberg-rust` at `bb1e4a4`, the revision Comet pins.
## Current behavior
The Parquet path routes each request to a policy location
(`parquet/objectstore/location_scoped.rs`). The Iceberg path does not:
- `load_file_io` (`execution/operators/iceberg_common.rs:226`) builds a
`FileIO` from a reference path. That path is the metadata location for scans
(`iceberg_scan.rs:181`) and the data location for writes
(`iceberg_write.rs:484`).
- `build_s3_credential_loader` (`iceberg_common.rs:308`) creates one
`CometS3CredentialBridge` from the reference path's bucket and `url.path()`
(`iceberg_common.rs:351`), then wraps it in a `CustomAwsCredentialLoader`.
- iceberg-rust's `s3_config_build` (`crates/storage/opendal/src/s3.rs:129`)
builds an operator for each file operation. It takes the bucket from the file's
path but pushes the same loader onto every operator's credential chain
(`s3.rs:149`).
- The bridge's `provide_credential` (`cloud/s3/credential_bridge.rs:376`)
always calls `getCredentialsForPath` with the reference bucket and path.
- The executor `FileIO` cache (`iceberg_common.rs:108`, 64 entries) keeps
`reference_path` in its key, so two tables never share a bridge.
The result is one credential identity per table and access mode. The
bridge's own doc comment describes this as "per-table-location (Iceberg)"
(`credential_bridge.rs:69`).
## Goals
- Every file a native Iceberg scan or write touches is served with the
credential of the longest policy location covering that file's path, in that
file's bucket. The provider supplies the locations through the existing
`getPolicyLocations`.
- The same refresh semantics as the Parquet path: after a 403, fetch the
bucket's locations again, at most once per snapshot generation, and retry once
if the path now routes to a different location.
- No change for providers that implement only `CometS3CredentialProvider`.
- No change to the public Java API.
## Non-goals
- Changing iceberg-rust or opendal. Both have the extension points this
needs (see Alternatives below).
- Caching credentials in Comet. As on the Parquet path, every request still
reaches `getCredentialsForPath`; what is kept is the location list and one
inner storage per location that has been used.
## Design
### Where it plugs in
`storage_factory_for` (`iceberg_common.rs:57`) already substitutes a Comet
factory for S3-compatible aliases. `BlobHostPromotingS3StorageFactory` and
`BlobHostPromotingS3Storage` (`parquet/objectstore/s3_blob_fs_support.rs:181`,
`:210`) wrap `OpenDalStorageFactory::S3` and delegate every `Storage` method by
path. The new pair follows the same pattern:
- `LocationScopedS3StorageFactory` is `#[typetag::serde(name =
"CometLocationScopedS3StorageFactory")]`. Its runtime state is
`#[serde(skip)]`, as `customized_credential_load` already is. Serde exists only
to satisfy the `Storage` supertraits; Comet never serializes storage during a
scan.
- `LocationScopedS3Storage` implements `Storage` (iceberg-rust
`crates/iceberg/src/io/storage/mod.rs:72`).
`build_s3_credential_loader` creates the reference bridge as it does today,
then asks it for `policy_locations()`, as `create_store` does on the Parquet
path.
- `None` means a base provider. Keep today's single-loader
`OpenDalStorageFactory::S3`.
- `Some(locations)` means a location-scoped provider. Return the new
factory, seeded with the reference bucket's snapshot.
For an S3-compatible alias scheme, the host-promotion wrapper stays
outermost, so routing sees the promoted host as the bucket.
### Routing
State per storage instance:
- A per-bucket map of location snapshots. Each snapshot is the Parquet
path's `LocationIndex`: canonical location mapped to credential path, plus a
generation counter and the last refresh failure.
- A per-`(bucket, credential_path)` map of inner `Arc<dyn Storage>`.
For each call, the storage:
1. Parses the call's `path` into a bucket and key.
2. Looks up or fetches the bucket's snapshot.
3. Routes the key with `LocationIndex::route` (`location_scoped.rs:110`):
longest covering location, compared segment by segment after percent-decoding,
with `/` as the implicit root.
4. Delegates to the inner storage for that `(bucket, credential_path)`.
`LocationIndex` and `route` are currently private to `location_scoped.rs`
and keyed by `object_store::path::Path`. Move them to a small shared module
used by both paths. Keep `Path` for canonicalization, since it is already the
decoding rule the SPI documents.
### Inner storages and bridges
An inner storage is `OpenDalStorageFactory::S3 { customized_credential_load:
Some(loader) }.build(config)`. Its `loader` wraps a bridge bound to `(bucket,
credential_path)`. It is built on first use and kept for the life of the outer
storage.
Bridges are derived from the reference bridge, so they reuse its dispatcher
handle and trigger no `ensureInitialized` call or class loading on the Tokio
worker that usually builds them:
- Within the reference bucket, derive with the existing `for_path`
(`credential_bridge.rs:144`).
- For another bucket, add a `for_location(bucket, path)` derivation. On the
Iceberg path the dispatch key is the catalog name, so one registration already
serves every bucket a table touches.
### Snapshot lifecycle and blocking
- The reference bucket's snapshot is fetched in
`build_s3_credential_loader`. That runs on the thread that calls
`load_file_io`, the same place the Parquet path's `create_store` fetches.
- A snapshot for any other bucket is fetched lazily, inside an async
`Storage` method, on a Tokio worker. That is a blocking JVM call of the kind
[#6293](https://github.com/apache/datafusion-comet/issues/6293) describes, so
wrap it in `tokio::task::block_in_place` as
[#6261](https://github.com/apache/datafusion-comet/pull/6261) does for
`acquireMemory`. The per-request `getCredentialsForPath` calls are already in
[#6293](https://github.com/apache/datafusion-comet/issues/6293)'s list and are
unchanged here.
- A snapshot lives as long as its storage, and the storage lives as long as
its `FileIO` in the executor cache. A cache hit therefore reuses the snapshot
until a 403 refreshes it, or until the entry is evicted and rebuilt.
### Refresh and retry
The semantics mirror `LocationScopedObjectStore`: `refresh` at
`location_scoped.rs:225` and the per-generation sharing described in the
contributor guide.
- **Trigger.** On this path the only trigger is a 403. The Parquet path also
refreshes on `CredentialProviderError`, but here reqsign's
`ProvideCredentialChain` swallows provider exceptions. This is the existing
"error message fidelity" caveat: the request goes out unsigned and S3 rejects
it with 403. The storage must detect permission-denied through the source chain
of iceberg's `Error`, which wraps opendal's `ErrorKind::PermissionDenied`.
- **Retried calls.** `exists`, `metadata`, `read`, and each `FileRead::read`.
- **Readers.** `reader()` and `new_input()` must stay retryable after they
return. `new_input(path)` therefore returns `InputFile::new(Arc<Self>, path)`
(iceberg-rust `file_io.rs:331`) rather than the inner storage's `InputFile`.
`reader(path)` returns a `FileRead` wrapper that holds the current route. On a
403 it refreshes, reopens through the new route if the route changed, and
retries the range once.
- **Budget.** The retry budget is per request, as on the Parquet path.
### Writes and deletes
- `write`, `writer`, and `new_output` route by path with `AccessMode::Write`
bridges, and do not retry, matching the Parquet store's non-read operations.
- `delete` and `delete_stream` route per path. `delete_stream` groups paths
by route before delegating.
- `delete_prefix` routes by the prefix itself, as `list` does on the Parquet
path. If the prefix spans several locations, the single credential may not
cover all of them. That is acceptable for the current callers, but see the open
questions.
### FileIO cache
The cache key stays as it is. The base provider still needs `reference_path`
in the key, and keeping it for location-scoped providers is harmless: at worst,
each cached `FileIO` holds its own snapshot. A follow-up could drop
`reference_path` from the key for location-scoped providers, so tables in one
catalog share a snapshot.
### Documentation to update
- `CometS3LocationScopedCredentialProvider` Javadoc, line 43: "Locations
apply to Comet's native Parquet reads. Iceberg reads do not use them."
- User guide `s3-credential-providers.md`, line 300: "Locations apply to
Comet's native Parquet reads only".
- Contributor guide `s3-credential-provider-design.md`: line 91 ("The
Iceberg path does not use locations"), the "Path-specific behavior" table, and
a new "Location-scoped credentials on the Iceberg path" section based on this
document.
- The `CometS3CredentialBridge` doc comment's "per-table-location (Iceberg)"
(`credential_bridge.rs:69`).
## Testing
**Rust unit tests** (`native/core`), using a fake inner `Storage` that
records `(bucket, credential_path, operation)`:
- Every `Storage` method routes by its path, including longest-match, root
fallback, and percent-decoded paths.
- `new_input` returns a file bound to the wrapper, and its reads route and
retry.
- Paths in a second bucket get their own snapshot through `for_location`,
with no second `ensureInitialized`.
- After a 403, one refresh per generation serves concurrent failures. A
retry happens only when the route changes, and a failed refresh fails every
request from that generation.
- `FileRead` reopens through the new route on a 403.
- Writes and deletes route by path without retry.
- A base provider (dispatcher returns `null`) still gets today's
single-loader factory.
**Integration tests** (`CometS3CredentialBridgeSuite`, MinIO,
`MinioLocationScopedCredentialProvider`). MinIO does not enforce per-prefix
policies here, so, as the Parquet test does, assert on the recorded credential
paths rather than on 403s.
- Create an Iceberg table in the scoped bucket, partitioned so its data
files fall under two installed locations, one of them nested under the table
location. Read it with the native scan. Assert that `credentialPaths()`
contains each location's credential path, not only the metadata file's.
- Create a table with `write.data.path` under a different installed
location, and assert its reads use that location's credential path.
- If native Iceberg writes are enabled in the suite, write both tables and
assert on the write-mode credential paths.
- A regression test that the base-provider Iceberg test
(`CometS3CredentialBridgeSuite.scala:108`) is unchanged.
## Alternatives
1. **A path-aware loader in iceberg-rust.** `s3_config_build` already
receives the path (`s3.rs:129`). `iceberg-storage-opendal` could accept a
loader factory keyed by bucket and path instead of one loader. That is cleaner
for other embedders, but it needs an upstream release and a revision bump, and
Comet would still own routing and refresh. It could be adopted later without
changing the SPI.
2. **Keep one credential per table and document the limitation.** This
leaves the multi-policy cases broken, including the case where a broader
credential reads a path that a narrower policy restricts.
3. **Have providers return a broader credential for the reference path.**
This pushes least-privilege violations onto vendors, and is not possible for
vendors whose policies differ within a table.
## Open questions
1. Should Iceberg routing be on by default for location-scoped providers, or
behind a config for one release? It extends documented behavior, although it
stays within the provider contract.
2. Should `reference_path` be dropped from the `FileIO` cache key for
location-scoped providers?
3. Should `delete_prefix` fan out across the locations under the prefix?
This matters only if a caller deletes a prefix that spans policies.
--
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]