sunchao commented on code in PR #6509:
URL: https://github.com/apache/datafusion-comet/pull/6509#discussion_r4156726091


##########
spark/src/test/spark-4.x/java/org/apache/comet/cloud/s3/IcebergRESTVendedS3Provider.java:
##########
@@ -62,11 +63,16 @@ public CometS3Credentials 
getCredentialsForPath(CometS3CredentialContext context
               + "Comet should always invoke initialize before 
getCredentialsForPath");
     }
     AwsCredentials c = p.resolveCredentials();
-    String sessionToken =
-        (c instanceof AwsSessionCredentials) ? ((AwsSessionCredentials) 
c).sessionToken() : null;
-    // Expiration is owned by VendedCredentialsProvider's CachedSupplier; we 
publish 0 so the
-    // native bridge applies its conservative floor to `opendal`'s cache while 
the inner
-    // CachedSupplier handles refresh on its own schedule.
-    return new CometS3Credentials(c.accessKeyId(), c.secretAccessKey(), 
sessionToken, 0L);
+    String sessionToken = null;
+    long expirationEpochMillis = 0L;
+    if (c instanceof AwsSessionCredentials) {
+      AwsSessionCredentials session = (AwsSessionCredentials) c;
+      sessionToken = session.sessionToken();
+      // Publish the vended credential's expiry, so Comet reuses it until 
shortly before and never
+      // after. VendedCredentialsProvider's CachedSupplier still decides when 
to fetch a new one.
+      expirationEpochMillis = 
session.expirationTime().map(Instant::toEpochMilli).orElse(0L);

Review Comment:
   [P2] Update the existing vended-credential test alongside this behavior 
change. 
`IcebergRESTVendedS3ProviderTest.initializeThenGetReturnsVendedCredentials` 
supplies an expiry one hour ahead but still asserts `getExpirationEpochMillis() 
== 0` at line 65. Returning the real expiry now deterministically fails that 
test and blocks JVM CI. Preserve the new behavior and have the test assert the 
exact expiry supplied in its properties.
   
   Evidence: Exact-head CI reproduced the assertion failure in the scans, 
shuffle, exec and expressions shards, plus TPC-H. For example, scans job 
110421192667 reports 
`IcebergRESTVendedS3ProviderTest.initializeThenGetReturnsVendedCredentials:65 
expected:<0> but was:<1790869267664>` and exits with a Maven test failure: 
https://github.com/apache/datafusion-comet/actions/runs/36874557330/job/110421192667.



##########
native/core/src/cloud/s3/credential_bridge.rs:
##########
@@ -41,15 +42,94 @@ use std::fmt;
 use std::sync::Arc;
 use std::time::Duration;
 
-/// Cap on opendal's credential cache when the provider does not report an 
expiry. Prevents the
-/// executor from holding a stale credential for the entire job lifetime. 
Shared with the IRSA
-/// web-identity provider (`super::web_identity`).
+/// Expiry the Iceberg path assumes when the provider does not report one. It 
bounds how long a
+/// long-lived reader or writer reuses a credential, but it can outlast a 
short-lived credential,
+/// so a provider that knows the expiry should report it. Shared with the IRSA 
web-identity
+/// provider (`super::web_identity`).
 pub(crate) const DEFAULT_EXPIRY_WHEN_UNKNOWN: Duration = 
Duration::from_secs(300);
 
+/// How long before its expiry a bridge stops reusing a credential and asks 
the provider again. It
+/// matches the refresh-ahead of the other credential caches Comet keeps, and 
it covers object_store,
+/// which signs a request once and sends that signature again on every retry 
for up to 3 minutes by
+/// default.
+pub(crate) const REFRESH_BEFORE_EXPIRY: Duration = Duration::from_secs(300);
+
+/// The earliest `expirationEpochMillis` taken at face value, 
2000-01-01T00:00:00Z. An earlier one is
+/// almost always seconds since the epoch sent as milliseconds.
+const EARLIEST_PLAUSIBLE_EXPIRY_MILLIS: i64 = 946_684_800_000;
+
 /// 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();
 
+/// Once-per-process latch for the "implausible expiry" warning.
+static WARNED_IMPLAUSIBLE_EXPIRY: OnceCell<()> = OnceCell::new();
+
+/// When a provider's credential stops working, as its `expirationEpochMillis` 
says.
+#[derive(Debug, PartialEq, Eq)]
+enum Expiry {
+    /// `0` or negative, or too early to be a real expiry: the provider does 
not know.
+    Unknown,
+    /// Too far ahead to represent, as `Long.MAX_VALUE` is: the credential 
does not expire.
+    Never,
+    At(Timestamp),
+}
+
+impl Expiry {
+    fn from_millis(millis: i64) -> Self {
+        if millis <= 0 {
+            return Expiry::Unknown;
+        }
+        if millis < EARLIEST_PLAUSIBLE_EXPIRY_MILLIS {
+            if WARNED_IMPLAUSIBLE_EXPIRY.set(()).is_ok() {
+                warn!(
+                    "CometS3CredentialProvider returned expirationEpochMillis 
{millis}, which is \
+                     before 2000 and probably in seconds; treating the expiry 
as unknown"
+                );
+            }
+            return Expiry::Unknown;
+        }
+        Timestamp::from_millisecond(millis).map_or(Expiry::Never, Expiry::At)
+    }
+}
+
+/// A credential with a known expiry, and when to stop reusing it.
+struct CachedCredential {
+    raw: RawCredentials,
+    refresh_at: Timestamp,
+}
+
+/// A bridge's last credential with a known expiry. It is reused until 
[`REFRESH_BEFORE_EXPIRY`]
+/// before that expiry, so the provider is asked about once per credential 
rather than once per
+/// request (Parquet) or storage call (Iceberg). A credential whose expiry is 
unknown, or that does
+/// not expire, is not kept, and the provider is asked for every time.
+#[derive(Default)]
+struct CredentialCache(Mutex<Option<CachedCredential>>);
+
+impl CredentialCache {
+    /// The kept credential if it is still fresh at `now`, or else the one 
`fetch` returns, kept if
+    /// it can be. Concurrent calls wait for one fetch.
+    fn get_or_fetch<E>(
+        &self,
+        now: Timestamp,
+        fetch: impl FnOnce() -> Result<RawCredentials, E>,
+    ) -> Result<RawCredentials, E> {
+        let mut cached = self.0.lock();

Review Comment:
   [P2] Preserve concurrent provider calls when credentials cannot be cached. 
For a provider that performs a per-request REST/token exchange and reports `0`, 
this mutex stays locked across `fetch()`, then the returned credential is 
discarded. Concurrent reads sharing a bridge therefore execute every exchange 
sequentially, whereas the previous direct-fetch path allowed them to overlap. 
The same applies to `Long.MAX_VALUE` and near-expiry credentials. This 
materially increases credential-resolution latency on the documented no-cache 
path. Could that path bypass fetch serialization while retaining coordinated 
refresh for reusable credentials?
   
   Evidence: A bounded Rust reproduction used the extracted `CredentialCache`, 
eight barrier-synchronized threads and a 50 ms simulated provider call. With 
expiry `0`, direct fetching completed in 55 ms with eight concurrent calls, 
versus 405 ms and one concurrent call through the new cache. Non-expiring and 
near-expiry cases likewise took 404–406 ms instead of 54–57 ms. The JVM 
dispatcher directly invokes the provider without another serialization layer.



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