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]