Savonitar commented on code in PR #28639:
URL: https://github.com/apache/flink/pull/28639#discussion_r3539520833


##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -416,13 +531,114 @@ long calculateRenewalDelay(Clock clock, long 
nextRenewal) {
         return renewalDelay;
     }
 
+    @VisibleForTesting
+    void setClock(Clock clock) {
+        this.clock = clock;
+    }
+
     /** Stops re-occurring token obtain task. */
     @Override
     public void stop() {
         LOG.info("Stopping credential renewal");
 
-        stopTokensUpdate();
+        synchronized (tokensUpdateFutureLock) {
+            // Mark stopped, cancel the pending cycle, and reset on-demand 
re-obtain bookkeeping
+            // atomically, so a concurrent reobtainDelegationTokens() cannot 
leave a live future
+            // orphaned after stop and a later start() does not inherit stale 
state.
+            stopped = true;
+            stopTokensUpdate();
+            reobtainScheduled = false;
+            lastReobtainAtMillis = NO_PREVIOUS_REOBTAIN;
+        }
+
+        for (DelegationTokenProvider provider : 
delegationTokenProviders.values()) {
+            try {
+                provider.stop();
+            } catch (Throwable t) {
+                LOG.error("Failed to stop delegation token provider {}", 
provider.serviceName(), t);
+            }
+        }
 
         LOG.info("Stopped credential renewal");
     }
+
+    @Override
+    public void reobtainDelegationTokens() {
+        synchronized (tokensUpdateFutureLock) {
+            if (scheduledExecutor == null || ioExecutor == null) {
+                LOG.debug(
+                        "A re-obtain of delegation tokens was requested but 
the manager was "
+                                + "constructed without executors (one-shot 
obtain path); the "
+                                + "request is ignored.");
+                return;
+            }
+            if (stopped) {
+                LOG.debug(
+                        "A re-obtain of delegation tokens was requested after 
the manager was "
+                                + "stopped; the request is ignored.");
+                return;
+            }
+            // Dedupe: if an on-demand re-obtain is already scheduled and has 
not started yet, the
+            // newly registered job(s) will be covered by it, so coalesce this 
request into it.
+            if (reobtainScheduled) {
+                LOG.debug("A re-obtain of delegation tokens is already 
scheduled; coalescing.");
+                return;
+            }
+            // Cooldown: bound how often on-demand re-obtains can run by 
deferring this cycle until
+            // at least reobtainCooldownMillis have passed since the previous 
on-demand re-obtain.
+            long now = clock.millis();
+            long delayMillis =
+                    lastReobtainAtMillis == NO_PREVIOUS_REOBTAIN
+                            ? 0L
+                            : Math.max(0L, lastReobtainAtMillis + 
reobtainCooldownMillis - now);
+            // Only bring the next cycle forward: if a cycle (e.g. the 
periodic renewal) is still
+            // pending and scheduled to fire sooner than the cooldown-deferred 
time, fire at that
+            // earlier time instead of pushing it later — otherwise a 
short-lived token could expire
+            // before it is renewed. The nextScheduledAtMillis > now guard 
skips an already-fired
+            // future that has not yet been re-armed, so this never bypasses 
the cooldown.
+            if (tokensUpdateFuture != null
+                    && nextScheduledAtMillis > now
+                    && nextScheduledAtMillis - now < delayMillis) {
+                delayMillis = nextScheduledAtMillis - now;
+            }
+            lastReobtainAtMillis = now;
+            reobtainScheduled = true;
+            LOG.debug(
+                    "Re-obtain of delegation tokens requested; scheduling an 
obtain cycle in {}",
+                    
TimeUtils.formatWithHighestUnit(Duration.ofMillis(delayMillis)));
+            scheduleRenewalLocked(delayMillis);
+        }
+    }
+
+    @Override
+    public void registerJob(JobID jobId, Configuration jobConfiguration) 
throws Exception {

Review Comment:
   Agreed. Added in  
https://github.com/apache/flink/pull/28639/changes/3c1dbd76bce6fbd345735c233c40b30d28824d6c



##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -416,13 +531,114 @@ long calculateRenewalDelay(Clock clock, long 
nextRenewal) {
         return renewalDelay;
     }
 
