parthchandra commented on code in PR #6023:
URL: https://github.com/apache/datafusion-comet/pull/6023#discussion_r4117802997


##########
spark/src/main/spark-3.x/org/apache/comet/cloud/s3/HadoopS3ACredentialProviderAdapter.java:
##########
@@ -0,0 +1,119 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.comet.cloud.s3;
+
+import java.net.URI;
+import java.util.Map;
+import java.util.concurrent.ConcurrentHashMap;
+
+import org.apache.hadoop.conf.Configuration;
+import org.apache.hadoop.fs.s3a.S3AUtils;
+
+import com.amazonaws.auth.AWSCredentialsProvider;
+
+import org.apache.comet.util.ClassLoaders;
+
+/**
+ * Delegates credential resolution to Hadoop S3A's own provider construction, 
so it accepts
+ * everything the {@code fs.s3a.aws.credentials.provider} chain accepts. This 
is the spark-3.x (AWS
+ * SDK v1) body; it calls {@link S3AUtils#createAWSCredentialProviderSet} and 
returns v1
+ * credentials.
+ *
+ * <p>Enable it (leaving {@code fs.s3a.aws.credentials.provider} untouched) 
with:
+ *
+ * <pre>
+ * 
spark.hadoop.fs.s3a.comet.credential.provider.class=org.apache.comet.cloud.s3.HadoopS3ACredentialProviderAdapter
+ * </pre>
+ */
+public class HadoopS3ACredentialProviderAdapter implements 
CometS3CredentialProvider {
+
+  private Map<String, String> properties;
+  // Captured on the thread that runs initialize() (the dispatcher calls it 
during planning, which
+  // has Spark's user-jar loader); native worker threads have a null context 
loader. Set on the
+  // Configuration so S3A's factory loads the named provider classes from it. 
This works on Hadoop
+  // 3.3.4 because it loads them through conf.getClasses, which honors the 
conf's loader.
+  private volatile ClassLoader classLoader;
+  // One delegate per bucket: on the Iceberg path the dispatch key is the 
catalog, so a single
+  // instance can serve multiple buckets; the Parquet path is per-bucket and 
uses a single entry.
+  private final ConcurrentHashMap<String, AWSCredentialsProvider> delegates =
+      new ConcurrentHashMap<>();
+
+  @Override
+  public void initialize(Map<String, String> catalogProperties) {
+    this.properties = catalogProperties;
+    this.classLoader = 
ClassLoaders.contextOrDefault(getClass().getClassLoader());
+  }
+
+  @Override
+  public CometS3Credentials getCredentialsForPath(CometS3CredentialContext 
context)
+      throws Exception {
+    AWSCredentialsProvider provider = ensureDelegate(context.getBucket());
+    return 
SdkCredentialExtraction.toCometCredentials(provider.getCredentials());
+  }
+
+  private AWSCredentialsProvider ensureDelegate(String bucket) throws 
Exception {
+    AWSCredentialsProvider existing = delegates.get(bucket);
+    if (existing != null) {
+      return existing;
+    }
+    synchronized (this) {
+      AWSCredentialsProvider delegate = delegates.get(bucket);
+      if (delegate == null) {
+        delegate = buildDelegate(bucket);
+        delegates.put(bucket, delegate);
+      }
+      return delegate;
+    }
+  }
+
+  private AWSCredentialsProvider buildDelegate(String bucket) throws Exception 
{
+    Configuration conf =
+        
S3AUtils.propagateBucketOptions(AdapterSupport.toConfiguration(properties), 
bucket);
+    if (classLoader != null) {
+      // Hadoop 3.3.4's factory loads named providers through conf.getClasses, 
which honors this
+      // loader, so a provider on the user-jar loader resolves even from a 
null-context worker
+      // thread.
+      conf.setClassLoader(classLoader);

Review Comment:
   Fixed. The spark-3.x body now does 
conf.setClassLoader(S3AFileSystem.class.getClassLoader()) just like 
S3AFileSystem.initialize, and I dropped the classLoader field. So named 
providers load through hadoop-aws's own loader on 3.3.4 too, and the 
userClassPathFirst bug stays fixed. Reworked the test into 
pinsHadoopAwsLoaderForNamedProvider (child-first case, fails at old head, 
passes now), removed the parenthetical from the spark-4.x javadoc, and left the 
two SDK adapters with their capture since S3A doesn't load their delegate.



##########
spark/src/main/spark-4.x/org/apache/comet/cloud/s3/SdkCredentialExtraction.java:
##########
@@ -0,0 +1,50 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.comet.cloud.s3;
+
+import java.time.Instant;
+import java.util.Optional;
+
+import software.amazon.awssdk.auth.credentials.AwsCredentials;
+import software.amazon.awssdk.auth.credentials.AwsSessionCredentials;
+
+/**
+ * Maps AWS SDK v2 {@link AwsCredentials} onto {@link CometS3Credentials}. 
Compiled only into
+ * spark-4.0+ builds (spark-4.x source set), directly against SDK v2.
+ */
+final class SdkCredentialExtraction {
+
+  private SdkCredentialExtraction() {}
+
+  static CometS3Credentials toCometCredentials(AwsCredentials creds) {
+    String sessionToken = null;
+    long expirationEpochMillis = 0L;
+    if (creds instanceof AwsSessionCredentials) {
+      AwsSessionCredentials session = (AwsSessionCredentials) creds;
+      sessionToken = session.sessionToken();
+      Optional<Instant> expiration = session.expirationTime();
+      if (expiration.isPresent()) {
+        expirationEpochMillis = expiration.get().toEpochMilli();
+      }
+    }
+    return new CometS3Credentials(

Review Comment:
   Fixed. Both extraction bodies now detect null-or-empty keys and throw a 
clear error naming the bucket instead of the NPE. Added 
refusesAnonymousCredentials to both adapter tests plus an empty-string case. 
The empty per-bucket class opt-out is intended and supported — I pointed the 
error and the user guide at it, and added 
test_empty_per_bucket_provider_class_opts_out. Filed 
https://github.com/apache/datafusion-comet/issues/6298 for real anonymous 
support through the SPI.



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