Skip to content

Fix supervisor crash on transient ZK connection loss - #9153

Open
jkrauss82 wants to merge 2 commits into
apache:masterfrom
jkrauss82:fix/connection-aware-zk-retry
Open

jkrauss82 wants to merge 2 commits into
apache:masterfrom
jkrauss82:fix/connection-aware-zk-retry

Conversation

@jkrauss82

@jkrauss82 jkrauss82 commented Oct 1, 2026 •

Copy link
Copy Markdown

Root cause: A transient ZooKeeper connection loss during the supervisor heartbeat cycle caused the process to die in 629ms via DefaultUncaughtExceptionHandler. The Curator RetryLoop blindly slept between retries, racing against the ZK client's SendThread reconnection instead of yielding to it.

Fix (two layers):

  1. ConnectionAwareRetryPolicy (new): A RetryPolicy wrapper that tracks the ZK ConnectionState via a ConnectionStateListener. When the connection is SUSPENDED or LOST, it calls CuratorFramework.blockUntilConnected(sessionTimeout) to yield to the SendThread's failover, instead of blind sleep+retry. On reconnection, the operation retries immediately on the new connection. Uses storm.zookeeper.session.timeout as the wait bound (no magic numbers).

  2. SupervisorHeartbeat safety net: Wrap the heartbeat body in try/catch so that a total ZK outage (beyond session timeout) skips the cycle instead of killing the process. A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (30s) allows 6 missed beats.

Wired up in CuratorUtils.newCurator() via AtomicReference late binding (the CuratorFramework doesn't exist until builder.build() returns).

What is the purpose of the change:

This fix prevents supervisor processes from dying due to transient ZooKeeper connection losses. During the incident, a single ZK ensemble member closing its socket caused the supervisor to exit in 629ms via DefaultUncaughtExceptionHandler, killing all topology workers on that node. The fix absorbs these transient blips by making the Curator retry policy aware of the ZK connection state — when SUSPENDED or LOST, it yields to the ZK client's SendThread reconnection via blockUntilConnected() instead of blind sleep+retry. This ensures supervisors survive network/ZK blips that are already being handled by the ZK client's failover mechanism.

How was the change tested:

The fix was deployed to a production cluster running 3 supervisors with 220+ active workers. During a subsequent network outage that caused multiple transient ZK connection losses, the supervisors survived without dying (previously they would have exited via DefaultUncaughtExceptionHandler). Log verification confirmed the ConnectionAwareRetryPolicy correctly detected SUSPENDED states, waited for reconnection via blockUntilConnected(), and retried operations on the new connection. No supervisors died, no topology workers were killed due to supervisor crashes.

Closes #9098


EDIT 2026-10-04: attach crash log

connectionloss-crash-anonymized.log

Root cause: A transient ZooKeeper connection loss during the supervisor
heartbeat cycle caused the process to die in 629ms via
DefaultUncaughtExceptionHandler. The Curator RetryLoop blindly slept
between retries, racing against the ZK client's SendThread reconnection
instead of yielding to it.

Fix (two layers):

1. ConnectionAwareRetryPolicy (new): A RetryPolicy wrapper that tracks
   the ZK ConnectionState via a ConnectionStateListener. When the
   connection is SUSPENDED or LOST, it calls
   CuratorFramework.blockUntilConnected(sessionTimeout) to yield to the
   SendThread's failover, instead of blind sleep+retry. On reconnection,
   the operation retries immediately on the new connection. Uses
   storm.zookeeper.session.timeout as the wait bound (no magic numbers).

2. SupervisorHeartbeat safety net: Wrap the heartbeat body in
   try/catch so that a total ZK outage (beyond session timeout) skips
   the cycle instead of killing the process. A missed heartbeat is
   harmless: nimbus.supervisor.timeout.secs (30s) allows 6 missed beats.

Wired up in CuratorUtils.newCurator() via AtomicReference late binding
(the CuratorFramework doesn't exist until builder.build() returns).
@reiabreu
reiabreu requested a review from rzo1 October 2, 2026 15:06
@reiabreu reiabreu added this to the 3.2.0 milestone Oct 2, 2026

@reiabreu reiabreu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This review was generated with the help of an LLM (Claude) and reviewed by me before posting.

Thanks for this — the #9098 diagnosis is solid and the two-layer approach (connection-aware retry + heartbeat safety net) targets the real root cause. I pulled the branch and ran the storm-client/storm-server suites. Findings below, with code where it helps. One blocker.

1. This breaks an existing test (blocker)

CuratorUtilsTest.newCuratorUsesExponentialBackoffTest casts the framework's retry policy directly to StormBoundedExponentialBackoffRetry. Since newCurator() now returns a ConnectionAwareRetryPolicy wrapper, that cast throws ClassCastException and mvn -pl storm-client test fails:

java.lang.ClassCastException: class ConnectionAwareRetryPolicy cannot be cast to class StormBoundedExponentialBackoffRetry  at CuratorUtilsTest.java:75

The fix keeps the original coverage (that newCurator wires the right interval/retries/ceiling) by unwrapping the delegate. Two small parts:

(a) add a test-only accessor to ConnectionAwareRetryPolicy (plus the import org.apache.storm.shade.com.google.common.annotations.VisibleForTesting;):

@VisibleForTesting
RetryPolicy getDelegate() {
    return delegate;
}

(b) unwrap in CuratorUtilsTest.newCuratorUsesExponentialBackoffTest (add import org.apache.storm.shade.org.apache.curator.RetryPolicy; and assertTrue):

RetryPolicy retryPolicy = curator.getZookeeperClient().getRetryPolicy();
assertTrue(retryPolicy instanceof ConnectionAwareRetryPolicy,
    "newCurator should wrap the retry policy in a ConnectionAwareRetryPolicy");
StormBoundedExponentialBackoffRetry policy =
    (StormBoundedExponentialBackoffRetry) ((ConnectionAwareRetryPolicy) retryPolicy).getDelegate();
// existing getBaseSleepTimeMs / getN / getSleepTimeMs assertions stay as-is

2. No unit tests for the new logic

ConnectionAwareRetryPolicy is 125 lines of connection-state/retry logic with no coverage (the PR's validation is a production anecdote). I wrote two classes — both pass locally (7/7 and 2/2). The policy tests use Mockito by capturing the ConnectionStateListener from bind(), so no real ZooKeeper is needed.

ConnectionAwareRetryPolicyTest.java (storm-client, 7 cases)
/*
 * Licensed to the Apache Software Foundation (ASF) under one
 * or more contributor license agreements.  See the NOTICE file
 * distributed with this work for additional information
 * regarding copyright ownership.  The ASF licenses this file
 * to you under the Apache License, Version 2.0 (the
 * "License"); you may not use this file except in compliance
 * with the License.  You may obtain a copy of the License at
 *
 * http://www.apache.org/licenses/LICENSE-2.0
 *
 * Unless required by applicable law or agreed to in writing, software
 * distributed under the License is distributed on an "AS IS" BASIS,
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
 * See the License for the specific language governing permissions and
 * limitations under the License.
 */

package org.apache.storm.utils;

import java.util.concurrent.TimeUnit;
import java.util.function.Supplier;
import org.apache.storm.shade.org.apache.curator.RetryPolicy;
import org.apache.storm.shade.org.apache.curator.RetrySleeper;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
import org.apache.storm.shade.org.apache.curator.framework.listen.Listenable;
import org.apache.storm.shade.org.apache.curator.framework.state.ConnectionState;
import org.apache.storm.shade.org.apache.curator.framework.state.ConnectionStateListener;
import org.junit.jupiter.api.Test;
import org.mockito.ArgumentCaptor;

import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertTrue;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
import static org.mockito.ArgumentMatchers.anyLong;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.never;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.verifyNoInteractions;
import static org.mockito.Mockito.when;

/**
 * Unit tests for {@link ConnectionAwareRetryPolicy}. The ZooKeeper connection state is driven by
 * capturing the {@link ConnectionStateListener} that {@link ConnectionAwareRetryPolicy#bind} registers
 * and firing state transitions at it, so no real ZooKeeper is required.
 */
public class ConnectionAwareRetryPolicyTest {

    private static final int SESSION_TIMEOUT_MS = 20_000;

    // JUnit 5 creates a fresh test instance per method, so these are clean for each test.
    private final RetryPolicy delegate = mock(RetryPolicy.class);
    private final RetrySleeper sleeper = mock(RetrySleeper.class);
    private final CuratorFramework zk = mock(CuratorFramework.class);
    private ConnectionStateListener listener;

    /**
     * Builds the policy, binds it to the mock framework, and captures the registered listener.
     *
     * @param zkSupplier what {@code zkSupplier.get()} returns inside allowRetry (the real wiring passes
     *                   {@code zkRef::get}; pass {@code () -> null} to exercise the not-yet-bound path)
     */
    @SuppressWarnings("unchecked")
    private ConnectionAwareRetryPolicy build(Supplier<CuratorFramework> zkSupplier) {
        Listenable<ConnectionStateListener> listenable = mock(Listenable.class);
        when(zk.getConnectionStateListenable()).thenReturn(listenable);

        ConnectionAwareRetryPolicy policy = new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS);
        policy.bind(zk);

        ArgumentCaptor<ConnectionStateListener> captor = ArgumentCaptor.forClass(ConnectionStateListener.class);
        verify(listenable).addListener(captor.capture());
        listener = captor.getValue();
        return policy;
    }

    private void fire(ConnectionState state) {
        listener.stateChanged(zk, state);
    }

    @Test
    public void connectedStateDelegatesToBackoffPolicy() {
        ConnectionAwareRetryPolicy policy = build(() -> zk);
        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);

        // Default state is CONNECTED (no event fired).
        boolean result = policy.allowRetry(0, 0L, sleeper);

        assertTrue(result, "when connected, must propagate the delegate's decision");
        verify(delegate).allowRetry(0, 0L, sleeper);
    }

    @Test
    public void reconnectedStateDelegatesToBackoffPolicy() {
        ConnectionAwareRetryPolicy policy = build(() -> zk);
        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);

        fire(ConnectionState.RECONNECTED);
        boolean result = policy.allowRetry(2, 100L, sleeper);

        assertFalse(result, "RECONNECTED is a healthy state and must delegate");
        verify(delegate).allowRetry(2, 100L, sleeper);
    }

    @Test
    public void suspendedBlocksUntilConnectedThenRetriesImmediately() throws Exception {
        ConnectionAwareRetryPolicy policy = build(() -> zk);
        when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);

        fire(ConnectionState.SUSPENDED);
        boolean result = policy.allowRetry(0, 0L, sleeper);

        assertTrue(result, "should retry once the SendThread has reconnected");
        verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
        verifyNoInteractions(sleeper); // retries immediately, no backoff sleep
        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
    }

    @Test
    public void lostBlocksUntilConnectedThenRetriesImmediately() throws Exception {
        ConnectionAwareRetryPolicy policy = build(() -> zk);
        when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);

        fire(ConnectionState.LOST);
        boolean result = policy.allowRetry(3, 1234L, sleeper);

        assertTrue(result, "LOST should also wait for reconnection then retry");
        verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
    }

    @Test
    public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
        ConnectionAwareRetryPolicy policy = build(() -> zk);
        when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(false);

        fire(ConnectionState.SUSPENDED);
        boolean result = policy.allowRetry(0, 0L, sleeper);

        assertFalse(result, "should abandon when not reconnected within the session timeout");
        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
    }

    @Test
    public void interruptedWhileWaitingReturnsFalseAndPreservesInterruptFlag() throws Exception {
        ConnectionAwareRetryPolicy policy = build(() -> zk);
        when(zk.blockUntilConnected(anyInt(), any())).thenThrow(new InterruptedException("test"));

        fire(ConnectionState.SUSPENDED);
        boolean result = policy.allowRetry(0, 0L, sleeper);

        assertFalse(result, "an interrupt while waiting should abandon the retry");
        assertTrue(Thread.interrupted(), "interrupt flag must be preserved (this check also clears it)");
    }

    @Test
    public void suspendedWithNullFrameworkFallsThroughToDelegate() {
        // zkSupplier returns null (framework not yet available), even though bind() ran on the mock.
        ConnectionAwareRetryPolicy policy = build(() -> null);
        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);

        fire(ConnectionState.SUSPENDED);
        boolean result = policy.allowRetry(1, 50L, sleeper);

        assertTrue(result, "with no framework available yet, must fall through to the delegate");
        verify(delegate).allowRetry(1, 50L, sleeper);
    }
}
SupervisorHeartbeatTest.java (storm-server, 2 cases)
/*
 * Licensed to the Apache Software Foundation (ASF) under one or more contributor license agreements.  See the NOTICE file distributed with
 * this work for additional information regarding copyright ownership.  The ASF licenses this file to you under the Apache License, Version
 * 2.0 (the "License"); you may not use this file except in compliance with the License.  You may obtain a copy of the License at
 *
 * http://www.apache.org/licenses/LICENSE-2.0
 *
 * Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS,
 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions
 * and limitations under the License.
 */