+    @VisibleForTesting
+    void setClock(Clock clock) {
+        this.clock = clock;
+    }
+
     /** Stops re-occurring token obtain task. */
     @Override
     public void stop() {
         LOG.info("Stopping credential renewal");
 
-        stopTokensUpdate();
+        synchronized (tokensUpdateFutureLock) {
+            // Mark stopped, cancel the pending cycle, and reset on-demand 
re-obtain bookkeeping
+            // atomically, so a concurrent reobtainDelegationTokens() cannot 
leave a live future
+            // orphaned after stop and a later start() does not inherit stale 
state.
+            stopped = true;
+            stopTokensUpdate();
+            reobtainScheduled = false;
+            lastReobtainAtMillis = NO_PREVIOUS_REOBTAIN;
+        }
+
+        for (DelegationTokenProvider provider : 
delegationTokenProviders.values()) {
+            try {
+                provider.stop();
+            } catch (Throwable t) {
+                LOG.error("Failed to stop delegation token provider {}", 
provider.serviceName(), t);
+            }
+        }
 
         LOG.info("Stopped credential renewal");
     }
+
+    @Override
+    public void reobtainDelegationTokens() {
+        synchronized (tokensUpdateFutureLock) {
+            if (scheduledExecutor == null || ioExecutor == null) {
+                LOG.debug(
+                        "A re-obtain of delegation tokens was requested but 
the manager was "
+                                + "constructed without executors (one-shot 
obtain path); the "
+                                + "request is ignored.");
+                return;
+            }
+            if (stopped) {
+                LOG.debug(
+                        "A re-obtain of delegation tokens was requested after 
the manager was "
+                                + "stopped; the request is ignored.");
+                return;
+            }
+            // Dedupe: if an on-demand re-obtain is already scheduled and has 
not started yet, the
+            // newly registered job(s) will be covered by it, so coalesce this 
request into it.
+            if (reobtainScheduled) {
+                LOG.debug("A re-obtain of delegation tokens is already 
scheduled; coalescing.");
+                return;
+            }
+            // Cooldown: bound how often on-demand re-obtains can run by 
deferring this cycle until
+            // at least reobtainCooldownMillis have passed since the previous 
on-demand re-obtain.
+            long now = clock.millis();
+            long delayMillis =
+                    lastReobtainAtMillis == NO_PREVIOUS_REOBTAIN
+                            ? 0L
+                            : Math.max(0L, lastReobtainAtMillis + 
reobtainCooldownMillis - now);
+            // Only bring the next cycle forward: if a cycle (e.g. the 
periodic renewal) is still
+            // pending and scheduled to fire sooner than the cooldown-deferred 
time, fire at that
+            // earlier time instead of pushing it later — otherwise a 
short-lived token could expire
+            // before it is renewed. The nextScheduledAtMillis > now guard 
skips an already-fired
+            // future that has not yet been re-armed, so this never bypasses 
the cooldown.
+            if (tokensUpdateFuture != null
+                    && nextScheduledAtMillis > now
+                    && nextScheduledAtMillis - now < delayMillis) {
+                delayMillis = nextScheduledAtMillis - now;
+            }
+            lastReobtainAtMillis = now;
+            reobtainScheduled = true;
+            LOG.debug(
+                    "Re-obtain of delegation tokens requested; scheduling an 
obtain cycle in {}",
+                    
TimeUtils.formatWithHighestUnit(Duration.ofMillis(delayMillis)));
+            scheduleRenewalLocked(delayMillis);

Review Comment:
   Yes, exactly right. Fixed in 9c1c49d: any throwable from schedule() now 
resets reobtainScheduled and nextScheduledAtMillis (same bookkeeping as the 
RejectedExecutionException branch) and is rethrown.



##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/ResourceManagerGateway.java:
##########
@@ -63,10 +64,42 @@ public interface ResourceManagerGateway
     /**
      * Register a {@link JobMaster} at the resource manager.
      *
+     * <p>Backward-compatible overload that registers without a job 
configuration. Equivalent to
+     * calling {@link #registerJobMaster(JobMasterId, ResourceID, String, 
JobID, Configuration,
+     * Duration)} with an empty configuration.
+     *
+     * @param jobMasterId The fencing token for the JobMaster leader
+     * @param jobMasterResourceId The resource ID of the JobMaster that 
registers
+     * @param jobMasterAddress The address of the JobMaster that registers
+     * @param jobId The Job ID of the JobMaster that registers
+     * @param timeout Timeout for the future to complete
+     * @return Future registration response
+     */
+    default CompletableFuture<RegistrationResponse> registerJobMaster(

Review Comment:
   Removed in c2e1b981eb6 and tests now pass the empty configuration explicitly 



##########
flink-runtime/src/test/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManagerTest.java:
##########
@@ -211,8 +231,9 @@ void startTokensUpdate() {
                     }
                 };
 
-        delegationTokenManager.startTokensUpdate();
+        // The first two cycles fail and schedule a retry each. The third 
succeeds.
         ExceptionThrowingDelegationTokenProvider.throwInUsage.set(true);
+        delegationTokenManager.start(tokens -> {});

Review Comment:
   addressed in 
https://github.com/apache/flink/pull/28639/changes/ce9efa9af8ab1307ec5e4a27eff7a6f2bc5d0910
 



##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -139,12 +240,15 @@ public DefaultDelegationTokenManager(
     private Map<String, DelegationTokenProvider> loadProviders() {
         LOG.info("Loading delegation token providers");
 
+        // Handed to every provider so it can request an immediate re-obtain 
later, from any
+        // thread, decoupled from the registerJob call stack.
+        final DelegationTokenManagerCallback callback = 
this::reobtainDelegationTokens;
         Map<String, DelegationTokenProvider> providers = new HashMap<>();
         Consumer<DelegationTokenProvider> loadProvider =
                 (provider) -> {
                     try {
                         if (isProviderEnabled(configuration, 
provider.serviceName())) {
-                            provider.init(configuration);
+                            provider.init(configuration, callback);

Review Comment:
   I can assume, you meant local variable? I removed it and inlined method ref 
in 
https://github.com/apache/flink/pull/28639/changes/8d5d41ad5c0de802898f742d94186d0d37475d4f



##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -377,13 +623,14 @@ void stopTokensUpdate() {
             if (tokensUpdateFuture != null) {
                 tokensUpdateFuture.cancel(true);
                 tokensUpdateFuture = null;
+                nextScheduledAtMillis = Long.MAX_VALUE;
             }
         }
     }
 
     @VisibleForTesting

Review Comment:
   adresed in 
https://github.com/apache/flink/pull/28639/changes/ddf100114cb5f209edd5930b5d99e9ce5313c198



##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -92,6 +101,22 @@ public class DefaultDelegationTokenManager implements 
DelegationTokenManager {
 
     @VisibleForTesting long lastKnownNextRenewal = Long.MAX_VALUE;
 
+    private final long reobtainCooldownMillis;
+
+    /**
+     * Clock used for renewal and cooldown timing. Renewal math reads absolute 
time (a token's
+     * validUntil is an absolute epoch), while scheduling and the cooldown 
read relative time, which
+     * wall-clock adjustments cannot distort. Never mix the two in one 
expression.
+     */
+    private final Clock clock;
+
+    /**
+     * Serializes the obtain-and-broadcast cycle so that, even though {@code 
cancel(true)} does not
+     * wait for an in-flight cycle and the IO executor is multi-threaded, two 
cycles can never run
+     * concurrently and broadcast tokens out of order.
+     */
+    private final Object obtainLock = new Object();

Review Comment:
   Renamed as suggested in b02b2245eac



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

Reply via email to