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


##########
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:
   I steered you wrong on this body last round, sorry. I said it needs the 
captured loader because Hadoop 3.3.4 loads the named providers through 
`conf.getClasses`. That part is true. But `S3AFileSystem.initialize` pins the 
conf's loader before it builds the chain, at `S3AFileSystem.java:442-445` in 
3.3.4:
   
   ```java
   // fix up the classloader of the configuration to be whatever
   // classloader loaded this filesystem.
   // See: HADOOP-17372
   conf.setClassLoader(this.getClass().getClassLoader());
   ```
   
   So plain Spark loads them through hadoop-aws's own loader on 3.3.4 too, the 
same as on 3.4. Putting the task's loader here instead brings back the bug 
HADOOP-17372 fixed. If that loader has its own copy of the AWS SDK, 
`createAWSCredentialProvider` gets a class that implements the other copy of 
`AWSCredentialsProvider` and rejects it. That's what happens with 
`spark.executor.userClassPathFirst=true` and an application jar that bundles 
the SDK.
   
   I reproduced it on 3.5. I used Spark's `ChildFirstURLClassLoader` over the 
`aws-java-sdk-core` jar as the context loader for `initialize()`, with 
`SystemPropertiesCredentialsProvider` as the named provider. A plain 
`S3AFileSystem.initialize` on that same thread resolves the keys. The adapter's 
first fetch fails with `IOException: Class class 
com.amazonaws.auth.SystemPropertiesCredentialsProvider does not implement 
AWSCredentialsProvider`. `DefaultAWSCredentialsProviderChain` from #6022 
behaves the same way.
   
   Could this body mirror `initialize` with 
`conf.setClassLoader(S3AFileSystem.class.getClassLoader())` and drop the 
`classLoader` field? With that one change the probe resolves and the other five 
tests here still pass. `usesLoaderCapturedAtInitializeOnNullContextThread` 
would need to become the child-first case, which fails at this head and passes 
with the pin. The parenthetical about the spark-3.x body in the spark-4.x 
javadoc can go as well, since both lines would then load the way Spark does. 
The two SDK adapters should keep their capture, because nothing in S3A loads 
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:
   What should happen when the chain resolves anonymous credentials? The guide 
has people name the Hadoop adapter globally. A bucket configured with just 
`AnonymousAWSCredentialsProvider` is common for public datasets, and the native 
reader handles it today by skipping the signature. With the adapter on 
globally, that bucket goes through the adapter instead. Both 
`AWSCredentialProviderList` versions deliberately return the anonymous 
credentials, which have null keys, and this constructor then throws 
`NullPointerException: accessKeyId`. The spark-3.x body does the same.
   
   I checked it on 3.5 and 4.1 with 
`fs.s3a.bucket.public-data.aws.credentials.provider` set to 
`AnonymousAWSCredentialsProvider`, and again with 
`SimpleAWSCredentialsProvider,AnonymousAWSCredentialsProvider` and no keys. All 
four cases fail the same way. Plain Spark reads that bucket anonymously, so 
someone who follows the guide turns a working public-bucket read into an opaque 
NPE.
   
   The SPI has no way to say "anonymous", so I don't think the adapter can 
match Spark here without an SPI change. Could the adapters detect the anonymous 
result and fail with a message that names the bucket and the cause? A test with 
the lone anonymous provider in each `HadoopS3ACredentialProviderAdapterTest` 
would pin it. Could the user guide's Hadoop adapter section mention it too? Its 
line about resolving "the same way it would under Spark" isn't true for these 
buckets yet.
   
   Separately, an empty 
`fs.s3a.bucket.<bucket>.comet.credential.provider.class` already opts a bucket 
out, since `lookup_provider_class` finds the per-bucket key and then filters 
out the empty value. Is that meant to be supported? If so, it would be the 
simplest thing for the error message and the guide to point at. And if you 
think anonymous support through the SPI is worth adding, could you file an 
issue and link it here?



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