package org.apache.storm.daemon.supervisor.timer;

import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.storm.cluster.IStormClusterState;
import org.apache.storm.daemon.supervisor.Supervisor;
import org.apache.storm.generated.SupervisorInfo;
import org.apache.storm.scheduler.ISupervisor;
import org.apache.storm.utils.Utils;
import org.junit.jupiter.api.Test;

import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyString;
import static org.mockito.ArgumentMatchers.eq;
import static org.mockito.Mockito.atLeastOnce;
import static org.mockito.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;

public class SupervisorHeartbeatTest {

    private static final String SUPERVISOR_ID = "supervisor-1";

    /**
     * Mocks a Supervisor with just enough state for {@code buildSupervisorInfo} to produce one
     * SupervisorInfo entry (the non-NUMA branch), so {@code run()} reaches the ZK heartbeat write.
     */
    private Supervisor mockSupervisor(IStormClusterState clusterState) {
        Supervisor supervisor = mock(Supervisor.class);
        ISupervisor iSupervisor = mock(ISupervisor.class);
        when(iSupervisor.getMetadata()).thenReturn(Arrays.asList(6700, 6701));

        when(supervisor.getStormClusterState()).thenReturn(clusterState);
        when(supervisor.getId()).thenReturn(SUPERVISOR_ID);
        when(supervisor.getiSupervisor()).thenReturn(iSupervisor);
        when(supervisor.getCurrAssignment()).thenReturn(new AtomicReference<>(new HashMap<>()));
        when(supervisor.getHostName()).thenReturn("host-1");
        when(supervisor.getAssignmentId()).thenReturn(SUPERVISOR_ID);
        when(supervisor.getThriftServerPort()).thenReturn(6627);
        when(supervisor.getUpTime()).thenReturn(Utils.makeUptimeComputer());
        when(supervisor.getStormVersion()).thenReturn("test-version");
        return supervisor;
    }

    private Map<String, Object> baseConf() {
        return new HashMap<>(Utils.readStormConfig());
    }

    @Test
    public void transientHeartbeatFailureIsSwallowedInsteadOfKillingTheProcess() {
        IStormClusterState clusterState = mock(IStormClusterState.class);
        // Simulate a transient ZK write failure during the heartbeat cycle (the real one is a wrapped
        // KeeperException$ConnectionLoss; the catch is on Exception, so a RuntimeException exercises the same path).
        doThrow(new RuntimeException("simulated transient ZK failure"))
            .when(clusterState).supervisorHeartbeat(anyString(), any(SupervisorInfo.class));

        Supervisor supervisor = mockSupervisor(clusterState);
        SupervisorHeartbeat heartbeat = new SupervisorHeartbeat(baseConf(), supervisor);

        // The whole point of the fix: run() must not propagate — otherwise the StormTimer's
        // DefaultUncaughtExceptionHandler would exit the supervisor process.
        assertDoesNotThrow(heartbeat::run);
        // It must have actually attempted the write (i.e. reached the failing ZK call).
        verify(clusterState, atLeastOnce()).supervisorHeartbeat(anyString(), any(SupervisorInfo.class));
    }

    @Test
    public void successfulCycleSendsTheHeartbeat() {
        IStormClusterState clusterState = mock(IStormClusterState.class);
        Supervisor supervisor = mockSupervisor(clusterState);
        SupervisorHeartbeat heartbeat = new SupervisorHeartbeat(baseConf(), supervisor);

        assertDoesNotThrow(heartbeat::run);
        verify(clusterState).supervisorHeartbeat(eq(SUPERVISOR_ID), any(SupervisorInfo.class));
    }
}

3. Blast radius — please confirm intent

newCurator() backs every Curator client (Nimbus, blobstore, Trident TransactionalState, ClientZookeeper), so this changes ZK retry behavior cluster-wide, not just the supervisor heartbeat. On SUSPENDED/LOST, operations that previously failed fast can now block up to storm.zookeeper.session.timeout (20s default). Is that acceptable for the Nimbus/leader paths, which weren't exercised here?

4. No overall bound on the SUSPENDED/LOST path (optional)

That branch ignores retryCount/elapsedTimeMs and the delegate's retry.times cap, so a flapping connection can retry indefinitely and wedge the calling thread (e.g. the heartbeat timer). Consider still consulting the delegate's budget before waiting, e.g.:

// honour the configured retry budget even while suspended
if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
    return false;
}

5. Thundering herd on reconnect (optional)

The deliberate immediate retry (no jitter) means many suspended clients retry in lockstep the instant the ensemble member recovers; the old backoff spread them out. A little jitter before return true; would help on large clusters.

6 & 7. The heartbeat catch is broader than the bug (see inline suggestion)

It wraps getNumaMap()/buildSupervisorInfo() (local, non-ZK logic) too, so a persistent local misconfig (e.g. a malformed supervisor.numa.meta) would be swallowed every cycle forever instead of surfacing. The inline suggestion scopes the catch to just the ZK write — and also corrects the comment: nimbus.supervisor.timeout.secs defaults to 60s, not 30s (~12 tolerated misses at a 5s frequency, not 6).

Comment on lines 162 to 174
try {
Map<String, Object> validatedNumaMap = SupervisorUtils.getNumaMap(conf);
Map<String, SupervisorInfo> supervisorInfoList = buildSupervisorInfo(conf, supervisor, validatedNumaMap);
for (Map.Entry<String, SupervisorInfo> supervisorInfoEntry : supervisorInfoList.entrySet()) {
stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue());
}
} catch (Exception e) {
// A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (default 30s)
// allows 6 missed beats before the supervisor is considered dead. The next
// heartbeat (5s later) will retry. Killing the entire supervisor process
// over one missed beat is not acceptable.
LOG.warn("Supervisor heartbeat failed, will retry next cycle", e);
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Scope the catch to the ZK write so a transient blip is skipped, but a persistent local failure (NUMA map, resource capacities) still surfaces instead of being swallowed forever. Also fixes the default in the comment (60s, not 30s).

Suggested change
try {
Map<String, Object> validatedNumaMap = SupervisorUtils.getNumaMap(conf);
Map<String, SupervisorInfo> supervisorInfoList = buildSupervisorInfo(conf, supervisor, validatedNumaMap);
for (Map.Entry<String, SupervisorInfo> supervisorInfoEntry : supervisorInfoList.entrySet()) {
stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue());
}
} catch (Exception e) {
// A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (default 30s)
// allows 6 missed beats before the supervisor is considered dead. The next
// heartbeat (5s later) will retry. Killing the entire supervisor process
// over one missed beat is not acceptable.
LOG.warn("Supervisor heartbeat failed, will retry next cycle", e);
}
Map<String, Object> validatedNumaMap = SupervisorUtils.getNumaMap(conf);
Map<String, SupervisorInfo> supervisorInfoList = buildSupervisorInfo(conf, supervisor, validatedNumaMap);
for (Map.Entry<String, SupervisorInfo> supervisorInfoEntry : supervisorInfoList.entrySet()) {
try {
stormClusterState.supervisorHeartbeat(supervisorInfoEntry.getKey(), supervisorInfoEntry.getValue());
} catch (Exception e) {
// A missed heartbeat is harmless: nimbus.supervisor.timeout.secs (default 60s) tolerates
// several missed beats; the next beat (5s later) retries. Local/config errors are left
// outside this catch so a persistent misconfiguration still surfaces instead of being
// swallowed every cycle.
LOG.warn("Supervisor heartbeat for {} failed, will retry next cycle",
supervisorInfoEntry.getKey(), e);
}
}

@reiabreu
reiabreu requested a review from GGraziadei October 2, 2026 15:58
@reiabreu

reiabreu commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

@rzo1 @GGraziadei the point highlighted above is more of a design decision. I'm inclined to say that yes, we want all processes to be resilient to zookeeper connection blips. What is your take on this?

  1. Blast radius — please confirm intent
    newCurator() backs every Curator client (Nimbus, blobstore, Trident TransactionalState, ClientZookeeper), so this changes ZK retry behavior cluster-wide, not just the supervisor heartbeat. On SUSPENDED/LOST, operations that previously failed fast can now block up to storm.zookeeper.session.timeout (20s default). Is that acceptable for the Nimbus/leader paths, which weren't exercised here?

@rzo1

rzo1 commented Oct 2, 2026

Copy link
Copy Markdown
Contributor

Please check the build failures. Thanks :)

…catch

- Fix CuratorUtilsTest ClassCastException (unwrap ConnectionAwareRetryPolicy)
- Add @VisibleForTesting getDelegate() accessor
- Add ConnectionAwareRetryPolicyTest (8 cases) and SupervisorHeartbeatTest (2 cases)
- Add delegate budget check before blocking on SUSPENDED/LOST
- Add RECONNECT_JITTER_MS constant (100ms) to avoid thundering herd
- Scope SupervisorHeartbeat catch to ZK write only (not getNumaMap/buildSupervisorInfo)
- Fix heartbeat comment (liveness via ephemeral ZK nodes, not heartbeat timestamps)
@jkrauss82

Copy link
Copy Markdown
Author

Thanks for the review @reiabreu I have addressed the comments/suggestions in the new commit I just pushed. Build failures should be gone now, too @rzo1

@reiabreu reiabreu left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This review was generated with the help of an LLM (Claude) and reviewed by me before posting.

