reiabreu commented on PR #9153:
URL: https://github.com/apache/storm/pull/9153#issuecomment-5973530933

   *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.
   
   ```java
   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.
   
   <details><summary>Patch (on top of ab67921da)</summary>
   
   ```diff
   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());
        }
    }
   ```
   </details>
   
   <details><summary>Failover test program (throwaway JUnit test, not for 
merging)</summary>
   
   ```java
   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);
                   }
               }
           }
       }
   }
   ```
   </details>
   


-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to