Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -359,6 +359,8 @@ public void shutdown(int exitCode) {
jettyWebServer.stop();
}
if (sharedServiceLocator != null) {
// Stop after Jetty so no new connection can restart the heartbeat scheduler.
sharedServiceLocator.getService(NotebookServer.class).stopHeartbeatScheduler();
if (!zConf.isRecoveryEnabled()) {
sharedServiceLocator.getService(InterpreterSettingManager.class).close();
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -150,7 +150,6 @@ String getKey() {
// Package-private (not private) so NotebookServerHeartbeatTest can observe scheduler
// lifecycle without exposing it as part of the public API.
ScheduledExecutorService heartbeatScheduler;
private Thread heartbeatShutdownHook;
private boolean heartbeatInitialized;

// TODO(jl): This will be removed by handling session directly
Expand Down Expand Up @@ -287,29 +286,18 @@ synchronized void startHeartbeatScheduler() {
});
heartbeatScheduler.scheduleAtFixedRate(
this::sendHeartbeat, intervalMs, intervalMs, TimeUnit.MILLISECONDS);
heartbeatShutdownHook = new Thread(this::stopHeartbeatScheduler);
Runtime.getRuntime().addShutdownHook(heartbeatShutdownHook);
LOGGER.info("Started websocket heartbeat scheduler with interval {} ms", intervalMs);
}

/**
* Stops the websocket heartbeat scheduler, if running, and deregisters its shutdown hook so
* repeated start/stop cycles do not accumulate hooks. Safe to call multiple times and safe
* to call when the scheduler was never started.
* Stops the websocket heartbeat scheduler, if running. Safe to call multiple times and safe
* to call when the scheduler was never started; a later connection starts it again.
*/
synchronized void stopHeartbeatScheduler() {
public synchronized void stopHeartbeatScheduler() {
if (heartbeatScheduler != null) {
heartbeatScheduler.shutdownNow();
heartbeatScheduler = null;
}
if (heartbeatShutdownHook != null && Thread.currentThread() != heartbeatShutdownHook) {
try {
Runtime.getRuntime().removeShutdownHook(heartbeatShutdownHook);
} catch (IllegalStateException e) {
// JVM is already shutting down; the hook will simply run (as a harmless no-op).
}
heartbeatShutdownHook = null;
}
heartbeatInitialized = false;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,11 +19,14 @@
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
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.Mockito.doThrow;
import static org.mockito.Mockito.mock;
import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;

import java.util.concurrent.ScheduledExecutorService;
import org.apache.zeppelin.MiniZeppelinServer;
import org.apache.zeppelin.conf.ZeppelinConfiguration;
import org.apache.zeppelin.notebook.AuthorizationService;
import org.junit.jupiter.api.AfterEach;
Expand Down Expand Up @@ -106,4 +109,47 @@ void startHeartbeatSchedulerDoesNotStartWhenIntervalIsNegative() {

assertNull(server.heartbeatScheduler);
}

@Test
void stopHeartbeatSchedulerAllowsRepeatedStartStopCycles() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Tests look good!

Right now the tests all call stopHeartbeatScheduler directly,
so they'd still pass even if the ZeppelinServer#shutdown call got removed.
The actual fix is that shutdown drives the teardown, and that path isn't covered.

Might be worth adding one that goes through the real shutdown path. Spin up MiniZeppelinServer, open a connection to start the heartbeat, shut down and check the scheduler's gone.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

+1. The two new tests also pass against the previous shutdown-hook code, so they don't pin down what this PR changes. There's no need to open a connection: calling startHeartbeatScheduler() directly, then checking scheduler.isShutdown() after MiniZeppelinServer.shutDown(), is enough.

NotebookServer server = buildNotebookServer(50L);

for (int i = 0; i < 3; i++) {
server.startHeartbeatScheduler();
ScheduledExecutorService scheduler = server.heartbeatScheduler;
assertNotNull(scheduler);

server.stopHeartbeatScheduler();

assertTrue(scheduler.isShutdown());
assertNull(server.heartbeatScheduler);
}
}

@Test
void stopHeartbeatSchedulerIsSafeWhenNeverStarted() {
NotebookServer server = buildNotebookServer(50L);

assertDoesNotThrow(server::stopHeartbeatScheduler);
assertDoesNotThrow(server::stopHeartbeatScheduler);
}

@Test
void zeppelinServerShutdownStopsHeartbeatScheduler() throws Exception {
MiniZeppelinServer zepServer =
new MiniZeppelinServer(NotebookServerHeartbeatTest.class.getSimpleName());
try {
zepServer.start();
NotebookServer server = zepServer.getService(NotebookServer.class);
server.startHeartbeatScheduler();
ScheduledExecutorService scheduler = server.heartbeatScheduler;
assertNotNull(scheduler);

zepServer.shutDown();

assertTrue(scheduler.isShutdown());
} finally {
zepServer.destroy();
}
}
}
Loading