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")));
+ }
}