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


##########
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:
   Thanks, good find. Rather than dropping the lock for credentials the cache 
can't keep, I kept one provider call at a time per bridge and let requests that 
waited for it share its outcome, a credential or an error. A burst now costs 
one call instead of one after another, which covers your eight-thread case. A 
bridge stands for one policy location on a location-scoped store, so requests 
for different locations still don't wait on each other. A request that arrives 
after a call finished still asks again when the credential can't be kept, so a 
provider that reports `0` still sees every request that doesn't overlap 
another. `concurrent_requests_share_a_fetch_they_cannot_keep` covers unknown, 
never-expiring and nearly expired credentials and a provider error, and without 
the sharing the burst makes eight calls. Fixed in d9e5ff4f2.



##########
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:
   Thanks, I missed that test. It now asserts the expiry its properties supply, 
and every JUnit test in `org.apache.comet.cloud.s3` passes locally. Fixed in 
d9e5ff4f2.



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