andygrove commented on code in PR #6478:
URL: https://github.com/apache/datafusion-comet/pull/6478#discussion_r4149193492


##########
native/core/src/execution/operators/iceberg_common.rs:
##########
@@ -348,30 +440,106 @@ pub(crate) fn build_s3_credential_loader(
     } else {
         catalog_name
     };
-    let bridge = CometS3CredentialBridge::new(
-        provider_class,
-        dispatch_key,
-        bucket,
-        url.path(),
-        access_mode,
-        catalog_properties,
-    );
-    match bridge {
-        Ok(b) => Ok((Some(CustomAwsCredentialLoader::new(b)), true)),
-        Err(e) => match access_mode {
-            AccessMode::Write => Err(DataFusionError::Execution(format!(
-                "Configured S3 credential provider {provider_class} failed to 
initialize: {e}; \
-                 refusing to write through the default opendal credential 
chain"
-            ))),
-            AccessMode::Read => {
-                log::warn!(
-                    "Failed to initialize CometS3CredentialBridge for 
{provider_class}: {e}; \
-                     falling back to default opendal credential chain"
-                );
-                Ok((None, false))
+    // A location-scoped provider's locations belong to its registration, not 
to this table, so
+    // every FileIO of the registration shares them, however many tables and 
commits it spans.
+    // Builds for one registration run one at a time: a registration already 
known to be
+    // location-scoped needs no bridge and no provider call, and otherwise the 
first build asks the
+    // provider while the rest wait for its answer.
+    let key = RegistrationKey::new(access_mode, dispatch_key, 
catalog_properties);
+    LOCATION_SCOPED.with_slot(&key, |slot| {
+        if let Some(shared) = slot.upgrade() {
+            shared
+                .ensure_bucket(bucket)
+                .map_err(|e| DataFusionError::Execution(e.to_string()))?;
+            return Ok((S3Access::LocationScoped(shared), true));
+        }
+        let bridge = CometS3CredentialBridge::new(
+            provider_class,
+            dispatch_key,
+            bucket,
+            url.path(),
+            access_mode,
+            catalog_properties,
+        );
+        match bridge {
+            Ok(bridge) => match bridge.policy_locations() {
+                Ok(None) => Ok((
+                    
S3Access::Loader(Some(CustomAwsCredentialLoader::new(bridge))),
+                    true,
+                )),
+                Ok(Some(locations)) => {
+                    let shared = Arc::new(location_scoped_state(bridge, 
bucket, locations)?);
+                    *slot = Arc::downgrade(&shared);
+                    Ok((S3Access::LocationScoped(shared), true))
+                }
+                Err(e) => Err(DataFusionError::Execution(format!(
+                    "Failed to get policy locations for {bucket} from 
{provider_class}: {e}"
+                ))),
+            },
+            Err(e) => match access_mode {
+                AccessMode::Write => Err(DataFusionError::Execution(format!(
+                    "Configured S3 credential provider {provider_class} failed 
to initialize: \
+                     {e}; refusing to write through the default opendal 
credential chain"
+                ))),
+                AccessMode::Read => {
+                    log::warn!(
+                        "Failed to initialize CometS3CredentialBridge for 
{provider_class}: {e}; \
+                         falling back to default opendal credential chain"
+                    );
+                    Ok((S3Access::Loader(None), false))
+                }
+            },
+        }
+    })
+}
+
+/// The shared locations of a location-scoped provider's registration, seeded 
with `locations` for
+/// `bucket`. A bucket's locations come from a bridge derived from `bridge` 
for that bucket, and
+/// each location's storage signs with a bridge derived for that location, so 
neither calls
+/// `ensureInitialized` again.
+fn location_scoped_state(
+    bridge: CometS3CredentialBridge,
+    bucket: &str,
+    locations: Vec<String>,
+) -> Result<SharedLocations, DataFusionError> {
+    let bridge = Arc::new(bridge);
+    let source_bridge = Arc::clone(&bridge);
+    let source: BucketLocationSource = Arc::new(move |bucket: &str| {
+        // A blocking JVM call, usually made inside an async storage call on a 
Tokio worker.
+        tokio::task::block_in_place(|| {

Review Comment:
   When a write task fed by JVM input is dropped mid-flight, 
`AbortOnDrop::drop` in `iceberg_write.rs` runs `delete_task_files` on a 
current-thread runtime. With this storage, a cleanup `delete` that gets a 403 
or a signing failure refreshes, reaches this `block_in_place`, and panics with 
`can call blocking only when running on the multi-threaded runtime`. 
`releasePlan` then throws, and the rest of the task's files are left behind. 
This only affects the opt-in native writer, but it is the failure this PR 
targets: a location's credential failing mid-job. Could this call 
`block_in_place` only on a multi-thread runtime, for example by checking 
`Handle::try_current()` and `runtime_flavor()`? A test that drives a refresh 
from a current-thread runtime would pin it.



-- 
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