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]

Reply via email to