From f829e11606699a4f66a107e8ca7be8a0c784e197 Mon Sep 17 00:00:00 2001 From: TimurRakhmatullin Date: Fri, 25 Sep 2026 07:48:31 -0700 Subject: [PATCH] ZOOKEEPER-XXXX: Restore interrupt status in quorum catch blocks MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Multiple catch(InterruptedException) blocks in the quorum package swallow the interrupt flag without calling Thread.currentThread().interrupt(). This violates the Java interrupt contract: any code higher in the call stack that checks Thread.interrupted() or calls a blocking method will never see the interruption, which can cause threads to hang during shutdown or fail to terminate promptly. This patch adds Thread.currentThread().interrupt() as the first statement in 9 catch blocks across 6 files: - QuorumCnxManager.halt() — after listener.join() - QuorumCnxManager.ListenerHandler.acceptConnections() — after Thread.sleep() in retry loop - QuorumCnxManager.SendWorker.run() — after pollSendQueue() - QuorumPeerMain.runFromConfig() — after quorumPeer.join() - LearnerHandler.shutdown() — after queuedPackets.put() - Observer.waitForReconnectDelayHelper() — after Thread.sleep() - Learner.connectToLeader() — after latch.await() - Learner.connectToLeader() finally — after awaitTermination() - FastLeaderElection.WorkerReceiver.run() — after manager.pollRecvQueue() No behavioral change: each catch block still logs the same message at the same level. The only addition is the single Thread.currentThread().interrupt() call so the flag is preserved for upstream callers. Co-Authored-By: Claude Opus 4.6 --- .../org/apache/zookeeper/server/quorum/FastLeaderElection.java | 1 + .../main/java/org/apache/zookeeper/server/quorum/Learner.java | 2 ++ .../org/apache/zookeeper/server/quorum/LearnerHandler.java | 1 + .../main/java/org/apache/zookeeper/server/quorum/Observer.java | 1 + .../org/apache/zookeeper/server/quorum/QuorumCnxManager.java | 3 +++ .../org/apache/zookeeper/server/quorum/QuorumPeerMain.java | 1 + 6 files changed, 9 insertions(+) diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java index 61c0eb60104..573b7ac7112 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/FastLeaderElection.java @@ -467,6 +467,7 @@ public void run() { } } } catch (InterruptedException e) { + Thread.currentThread().interrupt(); LOG.warn("Interrupted Exception while waiting for new message", e); } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java index adf0ef6e510..335e641f315 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Learner.java @@ -332,6 +332,7 @@ protected void connectToLeader(MultipleAddresses multiAddr, String hostname) thr try { latch.await(); } catch (InterruptedException e) { + Thread.currentThread().interrupt(); LOG.warn("Interrupted while trying to connect to Leader", e); } finally { executor.shutdown(); @@ -340,6 +341,7 @@ protected void connectToLeader(MultipleAddresses multiAddr, String hostname) thr LOG.error("not all the LeaderConnector terminated properly"); } } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); LOG.error("Interrupted while terminating LeaderConnector executor.", ie); } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java index 57947daa86c..f75a9d4bef6 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/LearnerHandler.java @@ -1047,6 +1047,7 @@ public void shutdown() { queuedPackets.clear(); queuedPackets.put(proposalOfDeath); } catch (InterruptedException e) { + Thread.currentThread().interrupt(); LOG.warn("Ignoring unexpected exception", e); } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java index 334fa54c1bc..ef73fd5521e 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/Observer.java @@ -258,6 +258,7 @@ private static void waitForReconnectDelayHelper(long delayValueMs) { try { Thread.sleep(randomDelay); } catch (InterruptedException e) { + Thread.currentThread().interrupt(); LOG.warn("Interrupted while waiting", e); } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumCnxManager.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumCnxManager.java index 7a16014848e..226b648f1ba 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumCnxManager.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumCnxManager.java @@ -819,6 +819,7 @@ public void halt() { try { listener.join(); } catch (InterruptedException ex) { + Thread.currentThread().interrupt(); LOG.warn("Got interrupted before joining the listener", ex); } softHalt(); @@ -1106,6 +1107,7 @@ private void acceptConnections() { } catch (IOException ie) { LOG.error("Error closing server socket", ie); } catch (InterruptedException ie) { + Thread.currentThread().interrupt(); LOG.error("Interrupted while sleeping. Ignoring exception", ie); } closeSocket(client); @@ -1282,6 +1284,7 @@ public void run() { send(b); } } catch (InterruptedException e) { + Thread.currentThread().interrupt(); LOG.warn("Interrupted while waiting for message on queue", e); } } diff --git a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumPeerMain.java b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumPeerMain.java index 3fb3d6bf313..0a4751a6b4a 100644 --- a/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumPeerMain.java +++ b/zookeeper-server/src/main/java/org/apache/zookeeper/server/quorum/QuorumPeerMain.java @@ -231,6 +231,7 @@ public void runFromConfig(QuorumPeerConfig config) throws IOException, AdminServ ZKAuditProvider.addZKStartStopAuditLog(); quorumPeer.join(); } catch (InterruptedException e) { + Thread.currentThread().interrupt(); // warn, but generally this is ok LOG.warn("Quorum Peer interrupted", e); } finally {