Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
10 changes: 8 additions & 2 deletions docs/source/contributor-guide/s3-credential-provider-design.md
Original file line number Diff line number Diff line change
Expand Up @@ -106,9 +106,15 @@ The 5-minute fallback is a safety net so a vendor that omits expiry cannot leave

## Property-bag handling on the Iceberg path

The full unfiltered FileIO property bag crosses JNI as `catalog_properties`. The storage-prefix filter (`s3.`/`gcs.`/`adls.`/`client.`) is applied native-side in `iceberg_scan.rs::load_file_io` immediately before `FileIOBuilder.with_prop`. This means the bridge sees `credentials.uri`, OAuth tokens, and any vendor-custom keys with no parallel field on the operator and no driver-side broadcast. Vendors set their own keys on the catalog config and read them back inside `initialize(Map)`.
The full unfiltered FileIO property bag crosses JNI as `catalog_properties`. The storage-prefix filter (`s3.`/`gcs.`/`adls.`/`client.`) is applied native-side in `iceberg_common.rs::build_file_io` immediately before `FileIOBuilder.with_prop`. This means the bridge sees `credentials.uri`, OAuth tokens, and any vendor-custom keys with no parallel field on the operator and no driver-side broadcast. Vendors set their own keys on the catalog config and read them back inside `initialize(Map)`.

`IcebergScanExec` derives a redacting `Debug` so plan dumps and tracing do not leak the property bag.
`IcebergScanExec` derives a redacting `Debug`, and the `FileIO` cache key's `Debug` omits the property bag, so plan dumps and tracing do not leak it.

## Executor `FileIO` cache on the Iceberg path

`load_file_io` in `iceberg_common.rs` serves clones from a per-executor cache instead of building a `FileIO` per task. An entry is keyed by access mode, catalog name, the full reference path and the whole catalog property bag, and it holds the `FileIO` together with the bridge it was built with, so a bridge and its dispatcher handle live as long as the entry. The cache holds 64 entries, evicts the least recently used one, and is drained with the Tokio runtime. The reference path has to stay in the key: the bridge is constructed with the bucket and `url.path()` of that location and the JVM provider is called with exactly that pair, so dropping the path from the key would hand one table's bridge to another table in the same bucket. Two builds are never cached: `memory:///`, whose namespace the write path uses per task, and a read whose configured provider failed to initialise and fell back to the default chain, so the next task retries the provider instead of inheriting the fallback.

This is consistent with [Why no Comet-side cache](#why-no-comet-side-cache): the cache holds the `FileIO` and its bridge, not credentials. `provide_credential` still reaches the vendor whenever `opendal`'s cached credential expires.

## Returns or throws, not a fall-through value

Expand Down
9 changes: 5 additions & 4 deletions native/core/src/cloud/s3/credential_bridge.rs
Original file line number Diff line number Diff line change
Expand Up @@ -45,8 +45,8 @@ use std::time::Duration;
/// executor from holding a stale credential for the entire job lifetime.
const DEFAULT_EXPIRY_WHEN_UNKNOWN: Duration = Duration::from_secs(300);

/// Once-per-process latch for the "missing expiry" warning. Bridges are per-scan, so a per-bridge
/// latch would re-log on every scan.
/// Once-per-process latch for the "missing expiry" warning. Bridges live as long as their entry in
/// the executor's FileIO cache, so a per-bridge latch would re-log for every new configuration.
static WARNED_MISSING_EXPIRY: OnceCell<()> = OnceCell::new();

/// Access intent forwarded to the Java SPI. Ordinal must match the JVM `CometS3AccessMode` enum.
Expand All @@ -56,8 +56,9 @@ pub enum AccessMode {
Write = 1,
}

/// Per-scan credential provider that delegates to the JVM SPI via JNI. `handle` is the JVM-side
/// identity for the `(provider_class, dispatch_key, catalog_properties)` triple returned by
/// Credential provider that delegates to the JVM SPI via JNI. Instances live in the executor's
/// FileIO and object store caches, so one serves many tasks. `handle` is the JVM-side identity
/// for the `(provider_class, dispatch_key, catalog_properties)` triple returned by
/// `ensureInitialized`. `bucket_jstr` / `path_jstr` are interned once at construction to avoid
/// per-call `new_string` allocations on the hot path.
///
Expand Down
1 change: 1 addition & 0 deletions native/core/src/execution/jni_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -425,6 +425,7 @@ pub fn get_runtime() -> Handle {
/// Must not be called from within the runtime's own worker threads, otherwise the shutdown
/// would deadlock/panic.
pub fn release_runtime() {
crate::execution::operators::clear_file_io_cache();
let runtime = TOKIO_RUNTIME.lock().take();
if let Some(runtime) = runtime {
runtime.shutdown_timeout(Duration::from_secs(3));
Expand Down
Loading
Loading