Follow-up on the retry-bound change (my earlier point 4). The current form opens the SUSPENDED/LOST branch with if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) return false; to honour the budget — but delegate.allowRetry (Curator's SleepingRetry) sleeps the exponential backoff as a side effect before returning true. So the suspended path now does backoff sleep → blockUntilConnected() → jitter, which reintroduces exactly the blind backoff this policy was built to avoid (the old "retry immediately" comment was dropped). The unit tests don't catch it because delegate is mocked, so the real sleep never runs.

Proposed fix: bound the suspended path by retry count (a pure check, no sleep) instead of calling the sleeping delegate. maxRetries is wired from the same storm.zookeeper.retry.times the delegate already uses, so both paths give up after the same number of attempts; the bounded jitter stays. The healthy (CONNECTED/RECONNECTED) path is unchanged and still uses the real backoff, which is correct there.

One API note: this adds a 4-arg constructor to ConnectionAwareRetryPolicy (your class) — (..., int sessionTimeoutMs, int maxRetries).

I ran it locally: ConnectionAwareRetryPolicyTest 8/8 and CuratorUtilsTest 5/5 pass (the only build failure is the unrelated python3-test CLI harness, which fails the same way without this change). The tests now assert verify(delegate, never()).allowRetry(...) on the suspended path — i.e. no blind backoff — plus a bounded-jitter check, and drive the budget by retry count.

The core production hunk is inline below as a suggestion; it depends on the constructor + CuratorUtils edits, so apply it together with the full patch:

diff --git a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
index eb041bdef..44d5a20fa 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
@@ -60,6 +60,7 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
     private final RetryPolicy delegate;
     private final Supplier<CuratorFramework> zkSupplier;
     private final int sessionTimeoutMs;
+    private final int maxRetries;
     private final AtomicReference<ConnectionState> connectionState =
         new AtomicReference<>(ConnectionState.CONNECTED);
 
@@ -69,13 +70,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
      *                         the framework does not exist until {@code builder.build()} returns
      * @param sessionTimeoutMs upper bound (in ms) for waiting on reconnection; typically
      *                         {@code storm.zookeeper.session.timeout}
+     * @param maxRetries       cap on the number of retries on the SUSPENDED/LOST path, applied without a
+     *                         backoff sleep. Should match the delegate's retry budget
+     *                         ({@code storm.zookeeper.retry.times}) so the two paths give up after the
+     *                         same number of attempts.
      */
     public ConnectionAwareRetryPolicy(RetryPolicy delegate,
                                       Supplier<CuratorFramework> zkSupplier,
-                                      int sessionTimeoutMs) {
+                                      int sessionTimeoutMs,
+                                      int maxRetries) {
         this.delegate = delegate;
         this.zkSupplier = zkSupplier;
         this.sessionTimeoutMs = sessionTimeoutMs;
+        this.maxRetries = maxRetries;
     }
 
     /**
@@ -104,14 +111,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
         ConnectionState state = connectionState.get();
 
         if (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST) {
-            // Honour the configured retry budget even while suspended.
-            if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
+            // Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On
+            // SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected()
+            // instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff
+            // as a side effect and reintroduce exactly the blind backoff this policy exists to avoid.
+            if (retryCount >= maxRetries) {
+                LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry",
+                    state, maxRetries, retryCount);
                 return false;
             }
 
             CuratorFramework zk = zkSupplier.get();
             if (zk == null) {
-                // Framework not yet available; delegate already approved the retry
+                // Framework not yet available; within budget, so allow the retry.
                 return true;
             }
 
diff --git a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
index d2365bd1a..d15e2b143 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
@@ -63,13 +63,15 @@ public class CuratorUtils {
         // Wrap with a connection-aware policy that yields to SendThread on SUSPENDED/LOST.
         AtomicReference<CuratorFramework> zkRef = new AtomicReference<>();
         int sessionTimeoutMs = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT));
+        int retryTimes = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES));
         ConnectionAwareRetryPolicy connectionAwarePolicy = new ConnectionAwareRetryPolicy(
                 new StormBoundedExponentialBackoffRetry(
                         ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL)),
                         ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL_CEILING)),
-                        ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES))),
+                        retryTimes),
                 zkRef::get,
-                sessionTimeoutMs);
+                sessionTimeoutMs,
+                retryTimes);
         builder.retryPolicy(connectionAwarePolicy);
 
         if (defaultAcl != null) {
Full patch (production + updated tests)
diff --git a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
index eb041bdef..44d5a20fa 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
@@ -60,6 +60,7 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
     private final RetryPolicy delegate;
     private final Supplier<CuratorFramework> zkSupplier;
     private final int sessionTimeoutMs;
+    private final int maxRetries;
     private final AtomicReference<ConnectionState> connectionState =
         new AtomicReference<>(ConnectionState.CONNECTED);
 
@@ -69,13 +70,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
      *                         the framework does not exist until {@code builder.build()} returns
      * @param sessionTimeoutMs upper bound (in ms) for waiting on reconnection; typically
      *                         {@code storm.zookeeper.session.timeout}
+     * @param maxRetries       cap on the number of retries on the SUSPENDED/LOST path, applied without a
+     *                         backoff sleep. Should match the delegate's retry budget
+     *                         ({@code storm.zookeeper.retry.times}) so the two paths give up after the
+     *                         same number of attempts.
      */
     public ConnectionAwareRetryPolicy(RetryPolicy delegate,
                                       Supplier<CuratorFramework> zkSupplier,
-                                      int sessionTimeoutMs) {
+                                      int sessionTimeoutMs,
+                                      int maxRetries) {
         this.delegate = delegate;
         this.zkSupplier = zkSupplier;
         this.sessionTimeoutMs = sessionTimeoutMs;
+        this.maxRetries = maxRetries;
     }
 
     /**
@@ -104,14 +111,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
         ConnectionState state = connectionState.get();
 
         if (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST) {
-            // Honour the configured retry budget even while suspended.
-            if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
+            // Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On
+            // SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected()
+            // instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff
+            // as a side effect and reintroduce exactly the blind backoff this policy exists to avoid.
+            if (retryCount >= maxRetries) {
+                LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry",
+                    state, maxRetries, retryCount);
                 return false;
             }
 
             CuratorFramework zk = zkSupplier.get();
             if (zk == null) {
-                // Framework not yet available; delegate already approved the retry
+                // Framework not yet available; within budget, so allow the retry.
                 return true;
             }
 
diff --git a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
index d2365bd1a..d15e2b143 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
@@ -63,13 +63,15 @@ public class CuratorUtils {
         // Wrap with a connection-aware policy that yields to SendThread on SUSPENDED/LOST.
         AtomicReference<CuratorFramework> zkRef = new AtomicReference<>();
         int sessionTimeoutMs = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT));
+        int retryTimes = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES));
         ConnectionAwareRetryPolicy connectionAwarePolicy = new ConnectionAwareRetryPolicy(
                 new StormBoundedExponentialBackoffRetry(
                         ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL)),
                         ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL_CEILING)),
-                        ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES))),
+                        retryTimes),
                 zkRef::get,
-                sessionTimeoutMs);
+                sessionTimeoutMs,
+                retryTimes);
         builder.retryPolicy(connectionAwarePolicy);
 
         if (defaultAcl != null) {
diff --git a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
index 38236fffa..b6259610c 100644
--- a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
+++ b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
@@ -34,6 +34,9 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.ArgumentMatchers.anyInt;
 import static org.mockito.ArgumentMatchers.anyLong;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.ArgumentMatchers.longThat;
+import static org.mockito.Mockito.atMost;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.never;
 import static org.mockito.Mockito.verify;
@@ -47,6 +50,9 @@ import static org.mockito.Mockito.when;
 public class ConnectionAwareRetryPolicyTest {
 
     private static final int SESSION_TIMEOUT_MS = 20_000;
+    private static final int MAX_RETRIES = 5;
+    // Must stay in sync with ConnectionAwareRetryPolicy's RECONNECT_JITTER_MS bound.
+    private static final long JITTER_BOUND_MS = 100L;
 
     private final RetryPolicy delegate = mock(RetryPolicy.class);
     private final RetrySleeper sleeper = mock(RetrySleeper.class);
@@ -58,7 +64,8 @@ public class ConnectionAwareRetryPolicyTest {
         Listenable<ConnectionStateListener> listenable = mock(Listenable.class);
         when(zk.getConnectionStateListenable()).thenReturn(listenable);
 
-        ConnectionAwareRetryPolicy policy = new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS);
+        ConnectionAwareRetryPolicy policy =
+            new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS, MAX_RETRIES);
         policy.bind(zk);
 
         ArgumentCaptor<ConnectionStateListener> captor = ArgumentCaptor.forClass(ConnectionStateListener.class);
@@ -76,7 +83,8 @@ public class ConnectionAwareRetryPolicyTest {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
         when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
 
-        // Default state is CONNECTED (no event fired).
+        // Default state is CONNECTED (no event fired). The backoff policy is correct here: a
+        // retryable error while connected is not a disconnect, so normal exponential backoff applies.
         boolean result = policy.allowRetry(0, 0L, sleeper);
 
         assertTrue(result, "when connected, must propagate the delegate's decision");
@@ -96,60 +104,64 @@ public class ConnectionAwareRetryPolicyTest {
     }
 
     @Test
-    public void suspendedBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+    public void suspendedWaitsForReconnectWithoutBlindBackoffThenRetries() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
         when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(0, 0L, sleeper);
+        boolean result = policy.allowRetry(0, 0L, sleeper); // retryCount 0 < MAX_RETRIES
 
         assertTrue(result, "should retry once the SendThread has reconnected");
         verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+        // The whole point of this policy: on SUSPENDED/LOST it must NOT invoke the delegate's
+        // exponential backoff sleep. Reaching the delegate here would reintroduce the blind backoff.
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
+        // Only bounded reconnection jitter may sleep — never a multi-second backoff.
+        verify(sleeper, atMost(1)).sleepFor(longThat(ms -> ms >= 0 && ms < JITTER_BOUND_MS),
+            eq(TimeUnit.MILLISECONDS));
     }
 
     @Test
-    public void lostBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+    public void lostWaitsForReconnectWithoutBlindBackoffThenRetries() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
         when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
 
         fire(ConnectionState.LOST);
-        boolean result = policy.allowRetry(3, 1234L, sleeper);
+        boolean result = policy.allowRetry(3, 1234L, sleeper); // 3 < MAX_RETRIES
 
         assertTrue(result, "LOST should also wait for reconnection then retry");
         verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 
     @Test
     public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
         when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(false);
 
         fire(ConnectionState.SUSPENDED);
         boolean result = policy.allowRetry(0, 0L, sleeper);
 
         assertFalse(result, "should abandon when not reconnected within the session timeout");
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 
     @Test
-    public void suspendedHonoursDelegateRetryBudget() throws Exception {
+    public void suspendedAbandonsWhenRetryBudgetExhausted() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        // Delegate says no more retries allowed (budget exhausted)
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(100, 99999L, sleeper);
+        // retryCount == MAX_RETRIES → budget exhausted; must give up without waiting or backing off.
+        boolean result = policy.allowRetry(MAX_RETRIES, 99999L, sleeper);
 
-        assertFalse(result, "should abandon when the delegate's retry budget is exhausted");
+        assertFalse(result, "should abandon when the retry budget is exhausted");
         verify(zk, never()).blockUntilConnected(anyInt(), any());
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 
     @Test
     public void interruptedWhileWaitingReturnsFalseAndPreservesInterruptFlag() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
         when(zk.blockUntilConnected(anyInt(), any())).thenThrow(new InterruptedException("test"));
 
         fire(ConnectionState.SUSPENDED);
@@ -160,15 +172,15 @@ public class ConnectionAwareRetryPolicyTest {
     }
 
     @Test
-    public void suspendedWithNullFrameworkFallsThroughToDelegate() {
-        // zkSupplier returns null (framework not yet available), even though bind() ran on the mock.
+    public void suspendedWithinBudgetButNullFrameworkAllowsRetry() {
+        // zkSupplier returns null (framework not yet bound), even though bind() ran on the mock.
         ConnectionAwareRetryPolicy policy = build(() -> null);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(1, 50L, sleeper);
+        boolean result = policy.allowRetry(1, 50L, sleeper); // 1 < MAX_RETRIES
 
-        assertTrue(result, "with no framework available yet, must fall through to the delegate");
-        verify(delegate).allowRetry(1, 50L, sleeper);
+        assertTrue(result, "within budget but framework not available yet → allow the retry");
+        // Still no blind backoff on the suspended path.
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 }

(Separately from this — the blast-radius design question still stands for @rzo1 / @GGraziadei.)

Comment on lines +107 to +116
// Honour the configured retry budget even while suspended.
if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
return false;
}

CuratorFramework zk = zkSupplier.get();
if (zk == null) {
// Framework not yet available; delegate already approved the retry
return true;
}

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Core change: replace the budget check that calls the sleeping delegate with a pure retry-count cap, so the suspended path no longer reintroduces the exponential backoff sleep. Depends on the new maxRetries field/constructor arg and the CuratorUtils wiring in the patch above — apply those together or this won't compile.

Suggested change
// Honour the configured retry budget even while suspended.
if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
return false;
}
CuratorFramework zk = zkSupplier.get();
if (zk == null) {
// Framework not yet available; delegate already approved the retry
return true;
}
// Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On
// SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected()
// instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff
// as a side effect and reintroduce exactly the blind backoff this policy exists to avoid.
if (retryCount >= maxRetries) {
LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry",
state, maxRetries, retryCount);
return false;
}
CuratorFramework zk = zkSupplier.get();
if (zk == null) {
// Framework not yet available; within budget, so allow the retry.
return true;
}

@GGraziadei

Copy link
Copy Markdown
Member

Thanks @jkrauss82 for the report and fix, and @reiabreu for the review.

I would leave the retry behaviour of the other Curator clients unchanged in this PR.

  • The supervisor crashes because a ConnectionLoss escapes on the HBTimer thread and the uncaught exception brings down the process. Catching it around the ZK write and letting the next heartbeat cycle retry seems enough to fix this specific issue.

  • newCurator() is used by all ClientZookeeper clients (Nimbus, supervisors and workers), the blobstore and Trident's TransactionalState. Changing the retry policy there could make operations that currently fail after the configured backoff block the calling thread for up to storm.zookeeper.session.timeout (20s) on each retry. I don't think we should introduce that change across the board based on one supervisor incident.

  • The exception happens 603 ms after the socket is closed, and just 21 ms after SUSPENDED. With the default settings, the first backoff in StormBoundedExponentialBackoffRetry is at least 1000 ms, so it looks like the retry loop didn't even get through its first sleep. This makes me wonder whether the retry policy was actually consulted, or whether it rejected the retry straight away. If that's the case, ConnectionAwareRetryPolicy wouldn't have made a difference here. The catch is what fixes the failure we're seeing.

I'd keep this PR focused on the SupervisorHeartbeat change and handle ConnectionAwareRetryPolicy separately, with tests against an actual ZooKeeper failover, including the Nimbus paths.

@jkrauss82, could you share the full stack trace for the ConnectionLoss? It would help us understand exactly where it comes from and clarify the last point.

Happy to reconsider if the stack trace shows something different. What do you think?

@jkrauss82

Copy link
Copy Markdown
Author

thanks for looking into this @GGraziadei

I don't have enough knowledge of the source and the actual inner working of the various components using Curator to contribute meaningful to the discussion regarding the scope of the PR. but I hope the logs below provide further insight. happy to provide more if necessary.

also sharing this summary my local AI gave me regarding the log content, just to provide further context as it includes the various relevant timestamps:

  1. EndOfStreamException at 08:39:53.907 — socket closed by ZK server, with the full stack trace from ClientCnxnSocketNIO.doIO

  2. Full ConnectionLossException stack trace at 08:39:54.510:

    • SupervisorHeartbeat.run(SupervisorHeartbeat.java:165)
    • StormClusterStateImpl.supervisorHeartbeat
    • PaceMakerStateStorage.set_ephemeral_node
    • ClientZookeeper.existsNode
    • Curator RetryLoop.callWithRetry — proving the retry loop was consulted
    • ... 8 more
  3. Process death at 08:39:54.536 — Utils.exitProcess(20)

  4. Timeline context showing:

    • ZK session established at 08:39:06 (RECONNECTED)
    • Socket closed at 08:39:53.907 (603ms before crash)
    • ConnectionState SUSPENDED at 08:39:54.489 (21ms before crash)
    • ConnectionLoss at 08:39:54.510
    • exitProcess at 08:39:54.536 (629ms total)

This confirms the retry loop was invoked but the underlying ZK operation failed before allowRetry() could be called.

2026-09-29 08:38:44.753 o.a.s.h.HealthChecker EventTimer [INFO] The supervisor healthchecks succeeded.
2026-09-29 08:38:42.214 o.a.s.s.o.a.c.f.s.ConnectionStateManager Curator-ConnectionStateManager-0 [WARN] Session timeout has elapsed while SUSPENDED. Injecting a session expiration. Elapsed ms: 36013. Adjusted session timeout ms: 20000
2026-09-29 08:39:06.405 o.a.s.s.o.a.z.ClientCnxnSocket main-EventThread [INFO] jute.maxbuffer value is 1048575 Bytes
2026-09-29 08:39:06.442 o.a.s.s.o.a.z.ClientCnxn main-EventThread [INFO] zookeeper.request.timeout value is 0. feature enabled=false
2026-09-29 08:39:06.549 o.a.s.s.o.a.z.ZooKeeperTestable Curator-ConnectionStateManager-0 [INFO] injectSessionExpiration() called
2026-09-29 08:39:06.550 o.a.s.s.o.a.z.ClientCnxn main-SendThread() [WARN] Session 0x0 for server nimbus-node-2/<ZK-node-2-ip>:2181, Closing socket connection. Attempting reconnect except it is a SessionExpiredException or SessionTimeoutException.
java.io.IOException: Connection has already been closed and reconnection is not allowed
        at org.apache.storm.shade.org.apache.zookeeper.ClientCnxn$SendThread.changeZkState(ClientCnxn.java:991)
        at org.apache.storm.shade.org.apache.zookeeper.ClientCnxn$SendThread.startConnect(ClientCnxn.java:1142)
        at org.apache.storm.shade.org.apache.zookeeper.ClientCnxn$SendThread.run(ClientCnxn.java:1200)
2026-09-29 08:39:06.551 o.a.s.s.o.a.c.ConnectionState main-EventThread [WARN] Session expired event received
2026-09-29 08:39:06.552 o.a.s.s.o.a.z.ZooKeeper main-EventThread [INFO] Initiating client connection, connectString=nimbus-node-3:2181,nimbus-node-2:2181,nimbus-node-1:2181/zookeeper sessionTimeout=20000 watcher=org.apache.storm.shade.org.apache.curator.ConnectionState@36fc05ff
2026-09-29 08:39:06.554 o.a.s.s.o.a.c.f.s.ConnectionStateManager main-EventThread [INFO] State change: LOST
2026-09-29 08:39:06.561 o.a.s.s.o.a.z.ClientCnxnSocket main-EventThread [INFO] jute.maxbuffer value is 1048575 Bytes
2026-09-29 08:39:06.577 o.a.s.s.o.a.z.ClientCnxn main-EventThread [INFO] zookeeper.request.timeout value is 0. feature enabled=false
2026-09-29 08:39:06.732 o.a.s.s.o.a.z.ClientCnxn main-EventThread [INFO] EventThread shut down for session: 0x0
2026-09-29 08:39:06.742 o.a.s.s.o.a.z.ClientCnxn main-EventThread [INFO] EventThread shut down for session: 0x300ff9e3d430001
2026-09-29 08:39:06.862 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-1:2181) [INFO] Opening socket connection to server nimbus-node-1/<ZK-node-1-ip>:2181.
2026-09-29 08:39:06.862 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-1:2181) [INFO] SASL config status: Will not attempt to authenticate using SASL (unknown error)
2026-09-29 08:39:06.914 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-1:2181) [INFO] Socket connection established, initiating session, client: /<supervisor-client-ip>:36856, server: nimbus-node-1/<ZK-node-1-ip>:2181
2026-09-29 08:39:06.931 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-1:2181) [INFO] Session establishment complete on server nimbus-node-1/<ZK-node-1-ip>:2181, session id = 0x300ff9e3d43024a, negotiated timeout = 20000
2026-09-29 08:39:06.932 o.a.s.s.o.a.c.f.s.ConnectionStateManager main-EventThread [INFO] State change: RECONNECTED
2026-09-29 08:39:53.911 o.a.s.d.s.t.SupervisorHealthCheck EventTimer [INFO] Running supervisor healthchecks...
2026-09-29 08:39:53.911 o.a.s.h.HealthChecker EventTimer [INFO] The supervisor healthchecks succeeded.
2026-09-29 08:39:53.907 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-1:2181) [WARN] Session 0x300ff9e3d43024a for server nimbus-node-1/<ZK-node-1-ip>:2181, Closing socket connection. Attempting reconnect except it is a SessionExpiredException or SessionTimeoutException.
EndOfStreamException: Unable to read additional data from server sessionid 0x300ff9e3d43024a, likely server has closed socket
        at org.apache.storm.shade.org.apache.zookeeper.ClientCnxnSocketNIO.doIO(ClientCnxnSocketNIO.java:77)
        at org.apache.storm.shade.org.apache.zookeeper.ClientCnxnSocketNIO.doTransport(ClientCnxnSocketNIO.java:356)
        at org.apache.storm.shade.org.apache.zookeeper.ClientCnxn$SendThread.run(ClientCnxn.java:1290)
2026-09-29 08:39:53.988 o.a.s.s.o.a.c.f.i.EnsembleTracker main-EventThread [INFO] New config event received: {server.2=nimbus-node-2:2888:3888:participant;0.0.0.0:2181, server.1=nimbus-node-3:2888:3888:participant;0.0.0.0:2181, server.3=nimbus-node-1:2888:3888:participant;0.0.0.0:2181, version=0}
2026-09-29 08:39:54.364 o.a.s.l.AsyncLocalizer AsyncLocalizer Task Executor - 2 [INFO] Finish cleanup
2026-09-29 08:39:54.365 o.a.s.l.AsyncLocalizer AsyncLocalizer Task Executor - 2 [INFO] Starting cleanup
2026-09-29 08:39:54.489 o.a.s.s.o.a.c.f.s.ConnectionStateManager main-EventThread [INFO] State change: SUSPENDED
2026-09-29 08:39:54.510 o.a.s.d.s.DefaultUncaughtExceptionHandler HBTimer [ERROR] Error when processing event
java.lang.RuntimeException: org.apache.storm.shade.org.apache.zookeeper.KeeperException$ConnectionLossException: KeeperErrorCode = ConnectionLoss for /supervisors
        at org.apache.storm.utils.Utils.wrapInRuntime(Utils.java:504)
        at org.apache.storm.zookeeper.ClientZookeeper.existsNode(ClientZookeeper.java:147)
        at org.apache.storm.zookeeper.ClientZookeeper.mkdirsImpl(ClientZookeeper.java:289)
        at org.apache.storm.zookeeper.ClientZookeeper.mkdirs(ClientZookeeper.java:70)
        at org.apache.storm.cluster.ZKStateStorage.set_ephemeral_node(ZKStateStorage.java:125)
        at org.apache.storm.cluster.PaceMakerStateStorage.set_ephemeral_node(PaceMakerStateStorage.java:79)
        at org.apache.storm.cluster.StormClusterStateImpl.supervisorHeartbeat(StormClusterStateImpl.java:531)
        at org.apache.storm.daemon.supervisor.timer.SupervisorHeartbeat.run(SupervisorHeartbeat.java:165)
        at org.apache.storm.StormTimer$1.run(StormTimer.java:110)
        at org.apache.storm.StormTimer$StormTimerTask.run(StormTimer.java:226)
Caused by: org.apache.storm.shade.org.apache.zookeeper.KeeperException$ConnectionLossException: KeeperErrorCode = ConnectionLoss for /supervisors
        at org.apache.storm.shade.org.apache.zookeeper.KeeperException.create(KeeperException.java:101)
        at org.apache.storm.shade.org.apache.zookeeper.KeeperException.create(KeeperException.java:53)
        at org.apache.storm.shade.org.apache.zookeeper.ZooKeeper.exists(ZooKeeper.java:1867)
        at org.apache.storm.shade.org.apache.curator.framework.imps.ExistsBuilderImpl$3.call(ExistsBuilderImpl.java:247)
        at org.apache.storm.shade.org.apache.curator.framework.imps.ExistsBuilderImpl$3.call(ExistsBuilderImpl.java:240)
        at org.apache.storm.shade.org.apache.curator.RetryLoop.callWithRetry(RetryLoop.java:88)
        at org.apache.storm.shade.org.apache.curator.framework.imps.ExistsBuilderImpl.pathInForegroundStandard(ExistsBuilderImpl.java:240)
        at org.apache.storm.shade.org.apache.curator.framework.imps.ExistsBuilderImpl.pathInForeground(ExistsBuilderImpl.java:235)
        at org.apache.storm.shade.org.apache.curator.framework.imps.ExistsBuilderImpl.forPath(ExistsBuilderImpl.java:202)
        at org.apache.storm.shade.org.apache.curator.framework.imps.ExistsBuilderImpl.forPath(ExistsBuilderImpl.java:35)
        at org.apache.storm.zookeeper.ClientZookeeper.existsNode(ClientZookeeper.java:144)
        ... 8 more
2026-09-29 08:39:54.536 o.a.s.u.Utils HBTimer [ERROR] Halting process: Error when processing an event
java.lang.RuntimeException: Halting process: Error when processing an event
        at org.apache.storm.utils.Utils.exitProcess(Utils.java:527)
        at org.apache.storm.daemon.supervisor.DefaultUncaughtExceptionHandler.uncaughtException(DefaultUncaughtExceptionHandler.java:25)
        at org.apache.storm.StormTimer$StormTimerTask.run(StormTimer.java:253)
2026-09-29 08:39:56.670 o.a.s.u.Utils ShutdownHook-sleepKill-1s [INFO] Halting after 1 seconds
2026-09-29 08:40:01.456 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-3:2181) [INFO] Opening socket connection to server nimbus-node-3/<ZK-node-3-ip>:2181.
2026-09-29 08:40:02.931 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-3:2181) [INFO] SASL config status: Will not attempt to authenticate using SASL (unknown error)
2026-09-29 08:40:14.005 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-3:2181) [WARN] Client session timed out, have not heard from server in 57673ms for session id 0x300ff9e3d43024a
2026-09-29 08:40:14.500 o.a.s.s.o.a.c.f.s.ConnectionStateManager Curator-ConnectionStateManager-0 [WARN] Session timeout has elapsed while SUSPENDED. Injecting a session expiration. Elapsed ms: 20011. Adjusted session timeout ms: 20000
2026-09-29 08:40:14.965 o.a.s.s.o.a.z.ZooKeeperTestable Curator-ConnectionStateManager-0 [INFO] injectSessionExpiration() called
2026-09-29 08:40:16.967 o.a.s.s.o.a.z.ClientCnxn main-SendThread(nimbus-node-3:2181) [WARN] Session 0x300ff9e3d43024a for server nimbus-node-3/<ZK-node-3-ip>:2181, Closing socket connection. Attempting reconnect except it is a SessionExpiredException or SessionTimeoutException.
org.apache.storm.shade.org.apache.zookeeper.ClientCnxn$SessionTimeoutException: Client session timed out, have not heard from server in 57673ms for session id 0x300ff9e3d43024a
        at org.apache.storm.shade.org.apache.zookeeper.ClientCnxn$SendThread.run(ClientCnxn.java:1252)
2026-09-29 08:40:20.908 o.a.s.d.s.Slot SLOT_6726 [WARN] SLOT 6726: HB is too old 122000 > 120000 for topology: topology-a-76-1790230326
2026-09-29 08:40:21.187 o.a.s.d.s.Slot SLOT_6826 [WARN] SLOT 6826: HB is too old 122000 > 120000 for topology: topology-h-376-1790234868
2026-09-29 08:40:23.729 o.a.s.s.o.a.c.ConnectionState main-EventThread [WARN] Session expired event received
2026-09-29 08:40:24.995 o.a.s.d.s.Slot SLOT_6709 [WARN] SLOT 6709: HB is too old 121000 > 120000 for topology: topology-f-26-1790229791
2026-09-29 08:40:25.071 o.a.s.d.s.Container SLOT_6826 [INFO] Killing be2391fd-02e0-4034-af67-fa9114bb8b23-127.0.1.1:d49090d8-1538-4beb-aba1-cc6f0462b37b
2026-09-29 08:40:24.421 o.a.s.d.s.Slot SLOT_6706 [WARN] SLOT 6706: HB is too old 174000 > 120000 for topology: topology-k-17-1790229697
2026-09-29 08:40:23.911 o.a.s.d.s.t.SupervisorHealthCheck EventTimer [INFO] Running supervisor healthchecks...
2026-09-29 08:40:27.905 o.a.s.h.HealthChecker EventTimer [INFO] The supervisor healthchecks succeeded.
2026-09-29 08:40:28.012 o.a.s.d.s.Slot SLOT_6708 [WARN] SLOT 6708: HB is too old 129000 > 120000 for topology: topology-g-23-1790229761
2026-09-29 08:40:28.805 o.a.s.s.o.a.z.ZooKeeper main-EventThread [INFO] Initiating client connection, connectString=nimbus-node-3:2181,nimbus-node-2:2181,nimbus-node-1:2181/zookeeper sessionTimeout=20000 watcher=org.apache.storm.shade.org.apache.curator.ConnectionState@36fc05ff
2026-09-29 08:40:32.373 o.a.s.d.s.Slot SLOT_6773 [WARN] SLOT 6773: HB is too old 132000 > 120000 for topology: topology-c-217-1790231999
2026-09-29 08:40:30.663 o.a.s.d.s.Slot SLOT_6832 [WARN] SLOT 6832: HB is too old 129000 > 120000 for topology: topology-i-394-1790235279
2026-09-29 08:40:30.682 o.a.s.d.s.Slot SLOT_6717 [WARN] SLOT 6717: HB is too old 122000 > 120000 for topology: topology-j-50-1790230043
2026-09-29 08:40:30.682 o.a.s.d.s.Slot SLOT_6812 [WARN] SLOT 6812: HB is too old 124000 > 120000 for topology: topology-e-335-1790234003
2026-09-29 08:40:50.664 o.a.s.s.o.a.c.f.s.ConnectionStateManager Curator-ConnectionStateManager-0 [WARN] Session timeout has elapsed while SUSPENDED. Injecting a session expiration. Elapsed ms: 36164. Adjusted session timeout ms: 20000
2026-09-29 08:40:30.682 o.a.s.d.s.Slot SLOT_6720 [WARN] SLOT 6720: HB is too old 130000 > 120000 for topology: topology-b-59-1790230141
2026-09-29 08:40:30.682 o.a.s.d.s.Slot SLOT_6803 [WARN] SLOT 6803: HB is too old 124000 > 120000 for topology: topology-d-459-1790338504

@reiabreu

reiabreu commented Oct 3, 2026

Copy link
Copy Markdown
Contributor

Written with Claude (an LLM). Claude ran all the tests described below and drafted this comment; I reviewed it before posting.

Summary: Claude's failover tests found that ConnectionAwareRetryPolicy as written can block ZooKeeper's event thread, which delays the "reconnected" event. A fix for this PR is below. With it, the policy behaved like stock retries in the same tests.

Problem. For background (inBackground) operations, the kind Nimbus's LeaderLatch uses, Curator calls allowRetry on ZooKeeper's event thread. blockUntilConnected waits for a flag that only an event on that same thread can set, so the wait can't end early. Claude ran the policy against a real ZooKeeper (Curator's TestingServer, Storm's shaded client), killing and restarting the server under load, next to a stock-backoff control. Steady background reads across a 4 s outage:

stock backoff as written with the fix
operations failed 0 2 0
"reconnected" notice after the server was back 0.5-1.4 s 16.3 s 0.47 s
event-thread blocking none 2 blocks, about 18 s none
session up went LOST up

A forced version is deterministic: allowRetry on the event thread blocked 10,014 ms and returned false although the server had been back for 8.5 s. With the fix it returns in 107 ms.

Fix for this PR. Wait for the reconnect only on foreground retries and delegate background retries to the wrapped policy. Curator passes RetryLoop.getDefaultRetrySleeper() for foreground retries and the operation itself for background ones (Curator 5.9.0 source). If that check ever stops matching, the policy falls back to stock behavior.

private static boolean isForegroundRetry(RetrySleeper sleeper) {
    return sleeper == RetryLoop.getDefaultRetrySleeper();
}

// in allowRetry(...):
if (isForegroundRetry(sleepSleeper) && (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST)) {
    ...   // wait for reconnect, as before
}
return delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper);

The patch below also bounds the suspended path by retry count instead of calling the delegate (which sleeps its backoff first), wires maxRetries in CuratorUtils, and updates ConnectionAwareRetryPolicyTest to 10 cases, including two that assert background retries never block. These tests and CuratorUtilsTest pass locally.

Trade-off. Background operations no longer survive a long outage. In an 8 s outage about 3,000 of 5,000 background reads gave up with the fix, the same as stock. The policy as written lost none, but only by blocking the event thread.

Limits.

  • Claude could not reproduce the original crash in a test. Stock retries survived a 2 s outage with 16 threads of synchronous reads, so the foreground benefit rests on the production report and the unit tests.
  • Traffic was synthetic (about 20,000 operations per second), one run per cell. Nimbus HA and session expiry are untested.

Re @GGraziadei's question (was the policy consulted?). The posted trace can't tell. RetryLoop.java:88 is proc.call(), and a refused retry rethrows the same exception object. The supervisor's effective storm.zookeeper.retry.times / .interval, or Curator's DEBUG logs from RetryLoopImpl, would. The log also shows SUSPENDED periods longer than the 20 s session timeout (36 s, 20 s, 36 s), where waiting can't help; the catch is what keeps the process alive. If maintainers prefer this PR to stay with the catch, the fix above can go in a follow-up.

Patch (on top of ab67921)
diff --git a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
index eb041bdef..e3f90eb6c 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicy.java
@@ -23,6 +23,7 @@ import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicReference;
 import java.util.function.Supplier;
 import org.apache.storm.shade.com.google.common.annotations.VisibleForTesting;
+import org.apache.storm.shade.org.apache.curator.RetryLoop;
 import org.apache.storm.shade.org.apache.curator.RetryPolicy;
 import org.apache.storm.shade.org.apache.curator.RetrySleeper;
 import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
@@ -44,7 +45,13 @@ import org.slf4j.LoggerFactory;
  * the operation on the new connection. If the connection cannot be
  * re-established within the session timeout, the retry is abandoned.
  *
- * <p>For all other connection states, this policy delegates to the wrapped
+ * <p>Waiting is only done for <em>foreground</em> retries, which run on the caller's own thread.
+ * Curator also retries <em>background</em> ({@code inBackground}) operations, and there the policy is
+ * invoked from ZooKeeper's single event thread or from Curator's background thread. Blocking the event
+ * thread delays delivery of the very "connected" event that {@code blockUntilConnected} is waiting for, so
+ * background retries always delegate to the wrapped policy, whose sleeper only schedules a re-queue.
+ *
+ * <p>For all other connection states, and for background retries, this policy delegates to the wrapped
  * delegate policy (typically a {@link StormBoundedExponentialBackoffRetry}).
  */
 public class ConnectionAwareRetryPolicy implements RetryPolicy {
@@ -60,6 +67,7 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
     private final RetryPolicy delegate;
     private final Supplier<CuratorFramework> zkSupplier;
     private final int sessionTimeoutMs;
+    private final int maxRetries;
     private final AtomicReference<ConnectionState> connectionState =
         new AtomicReference<>(ConnectionState.CONNECTED);
 
@@ -69,13 +77,19 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
      *                         the framework does not exist until {@code builder.build()} returns
      * @param sessionTimeoutMs upper bound (in ms) for waiting on reconnection; typically
      *                         {@code storm.zookeeper.session.timeout}
+     * @param maxRetries       cap on the number of retries on the SUSPENDED/LOST path, applied without a
+     *                         backoff sleep. Should match the delegate's retry budget
+     *                         ({@code storm.zookeeper.retry.times}) so the two paths give up after the
+     *                         same number of attempts.
      */
     public ConnectionAwareRetryPolicy(RetryPolicy delegate,
                                       Supplier<CuratorFramework> zkSupplier,
-                                      int sessionTimeoutMs) {
+                                      int sessionTimeoutMs,
+                                      int maxRetries) {
         this.delegate = delegate;
         this.zkSupplier = zkSupplier;
         this.sessionTimeoutMs = sessionTimeoutMs;
+        this.maxRetries = maxRetries;
     }
 
     /**
@@ -99,19 +113,33 @@ public class ConnectionAwareRetryPolicy implements RetryPolicy {
         return delegate;
     }
 
+    /**
+     * Curator's foreground retry loop passes {@link RetryLoop#getDefaultRetrySleeper()} (a real sleep on the
+     * caller's thread). Background operations pass the operation itself, whose {@code sleepFor} only records a
+     * re-queue time, and whose retry runs on a thread that must not block.
+     */
+    private static boolean isForegroundRetry(RetrySleeper sleeper) {
+        return sleeper == RetryLoop.getDefaultRetrySleeper();
+    }
+
     @Override
     public boolean allowRetry(int retryCount, long elapsedTimeMs, RetrySleeper sleepSleeper) {
         ConnectionState state = connectionState.get();
 
-        if (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST) {
-            // Honour the configured retry budget even while suspended.
-            if (!delegate.allowRetry(retryCount, elapsedTimeMs, sleepSleeper)) {
+        if (isForegroundRetry(sleepSleeper) && (state == ConnectionState.SUSPENDED || state == ConnectionState.LOST)) {
+            // Cap total attempts, but WITHOUT the delegate's exponential backoff sleep. On
+            // SUSPENDED/LOST this policy yields to the ZK SendThread via blockUntilConnected()
+            // instead of blind-sleeping; calling delegate.allowRetry() here would sleep the backoff
+            // as a side effect and reintroduce exactly the blind backoff this policy exists to avoid.
+            if (retryCount >= maxRetries) {
+                LOG.warn("ZK connection {} and retry budget ({}) exhausted on retry {}, abandoning retry",
+                    state, maxRetries, retryCount);
                 return false;
             }
 
             CuratorFramework zk = zkSupplier.get();
             if (zk == null) {
-                // Framework not yet available; delegate already approved the retry
+                // Framework not yet available; within budget, so allow the retry.
                 return true;
             }
 
diff --git a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
index d2365bd1a..d15e2b143 100644
--- a/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
+++ b/storm-client/src/jvm/org/apache/storm/utils/CuratorUtils.java
@@ -63,13 +63,15 @@ public class CuratorUtils {
         // Wrap with a connection-aware policy that yields to SendThread on SUSPENDED/LOST.
         AtomicReference<CuratorFramework> zkRef = new AtomicReference<>();
         int sessionTimeoutMs = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_SESSION_TIMEOUT));
+        int retryTimes = ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES));
         ConnectionAwareRetryPolicy connectionAwarePolicy = new ConnectionAwareRetryPolicy(
                 new StormBoundedExponentialBackoffRetry(
                         ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL)),
                         ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_INTERVAL_CEILING)),
-                        ObjectReader.getInt(conf.get(Config.STORM_ZOOKEEPER_RETRY_TIMES))),
+                        retryTimes),
                 zkRef::get,
