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


##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/ResourceManager.java:
##########
@@ -427,6 +430,14 @@ public CompletableFuture<RegistrationResponse> 
registerJobMaster(
                             jobMasterIdFuture,
                             (JobMasterGateway jobMasterGateway, JobMasterId 
leadingJobMasterId) -> {
                                 if (Objects.equals(leadingJobMasterId, 
jobMasterId)) {
+                                    // Register with the delegation token 
manager first; a
+                                    // provider failure rejects this 
registration so the job
+                                    // never starts without the tokens it 
requires.
+                                    try {
+                                        
delegationTokenManager.registerJob(jobId, jobConfiguration);
+                                    } catch (Exception e) {
+                                        return new 
RegistrationResponse.Failure(e);

Review Comment:
   > oncall guys
   
   I'm something of an on-call guy myself...
   
   Yes, it **was** logged in 
https://github.com/apache/flink/pull/28639/changes#diff-d32f89982ba30269ab9438d993d9deb3b29d046792b6cc5ea35a277eaacfb287R629
  
(flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java
 registerJob() 
   ```
   LOG.error("Failed to register job {}", jobId, e);
   ```
   And moreover, we already had the failing **provider** in stack trace.  The 
rollback path has its own `ERROR log too ("Failed to roll back registration of 
job {}")`.
   
   Update to: ^ Since then the failure surfaces improved further: the logs now 
name the failing provider directly (`Failed to register job {} for provider 
{}`, 
https://github.com/apache/flink/commit/3c1dbd76bce6fbd345735c233c40b30d28824d6c),
 and the JobMaster gets the failure back wrapped as a FlinkException naming the 
job and the delegation token manager in the RegistrationResponse 
(https://github.com/apache/flink/commit/485c10a47f6c884b15ec34ecfd2999f29cead82a).



##########
flink-runtime/src/main/java/org/apache/flink/runtime/resourcemanager/ResourceManager.java:
##########
@@ -427,6 +430,14 @@ public CompletableFuture<RegistrationResponse> 
registerJobMaster(
                             jobMasterIdFuture,
                             (JobMasterGateway jobMasterGateway, JobMasterId 
leadingJobMasterId) -> {
                                 if (Objects.equals(leadingJobMasterId, 
jobMasterId)) {
+                                    // Register with the delegation token 
manager first; a
+                                    // provider failure rejects this 
registration so the job
+                                    // never starts without the tokens it 
requires.
+                                    try {
+                                        
delegationTokenManager.registerJob(jobId, jobConfiguration);
+                                    } catch (Exception e) {
+                                        return new 
RegistrationResponse.Failure(e);

Review Comment:
   One addition to my note above: that holds for an `Exception`, but an `Error` 
slips past both catches. The realistic trigger is a `NoClassDefFoundError` (or 
other `LinkageError` from provider plugin code), the same failure class 
`loadProviders` already has:
   ```
   } catch (Exception | NoClassDefFoundError e) {
                           // The intentional general rule is that if a 
provider's init method throws
                           // exception
                           // then stop the workload
                           LOG.error(
                                   "Failed to initialize delegation token 
provider {}",
                                   provider.serviceName(),
                                   e);
                           throw new FlinkRuntimeException(e);
                       }
   ```
   
   ^ Covered in 
https://github.com/apache/flink/commit/485c10a47f6c884b15ec34ecfd2999f29cead82a.
 Both job hooks now treat Exception | LinkageError the same.



##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -310,64 +361,127 @@ public void start(Listener listener) throws Exception {
         this.listener = checkNotNull(listener, "Listener must not be null");
         synchronized (tokensUpdateFutureLock) {
             checkState(tokensUpdateFuture == null, "Manager is already 
started");
+            stopped = false;

Review Comment:
   Agreed on the rename and the early returns, implemented in 9c1c49d. However, 
I used the name `running` because it is used more often in the codebase.
   
   And another part I did differently: I set `running = true` before the inline 
first cycle in `start()`, exactly where `stopped = false` sat, not at the end 
of the method (if I understand your suggestion correctly). The first obtain 
runs inline and checks the flag twice, at the top of startTokensUpdate() and 
again in maybeScheduleRenewal() when it schedules the periodic renewal. With 
the flag flipped only at the end of start(), the inline cycle sees not-running, 
skips the whole first obtain, and never schedules the renewal. 
startShouldBeIdempotent and startTokensUpdateShouldScheduleRenewal lock that 
ordering in, both fail if the flag moves after the cycle.
   
   Same reason inverted on the stop() side: `running = false` flips before the 
per-job unregistration below it, so a re-obtain racing shutdown cannot schedule 
a cycle for a session that is shutting down. That cleanup also runs 
unconditionally (no early return) because ResourceManager calls stop() on the 
failed-startup path even when start() never ran.



##########
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()) {

Review Comment:
   > Start called from another thread and start some/all of them
   
   Providers have no start/restart hooks: init() runs once in the manager 
constructor, and start() only drives the obtain logic, so there is nothing 
start() could call to re-start a provider mid stop. An obtain cycle overlapping 
stop() is possible, but that overlap is already covered by the stop() javadoc.
   
   But another issue is possible. The manager instance is reused across 
leadership sessions, so it is possible to have stop() then start() on the same 
manager with the same provider instances and no re-init. A provider that closes 
its resources in stop() (as the javadoc tells it to) is broken in the next 
leadership term. The SPI javadoc says stop() is "called once during manager 
shutdown" and the wiring does not honor that.
   
   I see 2 options (taking into account the existing architecture/design):
   
   1.Fan out provider.stop() only at process shutdown and let manager stop() 
just stop scheduling. Providers keep a single init-to-stop lifecycle, matching 
their single init.
   Or
   2.Keep per-session stop() but redefine the contract: stop() may be followed 
by another start(), so providers must release only per-run resources and 
re-acquire lazy. Smaller diff, but it changes the agreed "called once" approach 
and pushes restart-stuff into every provider implementation.
   
   I implemented option 1 in 
[35495a6](https://github.com/apache/flink/pull/28639/changes/35495a68533d5da7f78ea83e2dec2ba077dd8c66).
 `DelegationTokenManager.close()` is the terminal teardown: it ends the session 
via stop(), then stops the providers exactly once, and a closed manager rejects 
a later start(). It is called by the component that created the manager: 
ClusterEntrypoint, MiniCluster, and YarnClusterDescriptor. The per-session 
stop() keeps the providers usable for the next leadership session. It is a 
single commit, so easy to rework if you prefer option 2.
   
   WDYT?



##########
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()) {

Review Comment:
   1. Renamed stop() to close() in [[FLINK-40019][core][runtime] Rename 
DelegationTokenProvider.stop() to 
close()](https://github.com/apache/flink/pull/28639/changes/4df526e2daaa3badbf2222fbd22343e02e4a151b)
 . Agree it meets naming API better.
   2.
   > correct me if I'm wrong
   > stop (together with start) can happen when HA kicks in during failover and 
it makes sure that new leader re-obtained tokens.
   > stop has nothing to do with freeing resources allocated by the provider 
and this function called regularly.
   
   yes, your understanding is correct 👍 
   
   



##########
flink-runtime/src/main/java/org/apache/flink/runtime/security/token/DefaultDelegationTokenManager.java:
##########
@@ -310,64 +361,127 @@ public void start(Listener listener) throws Exception {
         this.listener = checkNotNull(listener, "Listener must not be null");
         synchronized (tokensUpdateFutureLock) {
             checkState(tokensUpdateFuture == null, "Manager is already 
started");
+            stopped = false;
         }
 
         startTokensUpdate();
     }
 
     @VisibleForTesting
     void startTokensUpdate() {
-        try {
-            LOG.info("Starting tokens update task");
-            DelegationTokenContainer container = new 
DelegationTokenContainer();
-            Optional<Long> nextRenewal = 
obtainDelegationTokensAndGetNextRenewal(container);
-
-            if (container.hasTokens()) {
-                
delegationTokenReceiverRepository.onNewTokensObtained(container);
-
-                LOG.info("Notifying listener about new tokens");
-                checkNotNull(listener, "Listener must not be null");
-                
listener.onNewTokensObtained(InstantiationUtil.serializeObject(container));
-                LOG.info("Listener notified successfully");
-            } else {
-                LOG.warn("No tokens obtained so skipping notifications");
+        synchronized (tokensUpdateFutureLock) {
+            // The obtain cycle is starting: clear the dedupe flag so later 
on-demand requests can
+            // schedule a fresh cycle.
+            reobtainScheduled = false;
+            // If stop() ran before this cycle (already handed to the IO 
executor) began, skip the
+            // obtain/broadcast: the providers may already be stopped. Safe 
via this lock's
+            // happens-before with stop(). The dedupe flag is cleared above, 
so it is never stuck.
+            if (stopped) {
+                return;
             }
+        }
+        // Serialize the obtain-and-broadcast so a re-obtain racing the 
periodic renewal cannot run
+        // two cycles concurrently on the (multi-threaded) IO executor and 
broadcast out of order.
+        synchronized (obtainLock) {
+            try {
+                LOG.info("Starting tokens update task");
+                DelegationTokenContainer container = new 
DelegationTokenContainer();
+                Optional<Long> nextRenewal = 
obtainDelegationTokensAndGetNextRenewal(container);
+
+                if (container.hasTokens()) {
+                    
delegationTokenReceiverRepository.onNewTokensObtained(container);
+
+                    LOG.info("Notifying listener about new tokens");
+                    checkNotNull(listener, "Listener must not be null");
+                    
listener.onNewTokensObtained(InstantiationUtil.serializeObject(container));
+                    LOG.info("Listener notified successfully");
+                } else {
+                    LOG.warn("No tokens obtained so skipping notifications");
+                }
 
-            if (nextRenewal.isPresent()) {
-                lastKnownNextRenewal = nextRenewal.get();
-                currentRetryBackoff = renewalRetryInitialBackoff;
-                long renewalDelay =
-                        calculateRenewalDelay(Clock.systemDefaultZone(), 
nextRenewal.get());
-                synchronized (tokensUpdateFutureLock) {
-                    tokensUpdateFuture =
-                            scheduledExecutor.schedule(
-                                    () -> 
ioExecutor.execute(this::startTokensUpdate),
-                                    renewalDelay,
-                                    TimeUnit.MILLISECONDS);
+                if (nextRenewal.isPresent()) {
+                    lastKnownNextRenewal = nextRenewal.get();
+                    currentRetryBackoff = renewalRetryInitialBackoff;
+                    long renewalDelay = calculateRenewalDelay(clock, 
nextRenewal.get());
+                    maybeScheduleRenewal(renewalDelay);
+                    LOG.info(
+                            "Tokens update task started with {} delay",
+                            
TimeUtils.formatWithHighestUnit(Duration.ofMillis(renewalDelay)));
+                } else {
+                    LOG.warn(
+                            "Tokens update task not started because either no 
tokens obtained or none of the tokens specified its renewal date");
                 }
-                LOG.info(
-                        "Tokens update task started with {} delay",
-                        
TimeUtils.formatWithHighestUnit(Duration.ofMillis(renewalDelay)));
-            } else {
+            } catch (InterruptedException e) {
+                // Ignore, may happen if shutting down.
+                LOG.debug("Interrupted", e);
+            } catch (Exception e) {
+                long delay = calculateRetryDelay(clock);
+                maybeScheduleRenewal(delay);
                 LOG.warn(
-                        "Tokens update task not started because either no 
tokens obtained or none of the tokens specified its renewal date");
+                        "Failed to update tokens, will try again in {}",
+                        
TimeUtils.formatWithHighestUnit(Duration.ofMillis(delay)),
+                        e);

Review Comment:
   You are right, the log could claim a delay that was never scheduled. Fixed 
in 9c1c49d: maybeScheduleRenewal now returns the delay that actually took 
effect and the log prints that.



##########
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 {
+        try {
+            for (DelegationTokenProvider provider : 
delegationTokenProviders.values()) {
+                provider.registerJob(jobId, jobConfiguration);
+            }
+        } catch (Exception e) {
+            // If any of the providers fail to register, then unregister the 
job from them all.
+            // unregisterJob is idempotent, so it is safe to call it for 
providers that were never
+            // (or only partially) registered for this job before the failure. 
The rollback must
+            // never mask the original failure, so swallow any rollback 
exception.
+            try {
+                unregisterJob(jobId);
+            } catch (Exception rollbackException) {
+                LOG.error("Failed to roll back registration of job {}", jobId, 
rollbackException);
+            }
+            LOG.error("Failed to register job {}", jobId, e);
+            throw e;
+        }
+    }
+
+    @Override
+    public void unregisterJob(JobID jobId) throws Exception {
+        for (DelegationTokenProvider provider : 
delegationTokenProviders.values()) {
+            try {
+                provider.unregisterJob(jobId);
+            } catch (Exception e) {
+                LOG.error("Failed to unregister job for provider {}", 
provider.serviceName(), e);

Review Comment:
   Good catch, added in 
https://github.com/apache/flink/pull/28639/changes/3c1dbd76bce6fbd345735c233c40b30d28824d6c



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