This is an automated email from the ASF dual-hosted git repository.

tbonelee pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/zeppelin.git


The following commit(s) were added to refs/heads/master by this push:
     new ed0d731775 [ZEPPELIN-6694] Reap dead websocket connections via pong 
tracking
ed0d731775 is described below

commit ed0d731775b915c6fac5eb0b564c6c12694c5596
Author: JangAyeon <[email protected]>
AuthorDate: Sun Sep 27 16:47:39 2026 +0900

    [ZEPPELIN-6694] Reap dead websocket connections via pong tracking
    
    ### What is this PR for?
    Since [ZEPPELIN-6092](https://issues.apache.org/jira/browse/ZEPPELIN-6092), 
heartbeat pings keep resetting Jetty's idle timer, so the idle timeout no 
longer detects dead clients. A dead session stays open until TCP retransmission 
gives up.
    
    This PR tracks pongs to detect dead sessions explicitly:
    
    * Register a per-session PongMessage handler. A pong resets the session's 
ping count.
    * In sendHeartbeat, close a session after N consecutive missed pongs 
(GOING_AWAY, 1001). It is removed from bookkeeping first, since a dead peer may 
never trigger onClose.
    * Skip sessions that are handling a message. Jetty reads pongs only after 
onMessage returns, so a long-running op would otherwise look like missed pongs.
    * Add zeppelin.websocket.heartbeat.max.missed.pongs (default 3, <= 0 
disables). With the default 60s interval, a dead client is reaped in ~3-4 
minutes.
    
    
    ### What type of PR is it?
    Improvement
    
    ### Todos
    * [x] Track pongs and close sessions that miss N in a row
    * [x] Skip sessions that are handling a message
    * [x] Add config, docs and zeppelin-site.xml.template
    * [x] Update docs and `zeppelin-site.xml.template`
    * [x] Add unit tests
    
    ### What is the Jira issue?
    [ZEPPELIN-6694](https://issues.apache.org/jira/browse/ZEPPELIN-6694)
    
    ### How should this be tested?
    * Unit tests added:
       * NotebookSocketTest: ping counting, reset on pong, failed ping still 
counted, close swallows exceptions
       * NotebookServerHeartbeatTest:
           * a session with too many missed pongs is closed and removed from 
ConnectionManager, while a healthy session keeps receiving pings
           * nothing is closed when reaping is disabled
           * a session that is handling a message is not reaped
           * onMessage marks the session as handling a message until it returns
           * onOpen(Session, EndpointConfig) registers a pong handler that 
resets the ping count
      * ZeppelinConfigurationTest: default and override values for the new 
config
       * Run: `./mvnw -pl zeppelin-server test 
-Dtest='NotebookSocketTest,NotebookServerHeartbeatTest,ZeppelinConfigurationTest'`
    
    
    ### Questions:
    * Does the license files need to update? No
    * Is there breaking changes for older versions?  No. Reaping is on by 
default, but only affects sessions that stop answering pings, and can be 
disabled by setting the new property to 0.
    * Does this needs documentation? Yes, included in 
`docs/setup/operation/configuration.md`
    
    
    Closes #5500 from JangAyeon/ZEPPELIN-6694.
    
    Signed-off-by: ChanHo Lee <[email protected]>
---
 conf/zeppelin-site.xml.template                    |   6 +
 docs/setup/operation/configuration.md              |   6 +
 .../zeppelin/conf/ZeppelinConfiguration.java       |   6 +
 .../org/apache/zeppelin/socket/NotebookServer.java |  42 ++++++-
 .../org/apache/zeppelin/socket/NotebookSocket.java |  55 +++++++++
 .../zeppelin/conf/ZeppelinConfigurationTest.java   |  13 +++
 .../socket/NotebookServerHeartbeatTest.java        | 128 +++++++++++++++++++++
 .../apache/zeppelin/socket/NotebookSocketTest.java |  42 +++++++
 8 files changed, 295 insertions(+), 3 deletions(-)

diff --git a/conf/zeppelin-site.xml.template b/conf/zeppelin-site.xml.template
index 4107a6799c..3f1f7f8430 100755
--- a/conf/zeppelin-site.xml.template
+++ b/conf/zeppelin-site.xml.template
@@ -571,6 +571,12 @@
   <description>Interval in milliseconds at which the server sends a websocket 
ping frame to each session to keep it alive. Defaults to 60000 (1 minute). Set 
to 0 or a negative value to disable server-initiated heartbeats.</description>
 </property>
 
+<property>
+  <name>zeppelin.websocket.heartbeat.max.missed.pongs</name>
+  <value>3</value>
+  <description>Number of consecutive heartbeat pings left unanswered (no pong) 
before the server closes the websocket session as dead. Defaults to 3. Set to 0 
or a negative value to disable reaping.</description>
+</property>
+
 <property>
   <name>zeppelin.server.default.dir.allowed</name>
   <value>false</value>
diff --git a/docs/setup/operation/configuration.md 
b/docs/setup/operation/configuration.md
index 4215222c40..a6debaaa12 100644
--- a/docs/setup/operation/configuration.md
+++ b/docs/setup/operation/configuration.md
@@ -418,6 +418,12 @@ Sources descending by priority:
     <td>60000</td>
     <td>Interval(in milliseconds) at which the server sends a websocket ping 
frame to each session to keep it alive. Set to 0 or a negative value to disable 
server-initiated heartbeats.</td>
   </tr>
+  <tr>
+    <td><h6 
class="properties">ZEPPELIN_WEBSOCKET_HEARTBEAT_MAX_MISSED_PONGS</h6></td>
+    <td><h6 
class="properties">zeppelin.websocket.heartbeat.max.missed.pongs</h6></td>
+    <td>3</td>
+    <td>Number of consecutive heartbeat pings left unanswered (no pong) before 
the server closes the websocket session as dead. Set to 0 or a negative value 
to disable reaping. Has no effect when heartbeats are disabled.</td>
+  </tr>
   <tr>
     <td><h6 class="properties">ZEPPELIN_SERVER_DEFAULT_DIR_ALLOWED</h6></td>
     <td><h6 class="properties">zeppelin.server.default.dir.allowed</h6></td>
diff --git 
a/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java
 
b/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java
index 01ad388eba..6a7e1b7909 100644
--- 
a/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java
+++ 
b/zeppelin-server/src/main/java/org/apache/zeppelin/conf/ZeppelinConfiguration.java
@@ -743,6 +743,10 @@ public class ZeppelinConfiguration {
     return getLong(ConfVars.ZEPPELIN_WEBSOCKET_HEARTBEAT_INTERVAL);
   }
 
+  public int getWebsocketHeartbeatMaxMissedPongs() {
+    return getInt(ConfVars.ZEPPELIN_WEBSOCKET_HEARTBEAT_MAX_MISSED_PONGS);
+  }
+
   public String getJettyName() {
     return getString(ConfVars.ZEPPELIN_SERVER_JETTY_NAME);
   }
@@ -1105,6 +1109,8 @@ public class ZeppelinConfiguration {
     // per-connection traffic low. 60s gives 5 pings within the 300s default 
idle window.
     // <= 0 disables server-initiated heartbeats.
     
ZEPPELIN_WEBSOCKET_HEARTBEAT_INTERVAL("zeppelin.websocket.heartbeat.interval", 
60000L),
+    // Consecutive unanswered pings before a session is closed as dead. <= 0 
disables reaping.
+    
ZEPPELIN_WEBSOCKET_HEARTBEAT_MAX_MISSED_PONGS("zeppelin.websocket.heartbeat.max.missed.pongs",
 3),
     
ZEPPELIN_WEBSOCKET_PARAGRAPH_STATUS_PROGRESS("zeppelin.websocket.paragraph_status_progress.enable",
 true),
     ZEPPELIN_SERVER_DEFAULT_DIR_ALLOWED("zeppelin.server.default.dir.allowed", 
false),
     ZEPPELIN_SERVER_XFRAME_OPTIONS("zeppelin.server.xframe.options", 
"SAMEORIGIN"),
diff --git 
a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java 
b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java
index d55eb27196..ba5bffc76e 100644
--- 
a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java
+++ 
b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookServer.java
@@ -50,6 +50,7 @@ import jakarta.websocket.OnClose;
 import jakarta.websocket.OnError;
 import jakarta.websocket.OnMessage;
 import jakarta.websocket.OnOpen;
+import jakarta.websocket.PongMessage;
 import jakarta.websocket.Session;
 import jakarta.websocket.server.ServerEndpoint;
 import org.apache.commons.lang3.StringUtils;
@@ -250,6 +251,7 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
     if (checkOrigin(origin)) {
       NotebookSocket notebookSocket = sessionIdNotebookSocketMap
           .computeIfAbsent(session.getId(), unused -> new 
NotebookSocket(session, headers));
+      session.addMessageHandler(PongMessage.class, pong -> 
notebookSocket.onPong());
       onOpen(notebookSocket);
     } else {
       LOGGER.error("Websocket request is not allowed by {} settings. Origin: 
{}", ZEPPELIN_ALLOWED_ORIGINS,
@@ -307,12 +309,22 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
    * point of this heartbeat: it keeps connections alive even when the 
client-side application
    * keep-alive timer is throttled or stopped (e.g. a backgrounded browser 
tab). A single
    * session failing to receive a ping must not stop the remaining sessions 
from being pinged.
-   * Pong responses are not tracked; once the heartbeat is enabled, Jetty's 
idle timeout no
-   * longer determines connection liveness (see ZEPPELIN-6694).
+   *
+   * <p>Because these writes keep resetting Jetty's idle timer, the idle 
timeout can no longer
+   * detect dead clients. Liveness is therefore tracked explicitly 
(ZEPPELIN-6694): a session
+   * that has left {@code zeppelin.websocket.heartbeat.max.missed.pongs} 
consecutive pings
+   * unanswered is closed instead of being pinged again. A session that is 
still handling a
+   * message is never closed here, because Jetty reads its pongs only after 
onMessage returns.
    */
   void sendHeartbeat() {
+    int maxMissedPongs = zConf.getWebsocketHeartbeatMaxMissedPongs();
     for (NotebookSocket conn : connectionManager.connectedSockets) {
       try {
+        if (maxMissedPongs > 0 && !conn.isHandlingMessage()
+            && conn.getPingsSinceLastPong() >= maxMissedPongs) {
+          reapDeadConnection(conn, maxMissedPongs);
+          continue;
+        }
         conn.sendPing();
       } catch (RuntimeException e) {
         LOGGER.warn("Failed to send heartbeat ping to {}", conn, e);
@@ -320,10 +332,34 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
     }
   }
 
+  /**
+   * Removes a connection that stopped answering pings and closes its session. 
The connection is
+   * dropped from all bookkeeping before closing, because a dead peer never 
completes the close
+   * handshake and {@link #onClose} may therefore arrive late or not at all. 
Removing it from
+   * {@code sessionIdNotebookSocketMap} first also makes a later {@code 
onClose} a no-op.
+   */
+  private void reapDeadConnection(NotebookSocket conn, int maxMissedPongs) {
+    LOGGER.warn("Closing websocket to {}: {} consecutive heartbeat pings 
unanswered, last pong at {}",
+        conn, maxMissedPongs, new Date(conn.getLastPongTimestamp()));
+    sessionIdNotebookSocketMap.remove(conn.getSessionId());
+    removeConnection(conn);
+    conn.close(new CloseReason(CloseReason.CloseCodes.GOING_AWAY,
+        "No pong received for " + maxMissedPongs + " consecutive pings"));
+  }
+
   @OnMessage
   public void onMessage(Session session, String msg) {
     NotebookSocket conn = sessionIdNotebookSocketMap.get(session.getId());
-    onMessage(conn, msg);
+    if (conn == null) {
+      onMessage(conn, msg);
+      return;
+    }
+    conn.setHandlingMessage(true);
+    try {
+      onMessage(conn, msg);
+    } finally {
+      conn.setHandlingMessage(false);
+    }
   }
 
   public void onMessage(NotebookSocket conn, String msg) {
diff --git 
a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookSocket.java 
b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookSocket.java
index 57edf1d79b..db70d4fc2c 100644
--- 
a/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookSocket.java
+++ 
b/zeppelin-server/src/main/java/org/apache/zeppelin/socket/NotebookSocket.java
@@ -24,7 +24,9 @@ import org.slf4j.LoggerFactory;
 import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.Map;
+import java.util.concurrent.atomic.AtomicInteger;
 
+import jakarta.websocket.CloseReason;
 import jakarta.websocket.Session;
 
 /**
@@ -42,13 +44,26 @@ public class NotebookSocket {
   private Map<String, Object> headers;
   private String user;
 
+  // Liveness tracking (ZEPPELIN-6694). Written from the heartbeat thread 
(sendPing) and from
+  // the websocket container thread (onPong), hence atomic/volatile.
+  private final AtomicInteger unansweredPings = new AtomicInteger();
+  private volatile long lastPongTimestamp;
+  // True while onMessage is running for this session. Jetty reads the next 
frame (pongs
+  // included) only after onMessage returns, so pongs pile up unread during a 
long operation.
+  private volatile boolean handlingMessage;
+
   public NotebookSocket(Session session, Map<String, Object> headers) {
     this.session = session;
     this.headers = headers;
     this.user = StringUtils.EMPTY;
+    this.lastPongTimestamp = System.currentTimeMillis();
     LOGGER.debug("NotebookSocket created for session: {}", session.getId());
   }
 
+  public String getSessionId() {
+    return session.getId();
+  }
+
   public String getHeader(String key) {
     return String.valueOf(headers.get(key));
   }
@@ -67,8 +82,11 @@ public class NotebookSocket {
    * session resets Jetty's idle timeout as well as any intermediate proxy's 
idle timer, so no
    * application-level handling is required on the client. Exceptions are 
swallowed and logged
    * so a single dead session cannot break the caller's heartbeat loop over 
all sessions.
+   * Every call counts as one outstanding ping until {@link #onPong()} is 
called, including
+   * calls whose write failed, since a failed write is itself a sign the peer 
is gone.
    */
   public void sendPing() {
+    unansweredPings.incrementAndGet();
     try {
       session.getBasicRemote().sendPing(PING_PAYLOAD);
     } catch (IOException | IllegalArgumentException | IllegalStateException e) 
{
@@ -76,6 +94,43 @@ public class NotebookSocket {
     }
   }
 
+  /**
+   * Records a pong frame from the peer. Any pong proves the connection is 
alive, so the
+   * outstanding-ping counter is reset rather than decremented.
+   */
+  public void onPong() {
+    lastPongTimestamp = System.currentTimeMillis();
+    unansweredPings.set(0);
+  }
+
+  public int getPingsSinceLastPong() {
+    return unansweredPings.get();
+  }
+
+  public long getLastPongTimestamp() {
+    return lastPongTimestamp;
+  }
+
+  public boolean isHandlingMessage() {
+    return handlingMessage;
+  }
+
+  public void setHandlingMessage(boolean handlingMessage) {
+    this.handlingMessage = handlingMessage;
+  }
+
+  /**
+   * Closes the underlying session. Exceptions are swallowed and logged 
because this is used to
+   * reap connections that are already presumed dead.
+   */
+  public void close(CloseReason closeReason) {
+    try {
+      session.close(closeReason);
+    } catch (IOException | IllegalStateException e) {
+      LOGGER.debug("Failed to close session {}: {}", session.getId(), 
e.toString());
+    }
+  }
+
   public String getUser() {
     return user;
   }
diff --git 
a/zeppelin-server/src/test/java/org/apache/zeppelin/conf/ZeppelinConfigurationTest.java
 
b/zeppelin-server/src/test/java/org/apache/zeppelin/conf/ZeppelinConfigurationTest.java
index f1e4d0d4da..e230152b36 100644
--- 
a/zeppelin-server/src/test/java/org/apache/zeppelin/conf/ZeppelinConfigurationTest.java
+++ 
b/zeppelin-server/src/test/java/org/apache/zeppelin/conf/ZeppelinConfigurationTest.java
@@ -185,4 +185,17 @@ class ZeppelinConfigurationTest {
     
zConf.setProperty(ConfVars.ZEPPELIN_WEBSOCKET_HEARTBEAT_INTERVAL.getVarName(), 
"0");
     assertEquals(0L, zConf.getWebsocketHeartbeatInterval());
   }
+
+  @Test
+  void getWebsocketHeartbeatMaxMissedPongsDefaultTest() {
+    ZeppelinConfiguration zConf = 
ZeppelinConfiguration.load("zeppelin-test-site.xml");
+    assertEquals(3, zConf.getWebsocketHeartbeatMaxMissedPongs());
+  }
+
+  @Test
+  void getWebsocketHeartbeatMaxMissedPongsOverrideTest() {
+    ZeppelinConfiguration zConf = 
ZeppelinConfiguration.load("zeppelin-test-site.xml");
+    
zConf.setProperty(ConfVars.ZEPPELIN_WEBSOCKET_HEARTBEAT_MAX_MISSED_PONGS.getVarName(),
 "5");
+    assertEquals(5, zConf.getWebsocketHeartbeatMaxMissedPongs());
+  }
 }
diff --git 
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java
 
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java
index f932b1fa2d..fd7bac091d 100644
--- 
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java
+++ 
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerHeartbeatTest.java
@@ -17,20 +17,42 @@
 package org.apache.zeppelin.socket;
 
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
 import static org.junit.jupiter.api.Assertions.assertNull;
 import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.eq;
+import static org.mockito.Mockito.doAnswer;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.spy;
 import static org.mockito.Mockito.verify;
 import static org.mockito.Mockito.when;
 
+import java.io.IOException;
+import java.util.HashMap;
+import java.util.Map;
 import java.util.concurrent.ScheduledExecutorService;
+import java.util.concurrent.atomic.AtomicBoolean;
+
+import jakarta.websocket.CloseReason;
+import jakarta.websocket.EndpointConfig;
+import jakarta.websocket.MessageHandler;
+import jakarta.websocket.PongMessage;
+import jakarta.websocket.RemoteEndpoint;
+import jakarta.websocket.Session;
+
 import org.apache.zeppelin.MiniZeppelinServer;
 import org.apache.zeppelin.conf.ZeppelinConfiguration;
 import org.apache.zeppelin.notebook.AuthorizationService;
+import org.apache.zeppelin.utils.CorsUtils;
 import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
 
 class NotebookServerHeartbeatTest {
 
@@ -44,8 +66,13 @@ class NotebookServerHeartbeatTest {
   }
 
   private NotebookServer buildNotebookServer(long heartbeatIntervalMs) {
+    return buildNotebookServer(heartbeatIntervalMs, 0);
+  }
+
+  private NotebookServer buildNotebookServer(long heartbeatIntervalMs, int 
maxMissedPongs) {
     ZeppelinConfiguration zConf = mock(ZeppelinConfiguration.class);
     
when(zConf.getWebsocketHeartbeatInterval()).thenReturn(heartbeatIntervalMs);
+    
when(zConf.getWebsocketHeartbeatMaxMissedPongs()).thenReturn(maxMissedPongs);
     AuthorizationService authorizationService = 
mock(AuthorizationService.class);
     ConnectionManager connectionManager = new 
ConnectionManager(authorizationService, zConf);
 
@@ -83,6 +110,107 @@ class NotebookServerHeartbeatTest {
     verify(healthy).sendPing();
   }
 
+  @Test
+  void sendHeartbeatReapsSocketThatMissedTooManyPongs() {
+    NotebookServer server = buildNotebookServer(60000L, 3);
+    NotebookSocket dead = mock(NotebookSocket.class);
+    NotebookSocket alive = mock(NotebookSocket.class);
+    when(dead.getSessionId()).thenReturn("dead-session");
+    when(dead.getPingsSinceLastPong()).thenReturn(3);
+    when(alive.getPingsSinceLastPong()).thenReturn(2);
+    server.getConnectionManager().addConnection(dead);
+    server.getConnectionManager().addConnection(alive);
+
+    server.sendHeartbeat();
+
+    verify(dead).close(any(CloseReason.class));
+    verify(dead, never()).sendPing();
+    verify(alive).sendPing();
+    verify(alive, never()).close(any(CloseReason.class));
+    assertFalse(server.getConnectionManager().connectedSockets.contains(dead));
+    assertTrue(server.getConnectionManager().connectedSockets.contains(alive));
+  }
+
+  @Test
+  void sendHeartbeatDoesNotReapWhenReapingDisabled() {
+    NotebookServer server = buildNotebookServer(60000L, 0);
+    NotebookSocket silent = mock(NotebookSocket.class);
+    when(silent.getPingsSinceLastPong()).thenReturn(100);
+    server.getConnectionManager().addConnection(silent);
+
+    server.sendHeartbeat();
+
+    verify(silent).sendPing();
+    verify(silent, never()).close(any(CloseReason.class));
+  }
+
+  @Test
+  void sendHeartbeatDoesNotReapSocketThatIsHandlingMessage() {
+    NotebookServer server = buildNotebookServer(60000L, 3);
+    NotebookSocket busy = mock(NotebookSocket.class);
+    when(busy.isHandlingMessage()).thenReturn(true);
+    when(busy.getPingsSinceLastPong()).thenReturn(5);
+    server.getConnectionManager().addConnection(busy);
+
+    server.sendHeartbeat();
+
+    verify(busy, never()).close(any(CloseReason.class));
+    verify(busy).sendPing();
+    assertTrue(server.getConnectionManager().connectedSockets.contains(busy));
+  }
+
+  private Session openSession(NotebookServer server, String sessionId) throws 
IOException {
+    Session session = mock(Session.class);
+    when(session.getId()).thenReturn(sessionId);
+    
when(session.getBasicRemote()).thenReturn(mock(RemoteEndpoint.Basic.class));
+    Map<String, Object> headers = new HashMap<>();
+    headers.put(CorsUtils.HEADER_ORIGIN, "http://localhost:8080";);
+    EndpointConfig config = mock(EndpointConfig.class);
+    when(config.getUserProperties()).thenReturn(headers);
+    server.onOpen(session, config);
+    return session;
+  }
+
+  @Test
+  @SuppressWarnings("unchecked")
+  void onOpenRegistersPongHandlerThatResetsPingCount() throws IOException {
+    NotebookServer server = buildNotebookServer(60000L, 3);
+    Session session = openSession(server, "session-pong");
+    NotebookSocket conn = 
server.getConnectionManager().connectedSockets.peek();
+    assertNotNull(conn);
+
+    ArgumentCaptor<MessageHandler.Whole<PongMessage>> captor =
+        ArgumentCaptor.forClass(MessageHandler.Whole.class);
+    verify(session).addMessageHandler(eq(PongMessage.class), captor.capture());
+
+    conn.sendPing();
+    conn.sendPing();
+    assertEquals(2, conn.getPingsSinceLastPong());
+
+    captor.getValue().onMessage(mock(PongMessage.class));
+    assertEquals(0, conn.getPingsSinceLastPong());
+  }
+
+  @Test
+  void onMessageMarksSocketAsHandlingMessageUntilItReturns() throws 
IOException {
+    NotebookServer server = spy(buildNotebookServer(60000L, 3));
+    notebookServer = server;
+    Session session = openSession(server, "session-busy");
+    NotebookSocket conn = 
server.getConnectionManager().connectedSockets.peek();
+    assertNotNull(conn);
+
+    AtomicBoolean handlingDuringOp = new AtomicBoolean();
+    doAnswer(invocation -> {
+      handlingDuringOp.set(conn.isHandlingMessage());
+      return null;
+    }).when(server).onMessage(any(NotebookSocket.class), anyString());
+
+    server.onMessage(session, "{}");
+
+    assertTrue(handlingDuringOp.get());
+    assertFalse(conn.isHandlingMessage());
+  }
+
   @Test
   void startHeartbeatSchedulerStartsWhenIntervalPositive() {
     NotebookServer server = buildNotebookServer(50L);
diff --git 
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookSocketTest.java
 
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookSocketTest.java
index 4382e0dafe..53447b5d16 100644
--- 
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookSocketTest.java
+++ 
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookSocketTest.java
@@ -17,6 +17,7 @@
 package org.apache.zeppelin.socket;
 
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.mockito.ArgumentMatchers.any;
 import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
@@ -27,6 +28,7 @@ import java.io.IOException;
 import java.nio.ByteBuffer;
 import java.util.Collections;
 
+import jakarta.websocket.CloseReason;
 import jakarta.websocket.RemoteEndpoint;
 import jakarta.websocket.Session;
 
@@ -59,4 +61,44 @@ class NotebookSocketTest {
 
     assertDoesNotThrow(notebookSocket::sendPing);
   }
+
+  @Test
+  void sendPingCountsUnansweredPingsAndPongResetsCount() {
+    Session session = mock(Session.class);
+    when(session.getId()).thenReturn("session-3");
+    
when(session.getBasicRemote()).thenReturn(mock(RemoteEndpoint.Basic.class));
+    NotebookSocket notebookSocket = new NotebookSocket(session, 
Collections.emptyMap());
+
+    notebookSocket.sendPing();
+    notebookSocket.sendPing();
+    assertEquals(2, notebookSocket.getPingsSinceLastPong());
+
+    notebookSocket.onPong();
+    assertEquals(0, notebookSocket.getPingsSinceLastPong());
+  }
+
+  @Test
+  void failedPingStillCountsAsUnanswered() throws IOException {
+    Session session = mock(Session.class);
+    RemoteEndpoint.Basic basicRemote = mock(RemoteEndpoint.Basic.class);
+    when(session.getId()).thenReturn("session-4");
+    when(session.getBasicRemote()).thenReturn(basicRemote);
+    doThrow(new IOException("broken 
pipe")).when(basicRemote).sendPing(any(ByteBuffer.class));
+    NotebookSocket notebookSocket = new NotebookSocket(session, 
Collections.emptyMap());
+
+    notebookSocket.sendPing();
+
+    assertEquals(1, notebookSocket.getPingsSinceLastPong());
+  }
+
+  @Test
+  void closeSwallowsIOException() throws IOException {
+    Session session = mock(Session.class);
+    when(session.getId()).thenReturn("session-5");
+    doThrow(new IOException("already 
closed")).when(session).close(any(CloseReason.class));
+    NotebookSocket notebookSocket = new NotebookSocket(session, 
Collections.emptyMap());
+
+    assertDoesNotThrow(() -> notebookSocket.close(
+        new CloseReason(CloseReason.CloseCodes.GOING_AWAY, "test")));
+  }
 }

Reply via email to