-                sessionTimeoutMs);
+                sessionTimeoutMs,
+                retryTimes);
         builder.retryPolicy(connectionAwarePolicy);
 
         if (defaultAcl != null) {
diff --git a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
index 38236fffa..5f50cfc66 100644
--- a/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
+++ b/storm-client/test/jvm/org/apache/storm/utils/ConnectionAwareRetryPolicyTest.java
@@ -20,6 +20,7 @@ package org.apache.storm.utils;
 
 import java.util.concurrent.TimeUnit;
 import java.util.function.Supplier;
+import org.apache.storm.shade.org.apache.curator.RetryLoop;
 import org.apache.storm.shade.org.apache.curator.RetryPolicy;
 import org.apache.storm.shade.org.apache.curator.RetrySleeper;
 import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
@@ -43,13 +44,19 @@ import static org.mockito.Mockito.when;
  * Unit tests for {@link ConnectionAwareRetryPolicy}. The ZooKeeper connection state is driven by
  * capturing the {@link ConnectionStateListener} that {@link ConnectionAwareRetryPolicy#bind} registers
  * and firing state transitions at it, so no real ZooKeeper is required.
+ *
+ * <p>Curator passes {@link RetryLoop#getDefaultRetrySleeper()} for foreground retries and the operation
+ * itself for background retries. The policy waits only for the former, so the tests use the real default
+ * sleeper to exercise the waiting path and a mock sleeper to stand in for a background retry.
  */
 public class ConnectionAwareRetryPolicyTest {
 
     private static final int SESSION_TIMEOUT_MS = 20_000;
+    private static final int MAX_RETRIES = 5;
 
     private final RetryPolicy delegate = mock(RetryPolicy.class);
-    private final RetrySleeper sleeper = mock(RetrySleeper.class);
+    private final RetrySleeper foregroundSleeper = RetryLoop.getDefaultRetrySleeper();
+    private final RetrySleeper backgroundSleeper = mock(RetrySleeper.class);
     private final CuratorFramework zk = mock(CuratorFramework.class);
     private ConnectionStateListener listener;
 
@@ -58,7 +65,8 @@ public class ConnectionAwareRetryPolicyTest {
         Listenable<ConnectionStateListener> listenable = mock(Listenable.class);
         when(zk.getConnectionStateListenable()).thenReturn(listenable);
 
-        ConnectionAwareRetryPolicy policy = new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS);
+        ConnectionAwareRetryPolicy policy =
+            new ConnectionAwareRetryPolicy(delegate, zkSupplier, SESSION_TIMEOUT_MS, MAX_RETRIES);
         policy.bind(zk);
 
         ArgumentCaptor<ConnectionStateListener> captor = ArgumentCaptor.forClass(ConnectionStateListener.class);
@@ -76,11 +84,12 @@ public class ConnectionAwareRetryPolicyTest {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
         when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
 
-        // Default state is CONNECTED (no event fired).
-        boolean result = policy.allowRetry(0, 0L, sleeper);
+        // Default state is CONNECTED (no event fired). The backoff policy is correct here: a
+        // retryable error while connected is not a disconnect, so normal exponential backoff applies.
+        boolean result = policy.allowRetry(0, 0L, foregroundSleeper);
 
         assertTrue(result, "when connected, must propagate the delegate's decision");
-        verify(delegate).allowRetry(0, 0L, sleeper);
+        verify(delegate).allowRetry(0, 0L, foregroundSleeper);
     }
 
     @Test
@@ -89,86 +98,116 @@ public class ConnectionAwareRetryPolicyTest {
         when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
 
         fire(ConnectionState.RECONNECTED);
-        boolean result = policy.allowRetry(2, 100L, sleeper);
+        boolean result = policy.allowRetry(2, 100L, foregroundSleeper);
 
         assertFalse(result, "RECONNECTED is a healthy state and must delegate");
-        verify(delegate).allowRetry(2, 100L, sleeper);
+        verify(delegate).allowRetry(2, 100L, foregroundSleeper);
     }
 
     @Test
-    public void suspendedBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+    public void suspendedForegroundRetryWaitsForReconnectWithoutBlindBackoff() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
         when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(0, 0L, sleeper);
+        boolean result = policy.allowRetry(0, 0L, foregroundSleeper); // retryCount 0 < MAX_RETRIES
 
         assertTrue(result, "should retry once the SendThread has reconnected");
         verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+        // On SUSPENDED/LOST the delegate's exponential backoff sleep must not run: reaching it here
+        // would reintroduce the blind backoff this policy exists to avoid.
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 
     @Test
-    public void lostBlocksUntilConnectedThenRetriesImmediately() throws Exception {
+    public void lostForegroundRetryWaitsForReconnectWithoutBlindBackoff() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
         when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(true);
 
         fire(ConnectionState.LOST);
-        boolean result = policy.allowRetry(3, 1234L, sleeper);
+        boolean result = policy.allowRetry(3, 1234L, foregroundSleeper); // 3 < MAX_RETRIES
 
         assertTrue(result, "LOST should also wait for reconnection then retry");
         verify(zk).blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS);
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 
     @Test
-    public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
+    public void suspendedBackgroundRetryNeverBlocksAndDelegates() throws Exception {
+        // Background retries run on ZooKeeper's event thread (or Curator's background thread). Blocking there
+        // delays the "connected" event the wait is for, so they must go to the delegate, whose sleeper only
+        // schedules a re-queue.
         ConnectionAwareRetryPolicy policy = build(() -> zk);
         when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
+
+        fire(ConnectionState.SUSPENDED);
+        boolean result = policy.allowRetry(0, 0L, backgroundSleeper);
+
+        assertTrue(result, "a background retry must get the delegate's answer");
+        verify(delegate).allowRetry(0, 0L, backgroundSleeper);
+        verify(zk, never()).blockUntilConnected(anyInt(), any());
+    }
+
+    @Test
+    public void lostBackgroundRetryNeverBlocksAndDelegates() throws Exception {
+        ConnectionAwareRetryPolicy policy = build(() -> zk);
+        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
+
+        fire(ConnectionState.LOST);
+        boolean result = policy.allowRetry(2, 50L, backgroundSleeper);
+
+        assertFalse(result, "a background retry must get the delegate's answer, including a refusal");
+        verify(delegate).allowRetry(2, 50L, backgroundSleeper);
+        verify(zk, never()).blockUntilConnected(anyInt(), any());
+    }
+
+    @Test
+    public void suspendedAbandonsRetryWhenReconnectTimesOut() throws Exception {
+        ConnectionAwareRetryPolicy policy = build(() -> zk);
         when(zk.blockUntilConnected(SESSION_TIMEOUT_MS, TimeUnit.MILLISECONDS)).thenReturn(false);
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(0, 0L, sleeper);
+        boolean result = policy.allowRetry(0, 0L, foregroundSleeper);
 
         assertFalse(result, "should abandon when not reconnected within the session timeout");
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 
     @Test
-    public void suspendedHonoursDelegateRetryBudget() throws Exception {
+    public void suspendedAbandonsWhenRetryBudgetExhausted() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        // Delegate says no more retries allowed (budget exhausted)
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(false);
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(100, 99999L, sleeper);
+        // retryCount == MAX_RETRIES → budget exhausted; must give up without waiting or backing off.
+        boolean result = policy.allowRetry(MAX_RETRIES, 99999L, foregroundSleeper);
 
-        assertFalse(result, "should abandon when the delegate's retry budget is exhausted");
+        assertFalse(result, "should abandon when the retry budget is exhausted");
         verify(zk, never()).blockUntilConnected(anyInt(), any());
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 
     @Test
     public void interruptedWhileWaitingReturnsFalseAndPreservesInterruptFlag() throws Exception {
         ConnectionAwareRetryPolicy policy = build(() -> zk);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
         when(zk.blockUntilConnected(anyInt(), any())).thenThrow(new InterruptedException("test"));
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(0, 0L, sleeper);
+        boolean result = policy.allowRetry(0, 0L, foregroundSleeper);
 
         assertFalse(result, "an interrupt while waiting should abandon the retry");
         assertTrue(Thread.interrupted(), "interrupt flag must be preserved (this check also clears it)");
     }
 
     @Test
-    public void suspendedWithNullFrameworkFallsThroughToDelegate() {
-        // zkSupplier returns null (framework not yet available), even though bind() ran on the mock.
+    public void suspendedWithinBudgetButNullFrameworkAllowsRetry() {
+        // zkSupplier returns null (framework not yet bound), even though bind() ran on the mock.
         ConnectionAwareRetryPolicy policy = build(() -> null);
-        when(delegate.allowRetry(anyInt(), anyLong(), any())).thenReturn(true);
 
         fire(ConnectionState.SUSPENDED);
-        boolean result = policy.allowRetry(1, 50L, sleeper);
+        boolean result = policy.allowRetry(1, 50L, foregroundSleeper); // 1 < MAX_RETRIES
 
-        assertTrue(result, "with no framework available yet, must fall through to the delegate");
-        verify(delegate).allowRetry(1, 50L, sleeper);
+        assertTrue(result, "within budget but framework not available yet → allow the retry");
+        // Still no blind backoff on the suspended path.
+        verify(delegate, never()).allowRetry(anyInt(), anyLong(), any());
     }
 }
Failover test program (throwaway JUnit test, not for merging)
package org.apache.storm.utils;

import java.io.PrintWriter;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicReference;
import org.apache.curator.test.TestingServer;
import org.apache.storm.shade.org.apache.curator.RetryPolicy;
import org.apache.storm.shade.org.apache.curator.RetrySleeper;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFramework;
import org.apache.storm.shade.org.apache.curator.framework.CuratorFrameworkFactory;
import org.apache.storm.shade.org.apache.curator.framework.api.BackgroundCallback;
import org.apache.storm.shade.org.apache.curator.framework.state.ConnectionState;
import org.apache.storm.shade.org.apache.zookeeper.KeeperException;
import org.junit.jupiter.api.Test;

/**
 * THROWAWAY experiment (not for commit): does ConnectionAwareRetryPolicy.allowRetry() ever block on ZooKeeper's
 * event thread, and does that delay delivery of the RECONNECTED notification?
 */
public class ConnectionAwareRetryPolicyReproTest {
    private static final String OUT = "/tmp/claude-1000/-home-rui-workspace/a5338b11-f59f-4b5e-aa07-12c2598cd7bd/scratchpad/storm9153/repro_out.txt";
    private static final int SESSION_MS = 10_000;
    private static final int CONN_MS = 5_000;
    private static final int MAX_RETRIES = 5;
    private static final int BURST = 5000;

    private final List<String> out = Collections.synchronizedList(new ArrayList<>());
    private long t0;

    private void log(String s) {
        out.add(String.format("%7d ms | %s", System.currentTimeMillis() - t0, s));
    }

    /** Records every allowRetry call: thread, the policy's own view of the state, duration, result. */
    static class RecordingPolicy extends ConnectionAwareRetryPolicy {
        final List<String> calls = Collections.synchronizedList(new ArrayList<>());
        final AtomicInteger total = new AtomicInteger();
        final AtomicInteger onEventThread = new AtomicInteger();
        final AtomicInteger onEventThreadWhileDown = new AtomicInteger();
        final AtomicInteger blockedOnEventThread = new AtomicInteger();
        final AtomicInteger maxMs = new AtomicInteger();
        final AtomicInteger offEventThread = new AtomicInteger();
        final AtomicInteger offEventThreadWhileDown = new AtomicInteger();
        final AtomicInteger maxMsOffEventWhileDown = new AtomicInteger();
        final java.util.Set<String> offEventThreadNames = java.util.concurrent.ConcurrentHashMap.newKeySet();
        final long start;
        final Field stateField;

        RecordingPolicy(RetryPolicy delegate, java.util.function.Supplier<CuratorFramework> sup, long start) throws Exception {
            super(delegate, sup, SESSION_MS, MAX_RETRIES);
            this.start = start;
            stateField = ConnectionAwareRetryPolicy.class.getDeclaredField("connectionState");
            stateField.setAccessible(true);
        }

        @Override
        @SuppressWarnings("unchecked")
        public boolean allowRetry(int retryCount, long elapsedTimeMs, RetrySleeper sleeper) {
            String thread = Thread.currentThread().getName();
            boolean event = thread.endsWith("-EventThread");
            Object state;
            try {
                state = ((AtomicReference<ConnectionState>) stateField.get(this)).get();
            } catch (Exception e) {
                state = "?";
            }
            long a = System.currentTimeMillis();
            boolean r = super.allowRetry(retryCount, elapsedTimeMs, sleeper);
            long d = System.currentTimeMillis() - a;
            total.incrementAndGet();
            boolean down = state == ConnectionState.SUSPENDED || state == ConnectionState.LOST;
            if (!event) {
                offEventThread.incrementAndGet();
                offEventThreadNames.add(thread);
                if (down) {
                    offEventThreadWhileDown.incrementAndGet();
                    maxMsOffEventWhileDown.accumulateAndGet((int) d, Math::max);
                }
            }
            if (event) {
                onEventThread.incrementAndGet();
                if (down) {
                    onEventThreadWhileDown.incrementAndGet();
                }
                if (d > 500) {
                    blockedOnEventThread.incrementAndGet();
                }
            }
            maxMs.accumulateAndGet((int) d, Math::max);
            if (down || d > 200) {
                calls.add(String.format("%7d ms | thread=%s event=%s state=%s retry=%d sleeperClass=%s took=%dms -> %s",
                    System.currentTimeMillis() - start, thread, event, state, retryCount,
                    sleeper.getClass().getSimpleName(), d, r));
            }
            return r;
        }
    }

    private void experiment(String label, boolean aware, int blips, long downMs, long upMs) throws Exception {
        t0 = System.currentTimeMillis();
        out.add("");
        out.add("===== " + label + " | aware=" + aware + " blips=" + blips + " down=" + downMs + "ms up=" + upMs + "ms =====");
        try (TestingServer server = new TestingServer(true)) {
            AtomicReference<CuratorFramework> ref = new AtomicReference<>();
            RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
            RecordingPolicy rec = aware ? new RecordingPolicy(stock, ref::get, t0) : null;
            CuratorFramework zk = CuratorFrameworkFactory.builder()
                .connectString(server.getConnectString())
                .sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS)
                .retryPolicy(aware ? rec : stock).build();
            ref.set(zk);
            if (aware) {
                rec.bind(zk);
            }
            final long[] lastRestartDone = {0};
            final long[] reconnectedAt = {0};
            zk.getConnectionStateListenable().addListener((c, s) -> {
                log("state -> " + s);
                if (s == ConnectionState.RECONNECTED) {
                    reconnectedAt[0] = System.currentTimeMillis();
                }
            });
            zk.start();
            zk.blockUntilConnected(10, TimeUnit.SECONDS);
            zk.create().forPath("/r");
            out.add("negotiated session timeout = " + zk.getZookeeperClient().getZooKeeper().getSessionTimeout() + " ms");

            AtomicInteger connLoss = new AtomicInteger();
            AtomicInteger ok = new AtomicInteger();
            AtomicInteger other = new AtomicInteger();
            final AtomicInteger done = new AtomicInteger();
            final long[] allDoneAt = {0};
            final long[] probeSentAt = {0};
            final long[] probeDoneAt = {0};
            BackgroundCallback cb = (c, e) -> {
                int rc = e.getResultCode();
                if (rc == KeeperException.Code.CONNECTIONLOSS.intValue()) {
                    connLoss.incrementAndGet();
                } else if (rc == KeeperException.Code.OK.intValue()) {
                    ok.incrementAndGet();
                } else {
                    other.incrementAndGet();
                }
                if (done.incrementAndGet() == BURST) {
                    allDoneAt[0] = System.currentTimeMillis();
                }
            };

            for (int b = 0; b < blips; b++) {
                reconnectedAt[0] = 0;
                for (int i = 0; i < BURST; i++) {
                    zk.getData().inBackground(cb).forPath("/r");   // many in flight when the socket drops
                }
                log("blip " + b + ": stopping server");
                server.stop();
                if (downMs > 2000) {
                    Thread.sleep(downMs - 1000);
                    probeSentAt[0] = System.currentTimeMillis();
                    zk.getData().inBackground((c, e) -> probeDoneAt[0] = System.currentTimeMillis()).forPath("/r");
                    Thread.sleep(1000);
                } else {
                    Thread.sleep(downMs);
                }
                server.restart();
                lastRestartDone[0] = System.currentTimeMillis();
                log("blip " + b + ": server restarted");
                Thread.sleep(upMs);
            }
            long waitUntil = System.currentTimeMillis() + 40_000;
            while (reconnectedAt[0] == 0 && System.currentTimeMillis() < waitUntil) {
                Thread.sleep(50);
            }
            long doneDeadline = System.currentTimeMillis() + 40_000;
            while (done.get() < BURST && System.currentTimeMillis() < doneDeadline) {
                Thread.sleep(50);
            }
            Thread.sleep(500);
            out.add("callbacks: CONNECTIONLOSS=" + connLoss + " OK=" + ok + " other=" + other + "  (CONNECTIONLOSS = ops whose retries were EXHAUSTED; Curator hides retried failures from this callback)");
            out.add("burst callbacks completed: " + done + "/" + BURST
                + (allDoneAt[0] == 0 ? " (some NEVER completed)" : "; last one " + (allDoneAt[0] - lastRestartDone[0]) + " ms after restart (negative = before)"));
            if (probeSentAt[0] != 0) {
                out.add("probe op sent ~1s before restart: " + (probeDoneAt[0] == 0 ? "NOT COMPLETED"
                    : "completed " + (probeDoneAt[0] - lastRestartDone[0]) + " ms after restart"));
            }
            out.add("last restart -> RECONNECTED notification latency = "
                + (reconnectedAt[0] == 0 ? "NEVER (40s)" : (reconnectedAt[0] - lastRestartDone[0]) + " ms"));
            if (aware) {
                out.add("policy.allowRetry calls: total=" + rec.total + " onEventThread=" + rec.onEventThread
                    + " onEventThreadWhileSuspendedOrLost=" + rec.onEventThreadWhileDown
                    + " blockedOver500msOnEventThread=" + rec.blockedOnEventThread + " maxCallMs=" + rec.maxMs);
                out.add("   OFF event thread: calls=" + rec.offEventThread + " whileSuspendedOrLost=" + rec.offEventThreadWhileDown
                    + " longestWhileDownMs=" + rec.maxMsOffEventWhileDown + " threads=" + rec.offEventThreadNames);
                synchronized (rec.calls) {
                    int n = 0;
                    for (String c : rec.calls) {
                        if (n++ < 12) {
                            out.add("   " + c);
                        }
                    }
                    if (rec.calls.size() > 12) {
                        out.add("   ... (" + rec.calls.size() + " notable calls total)");
                    }
                }
            }
            zk.close();
        }
    }


    /**
     * Forces the suspected situation: the policy believes SUSPENDED, Curator is disconnected, and a failed-operation
     * callback runs on ZooKeeper's event thread and calls allowRetry(). Does that stall reconnect delivery?
     */
    private void mechanism() throws Exception {
        t0 = System.currentTimeMillis();
        out.add("");
        out.add("===== MECHANISM (forced): allowRetry() called from a callback on ZK's event thread while policy=SUSPENDED and Curator is disconnected =====");
        try (TestingServer server = new TestingServer(true)) {
            AtomicReference<CuratorFramework> ref = new AtomicReference<>();
            RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
            RecordingPolicy rec = new RecordingPolicy(stock, ref::get, t0);
            CuratorFramework zk = CuratorFrameworkFactory.builder().connectString(server.getConnectString())
                .sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS).retryPolicy(rec).build();
            ref.set(zk);
            rec.bind(zk);
            final java.util.concurrent.CountDownLatch suspended = new java.util.concurrent.CountDownLatch(1);
            final long[] reconnectedAt = {0};
            zk.getConnectionStateListenable().addListener((c, st) -> {
                log("state -> " + st);
                if (st == ConnectionState.SUSPENDED) {
                    suspended.countDown();
                }
                if (st == ConnectionState.RECONNECTED) {
                    reconnectedAt[0] = System.currentTimeMillis();
                }
            });
            zk.start();
            zk.blockUntilConnected(10, TimeUnit.SECONDS);
            zk.create().forPath("/r");

            final long[] cbStart = {0};
            final long[] cbEnd = {0};
            final boolean[] cbResult = {false};
            final String[] cbThread = {null};
            final int[] cbRc = {99};
            log("stopping server");
            server.stop();
            suspended.await(10, TimeUnit.SECONDS);
            Thread.sleep(300);
            // Raw ZooKeeper async call: queued at the ZK level while disconnected, fails with CONNECTIONLOSS on a later
            // failed connect attempt. Its callback runs on the event thread and calls allowRetry, as Curator's would.
            org.apache.storm.shade.org.apache.zookeeper.ZooKeeper raw = zk.getZookeeperClient().getZooKeeper();
            raw.getData("/r", false, (rc, path, ctx, data, stat) -> {
                cbThread[0] = Thread.currentThread().getName();
                cbRc[0] = rc;
                cbStart[0] = System.currentTimeMillis();
                log("raw callback on " + cbThread[0] + " rc=" + rc + " -> calling policy.allowRetry");
                cbResult[0] = rec.allowRetry(0, 0L, (time, unit) -> unit.sleep(time));
                cbEnd[0] = System.currentTimeMillis();
                log("policy.allowRetry returned " + cbResult[0] + " after " + (cbEnd[0] - cbStart[0]) + " ms");
            }, null);
            long until = System.currentTimeMillis() + 8000;
            while (cbStart[0] == 0 && System.currentTimeMillis() < until) {
                Thread.sleep(20);
            }
            out.add("callback started: " + (cbStart[0] != 0) + " on thread " + cbThread[0] + " rc=" + cbRc[0]);
            Thread.sleep(1500);
            long restartedAt = System.currentTimeMillis();
            server.restart();
            log("server restarted (callback has been parked " + (restartedAt - cbStart[0]) + " ms)");
            long end = System.currentTimeMillis() + 40_000;
            while ((cbEnd[0] == 0 || reconnectedAt[0] == 0) && System.currentTimeMillis() < end) {
                Thread.sleep(50);
            }
            Thread.sleep(300);
            out.add("RESULT: allowRetry blocked " + (cbEnd[0] == 0 ? "NEVER RETURNED" : (cbEnd[0] - cbStart[0]) + " ms")
                + " and returned " + cbResult[0]);
            out.add("RESULT: server was back " + (cbEnd[0] == 0 ? "?" : (cbEnd[0] - restartedAt) + " ms") + " before allowRetry returned");
            out.add("RESULT: RECONNECTED notification " + (reconnectedAt[0] == 0 ? "NEVER" : (reconnectedAt[0] - restartedAt) + " ms after server restart"));
            synchronized (rec.calls) {
                rec.calls.forEach(c -> out.add("   " + c));
            }
            zk.close();
        }
    }

    /** Steady traffic across the drop, to see whether the risky situation arises on its own. */
    private void natural(boolean aware) throws Exception {
        t0 = System.currentTimeMillis();
        out.add("");
        out.add("===== NATURAL: steady background traffic across a 4s outage | aware=" + aware + " =====");
        try (TestingServer server = new TestingServer(true)) {
            AtomicReference<CuratorFramework> ref = new AtomicReference<>();
            RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
            RecordingPolicy rec = aware ? new RecordingPolicy(stock, ref::get, t0) : null;
            CuratorFramework zk = CuratorFrameworkFactory.builder().connectString(server.getConnectString())
                .sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS).retryPolicy(aware ? rec : stock).build();
            ref.set(zk);
            if (aware) {
                rec.bind(zk);
            }
            final long[] reconnectedAt = {0};
            zk.getConnectionStateListenable().addListener((c, st) -> {
                log("state -> " + st);
                if (st == ConnectionState.RECONNECTED) {
                    reconnectedAt[0] = System.currentTimeMillis();
                }
            });
            zk.start();
            zk.blockUntilConnected(10, TimeUnit.SECONDS);
            zk.create().forPath("/r");

            AtomicInteger issued = new AtomicInteger();
            AtomicInteger done = new AtomicInteger();
            AtomicInteger ok = new AtomicInteger();
            AtomicInteger lost = new AtomicInteger();
            BackgroundCallback cb = (c, e) -> {
                if (e.getResultCode() == KeeperException.Code.OK.intValue()) {
                    ok.incrementAndGet();
                } else {
                    lost.incrementAndGet();
                }
                done.incrementAndGet();
            };
            java.util.concurrent.atomic.AtomicBoolean go = new java.util.concurrent.atomic.AtomicBoolean(true);
            Thread issuer = new Thread(() -> {
                while (go.get()) {
                    try {
                        zk.getData().inBackground(cb).forPath("/r");
                        issued.incrementAndGet();
                    } catch (Exception e) {
                        break;
                    }
                    java.util.concurrent.locks.LockSupport.parkNanos(20_000);
                }
            }, "issuer");
            issuer.start();
            Thread.sleep(300);
            log("stopping server (traffic continues)");
            server.stop();
            Thread.sleep(600);
            go.set(false);
            issuer.join();
            Thread.sleep(3400);
            long restartedAt = System.currentTimeMillis();
            server.restart();
            log("server restarted");
            long end = System.currentTimeMillis() + 40_000;
            while ((reconnectedAt[0] == 0 || done.get() < issued.get()) && System.currentTimeMillis() < end) {
                Thread.sleep(50);
            }
            Thread.sleep(300);
            out.add("operations issued=" + issued + " completed=" + done + " ok=" + ok + " failed(retries exhausted)=" + lost);
            out.add("RECONNECTED notification " + (reconnectedAt[0] == 0 ? "NEVER" : (reconnectedAt[0] - restartedAt) + " ms after server restart"));
            if (aware) {
                out.add("policy.allowRetry calls: total=" + rec.total + " onEventThread=" + rec.onEventThread
                    + " onEventThreadWhileSuspendedOrLost=" + rec.onEventThreadWhileDown
                    + " blockedOver500msOnEventThread=" + rec.blockedOnEventThread + " maxCallMs=" + rec.maxMs);
                out.add("   OFF event thread: calls=" + rec.offEventThread + " whileSuspendedOrLost=" + rec.offEventThreadWhileDown
                    + " longestWhileDownMs=" + rec.maxMsOffEventWhileDown + " threads=" + rec.offEventThreadNames);
                synchronized (rec.calls) {
                    int n = 0;
                    for (String c : rec.calls) {
                        if (n++ < 8) {
                            out.add("   " + c);
                        }
                    }
                    if (rec.calls.size() > 8) {
                        out.add("   ... (" + rec.calls.size() + " notable calls total)");
                    }
                }
            }
            zk.close();
        }
    }


    /** Synchronous (foreground) operations across an outage that outlasts the (shortened) stock retry budget. */
    private void foreground(boolean aware) throws Exception {
        t0 = System.currentTimeMillis();
        out.add("");
        out.add("===== FOREGROUND: 16 threads doing synchronous getData across a 2s outage | aware=" + aware + " =====");
        try (TestingServer server = new TestingServer(true)) {
            AtomicReference<CuratorFramework> ref = new AtomicReference<>();
            RetryPolicy stock = new StormBoundedExponentialBackoffRetry(100, 1000, MAX_RETRIES);
            RecordingPolicy rec = aware ? new RecordingPolicy(stock, ref::get, t0) : null;
            CuratorFramework zk = CuratorFrameworkFactory.builder().connectString(server.getConnectString())
                .sessionTimeoutMs(SESSION_MS).connectionTimeoutMs(CONN_MS).retryPolicy(aware ? rec : stock).build();
            ref.set(zk);
            if (aware) {
                rec.bind(zk);
            }
            final long[] reconnectedAt = {0};
            zk.getConnectionStateListenable().addListener((c, st) -> {
                log("state -> " + st);
                if (st == ConnectionState.RECONNECTED) {
                    reconnectedAt[0] = System.currentTimeMillis();
                }
            });
            zk.start();
            zk.blockUntilConnected(10, TimeUnit.SECONDS);
            zk.create().forPath("/r");

            AtomicInteger ok = new AtomicInteger();
            AtomicInteger failed = new AtomicInteger();
            AtomicInteger maxMs = new AtomicInteger();
            java.util.Map<String, Integer> failTypes = new java.util.concurrent.ConcurrentHashMap<>();
            java.util.concurrent.atomic.AtomicBoolean go = new java.util.concurrent.atomic.AtomicBoolean(true);
            List<Thread> threads = new ArrayList<>();
            for (int i = 0; i < 16; i++) {
                Thread t = new Thread(() -> {
                    while (go.get()) {
                        long a = System.currentTimeMillis();
                        try {
                            zk.getData().forPath("/r");
                            ok.incrementAndGet();
                        } catch (Exception e) {
                            failed.incrementAndGet();
                            failTypes.merge(e.getClass().getSimpleName(), 1, Integer::sum);
                        }
                        maxMs.accumulateAndGet((int) (System.currentTimeMillis() - a), Math::max);
                        java.util.concurrent.locks.LockSupport.parkNanos(1_000_000);
                    }
                }, "fg-" + i);
                t.start();
                threads.add(t);
            }
            Thread.sleep(500);
            log("stopping server");
            server.stop();
            Thread.sleep(2000);
            long restartedAt = System.currentTimeMillis();
            server.restart();
            log("server restarted");
            Thread.sleep(4000);
            go.set(false);
            for (Thread t : threads) {
                t.join(15_000);
            }
            out.add("synchronous calls: ok=" + ok + " failed=" + failed + " " + failTypes + " longestCallMs=" + maxMs);
            out.add("RECONNECTED notification " + (reconnectedAt[0] == 0 ? "NEVER" : (reconnectedAt[0] - restartedAt) + " ms after server restart"));
            if (aware) {
                out.add("policy.allowRetry calls: total=" + rec.total + " onEventThread=" + rec.onEventThread
                    + " onEventThreadWhileSuspendedOrLost=" + rec.onEventThreadWhileDown
                    + " blockedOver500msOnEventThread=" + rec.blockedOnEventThread + " maxCallMs=" + rec.maxMs);
            }
            zk.close();
        }
    }

    @Test
    public void repro() throws Exception {
        try {
            mechanism();
            natural(false);
            natural(true);
            foreground(false);
            foreground(true);
            experiment("long outage (8s > 5s connection timeout)", false, 1, 8000, 500);
            experiment("long outage (8s > 5s connection timeout)", true, 1, 8000, 500);
        } finally {
            try (PrintWriter w = new PrintWriter(OUT)) {
                synchronized (out) {
                    out.forEach(w::println);
                }
            }
        }
    }
}

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Supervisor process exits on transient ZooKeeper connection loss

4 participants