diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java new file mode 100644 index 00000000000..7d311aa35b4 --- /dev/null +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterProcessStatus.java @@ -0,0 +1,111 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.zeppelin.interpreter; + +import org.apache.zeppelin.interpreter.remote.RemoteInterpreterProcess; + +/** + * Point-in-time status snapshot of a single interpreter process as seen by the Zeppelin server. + * Built purely from in-memory server state without contacting the process, so {@code started} + * reflects whether a process handle exists, not whether the process is currently reachable. + * Reachability is intentionally out of scope here to keep the read path non-blocking. + * + *

Every value below must be readable without leaving the JVM. In particular do not call + * {@code isRunning()}, {@code isAlive()} or {@code getErrorMessage()} on the process from here: + * those are failure-path diagnostics that contact the container runtime on some launchers, so + * calling them would let a slow or unreachable runtime block this endpoint. + * + *

{@code started} and {@code launching} are read one after the other rather than under the + * lock that guards a launch, so this is a best-effort view of a group that is starting up: the + * pair can straddle the moment a launch finishes. What it does buy is that the window in which + * a handle carries no {@code host} or {@code port} yet is reported as such instead of looking + * like a fully started process. + */ +public class InterpreterProcessStatus { + private final String settingId; + private final String settingName; + private final String groupId; + private final int numSessions; + private final boolean launching; + private final boolean started; + private String host; + private int port = -1; + private String startTime; + private long attachedForSeconds; + + public InterpreterProcessStatus(ManagedInterpreterGroup group) { + InterpreterSetting setting = group.getInterpreterSetting(); + this.settingId = setting.getId(); + this.settingName = setting.getName(); + this.groupId = group.getId(); + this.numSessions = group.getSessionNum(); + this.launching = group.isLaunchingInterpreterProcess(); + RemoteInterpreterProcess process = group.getInterpreterProcess(); + this.started = process != null; + if (started) { + this.host = process.getHost(); + this.port = process.getPort(); + this.startTime = process.getStartTime(); + this.attachedForSeconds = (System.currentTimeMillis() - process.getStartTimeMs()) / 1000; + } + } + + public String getSettingId() { + return settingId; + } + + public String getSettingName() { + return settingName; + } + + public String getGroupId() { + return groupId; + } + + public int getNumSessions() { + return numSessions; + } + + /** + * @return whether a process is currently being launched for this group, in which case + * {@code host} and {@code port} may not be filled in yet even when {@code started} + */ + public boolean isLaunching() { + return launching; + } + + public boolean isStarted() { + return started; + } + + public String getHost() { + return host; + } + + public int getPort() { + return port; + } + + public String getStartTime() { + return startTime; + } + + public long getAttachedForSeconds() { + return attachedForSeconds; + } +} diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java index 8b4c2fe55be..d0cff37f2b8 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSettingManager.java @@ -711,6 +711,17 @@ public List getAllInterpreterGroup() { return interpreterGroups; } + /** + * Snapshot the status of every running interpreter group. Uses in-memory state only + */ + public List getInterpreterProcessStatuses() { + List statuses = new ArrayList<>(); + for (ManagedInterpreterGroup group : getAllInterpreterGroup()) { + statuses.add(new InterpreterProcessStatus(group)); + } + return statuses; + } + // TODO(zjffdu) Current approach is not optimized. we have to iterate all interpreter settings. public void removeInterpreterGroup(String intpGroupId) { for (InterpreterSetting interpreterSetting : interpreterSettings.values()) { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java index 3a2f78af895..3a8b14ee81e 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/ManagedInterpreterGroup.java @@ -41,7 +41,7 @@ public class ManagedInterpreterGroup extends InterpreterGroup { private static final Logger LOGGER = LoggerFactory.getLogger(ManagedInterpreterGroup.class); private InterpreterSetting interpreterSetting; - private RemoteInterpreterProcess remoteInterpreterProcess; // attached remote interpreter process + private volatile RemoteInterpreterProcess remoteInterpreterProcess; private Object interpreterProcessCreationLock = new Object(); private final ZeppelinConfiguration zConf; private volatile long lastUsedTimeInMillis = System.currentTimeMillis(); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterManagedProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterManagedProcess.java index 02cedb322fe..d9c33e40055 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterManagedProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterManagedProcess.java @@ -32,14 +32,14 @@ public abstract class RemoteInterpreterManagedProcess extends RemoteInterpreterP private final String interpreterPortRange; - private String host = null; - private int port = -1; + private volatile String host = null; + private volatile int port = -1; private final String interpreterDir; private final String localRepoDir; private final String interpreterSettingName; private final String interpreterGroupId; private final boolean isUserImpersonated; - private String errorMessage; + private volatile String errorMessage; private Map env; diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java index e994439c890..89c27a418c0 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcess.java @@ -45,7 +45,8 @@ public abstract class RemoteInterpreterProcess implements InterpreterClient, Aut protected String intpEventServerHost; protected int intpEventServerPort; private PooledRemoteClient remoteClient; - private String startTime; + private final long startTimeMs; + private final String startTime; public RemoteInterpreterProcess(int connectTimeout, int connectionPoolSize, @@ -54,7 +55,8 @@ public RemoteInterpreterProcess(int connectTimeout, this.connectTimeout = connectTimeout; this.intpEventServerHost = intpEventServerHost; this.intpEventServerPort = intpEventServerPort; - this.startTime = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date()); + this.startTimeMs = System.currentTimeMillis(); + this.startTime = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss").format(new Date(startTimeMs)); this.remoteClient = new PooledRemoteClient<>(() -> { TSocket transport = new TSocket(getHost(), getPort()); try { @@ -71,9 +73,21 @@ public int getConnectTimeout() { return connectTimeout; } + /** + * When the server created this object, formatted for display. This is not necessarily when the + * interpreter itself started: {@link RemoteInterpreterRunningProcess} is constructed fresh when + * the server recovers a process that outlived it, and when it attaches to an interpreter that + * was already running, so on those paths the stamp is the moment of attachment. + * + * @return the creation instant of this object as {@code yyyy-MM-dd HH:mm:ss} + */ public String getStartTime() { return startTime; } + + public long getStartTimeMs() { + return startTimeMs; + } @Override public void close() { diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java index 3b9d754e919..163fad96145 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/rest/InterpreterRestApi.java @@ -102,6 +102,17 @@ public Response listSettings() { return new JsonResponse<>(Status.OK, "", interpreterSettingManager.get()).build(); } + /** + * List the runtime status of all running interpreter processes. + */ + @GET + @Path("status") + @ZeppelinApi + public Response getInterpreterProcessStatus() { + return new JsonResponse<>(Status.OK, "", + interpreterSettingManager.getInterpreterProcessStatuses()).build(); + } + /** * Get a setting. */ diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java index 95e126fedff..ab489c278cd 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/InterpreterSettingManagerTest.java @@ -40,7 +40,9 @@ import java.util.Map; 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.junit.jupiter.api.Assertions.fail; import static org.mockito.Mockito.mock; @@ -262,6 +264,29 @@ void testRestartShared() throws InterpreterException { assertEquals(0, interpreterSetting.getAllInterpreterGroups().size()); } + @Test + void testGetInterpreterProcessStatuses() throws InterpreterException { + // no interpreter group has been created yet + assertTrue(interpreterSettingManager.getInterpreterProcessStatuses().isEmpty()); + + InterpreterSetting interpreterSetting = interpreterSettingManager.getByName("test"); + interpreterSetting.getOption().setPerUser("shared"); + interpreterSetting.getOption().setPerNote("shared"); + interpreterSetting.getOrCreateSession("user1", note1Id); + + List statuses = + interpreterSettingManager.getInterpreterProcessStatuses(); + assertEquals(1, statuses.size()); + InterpreterProcessStatus status = statuses.get(0); + assertEquals("test", status.getSettingName()); + assertEquals(1, status.getNumSessions()); + // process starts lazily on first interpret, so it is not started at this point + assertFalse(status.isStarted()); + assertFalse(status.isLaunching()); + assertNull(status.getHost()); + assertEquals(-1, status.getPort()); + } + @Test void testRestartPerUserIsolated() throws InterpreterException { InterpreterSetting interpreterSetting = interpreterSettingManager.getByName("test"); diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java index 19435b32112..66d266b703d 100644 --- a/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/rest/InterpreterRestApiTest.java @@ -106,6 +106,17 @@ void getSettings() throws IOException { get.close(); } + @Test + void testGetInterpreterProcessStatus() throws IOException { + // when + CloseableHttpResponse get = httpGet("/interpreter/status"); + // then + assertThat(get, isAllowed()); + JsonArray body = getArrayBodyFieldFromResponse(EntityUtils.toString(get.getEntity(), StandardCharsets.UTF_8)); + assertNotNull(body); + get.close(); + } + @Test void testGetNonExistInterpreterSetting() throws IOException { // when