From df61eb9d1c0c726601ce660af43236c216ee8adc Mon Sep 17 00:00:00 2001 From: dev-donghwan Date: Wed, 30 Sep 2026 17:36:37 +0900 Subject: [PATCH] [ZEPPELIN-6723] Close an interpreter group on unregister only when the sender is its registered process --- .../remote/RemoteInterpreterEventClient.java | 8 +- .../remote/RemoteInterpreterServer.java | 15 +- .../interpreter/thrift/AngularObjectId.java | 2 +- .../thrift/AppOutputAppendEvent.java | 2 +- .../thrift/AppOutputUpdateEvent.java | 2 +- .../thrift/AppStatusUpdateEvent.java | 2 +- .../thrift/InterpreterCompletion.java | 2 +- .../thrift/InterpreterRPCException.java | 2 +- .../interpreter/thrift/LibraryMetadata.java | 2 +- .../interpreter/thrift/OutputAppendEvent.java | 2 +- .../thrift/OutputUpdateAllEvent.java | 2 +- .../interpreter/thrift/OutputUpdateEvent.java | 2 +- .../interpreter/thrift/ParagraphInfo.java | 2 +- .../interpreter/thrift/RegisterInfo.java | 2 +- .../thrift/RemoteApplicationResult.java | 2 +- .../thrift/RemoteInterpreterContext.java | 2 +- .../thrift/RemoteInterpreterEvent.java | 2 +- .../thrift/RemoteInterpreterEventService.java | 144 ++++++++++++++-- .../thrift/RemoteInterpreterEventType.java | 2 +- .../thrift/RemoteInterpreterResult.java | 2 +- .../RemoteInterpreterResultMessage.java | 2 +- .../thrift/RemoteInterpreterService.java | 2 +- .../thrift/RunParagraphsEvent.java | 2 +- .../interpreter/thrift/ServiceException.java | 2 +- .../interpreter/thrift/WebUrlInfo.java | 2 +- .../RemoteInterpreterEventService.thrift | 2 +- .../remote/RemoteInterpreterServerTest.java | 20 +++ .../interpreter/InterpreterSetting.java | 16 ++ .../InterpreterSettingManager.java | 7 - .../interpreter/ManagedInterpreterGroup.java | 11 ++ .../RemoteInterpreterEventServer.java | 37 +++- ...nterpreterEventServerRegistrationTest.java | 163 ++++++++++++++++++ 32 files changed, 414 insertions(+), 53 deletions(-) create mode 100644 zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerRegistrationTest.java diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java index df281c84d44..75b66bb8e8b 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterEventClient.java @@ -91,9 +91,13 @@ public void registerInterpreterProcess(RegisterInfo registerInfo) { }); } - public void unRegisterInterpreterProcess() { + /** + * @param registerInfo what this process registered with, so that the server can tell whether + * it is the process of the interpreter group; null if it does not register + */ + public void unRegisterInterpreterProcess(RegisterInfo registerInfo) { callRemoteFunction(client -> { - client.unRegisterInterpreterProcess(intpGroupId); + client.unRegisterInterpreterProcess(intpGroupId, registerInfo); return null; }); } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServer.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServer.java index 346a8641ea1..ca20b1db3f6 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServer.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServer.java @@ -228,7 +228,7 @@ public void shutdown() throws InterpreterRPCException, TException { if (intpEventClient != null) { try { LOGGER.info("Unregister interpreter process"); - intpEventClient.unRegisterInterpreterProcess(); + intpEventClient.unRegisterInterpreterProcess(getRegisterInfo()); } catch (Exception e) { LOGGER.error("Fail to unregister remote interpreter process", e); } @@ -250,6 +250,15 @@ public Properties getProperties() { return this.zProperties; } + /** + * What this process registers with. host and port are fixed in the constructor, so an + * unregister carries the same values as the registration. null for a DevInterpreter, which does + * not register. + */ + private RegisterInfo getRegisterInfo() { + return host == null ? null : new RegisterInfo(host, port, interpreterGroupId); + } + public LifecycleManager getLifecycleManager() { return this.lifecycleManager; } @@ -593,7 +602,7 @@ public void run() { } } if (!Thread.currentThread().isInterrupted()) { - RegisterInfo registerInfo = new RegisterInfo(host, port, interpreterGroupId); + RegisterInfo registerInfo = getRegisterInfo(); try { intpEventClient = new RemoteInterpreterEventClient(intpEventServerHost, intpEventServerPort, 10); LOGGER.info("Registering interpreter process"); @@ -668,7 +677,7 @@ public void run() { if (intpEventClient != null && CAUSE_SHUTDOWN_HOOK.equals(cause)) { try { LOGGER.info("Unregister interpreter process"); - intpEventClient.unRegisterInterpreterProcess(); + intpEventClient.unRegisterInterpreterProcess(getRegisterInfo()); } catch (Exception e) { LOGGER.error("Fail to unregister remote interpreter process", e); } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AngularObjectId.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AngularObjectId.java index 7a62542b34e..1f7588311ae 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AngularObjectId.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AngularObjectId.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class AngularObjectId implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("AngularObjectId"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputAppendEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputAppendEvent.java index 5e19d8d3154..85bda763d7a 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputAppendEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputAppendEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class AppOutputAppendEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("AppOutputAppendEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputUpdateEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputUpdateEvent.java index 6890bffa66f..813ba0dbd71 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputUpdateEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppOutputUpdateEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class AppOutputUpdateEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("AppOutputUpdateEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppStatusUpdateEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppStatusUpdateEvent.java index af9c8292725..d550b68a931 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppStatusUpdateEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/AppStatusUpdateEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class AppStatusUpdateEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("AppStatusUpdateEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterCompletion.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterCompletion.java index af96cedca96..a7bd141cb00 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterCompletion.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterCompletion.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class InterpreterCompletion implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("InterpreterCompletion"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterRPCException.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterRPCException.java index 90f1d943a41..f35d1aaed8c 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterRPCException.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/InterpreterRPCException.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class InterpreterRPCException extends org.apache.thrift.TException implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("InterpreterRPCException"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/LibraryMetadata.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/LibraryMetadata.java index 381b8899750..07222c37d59 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/LibraryMetadata.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/LibraryMetadata.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class LibraryMetadata implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("LibraryMetadata"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputAppendEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputAppendEvent.java index fbcf43da983..6789a391258 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputAppendEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputAppendEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-20") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class OutputAppendEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("OutputAppendEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateAllEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateAllEvent.java index 7856a361e36..2d71fd6beb3 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateAllEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateAllEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-20") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class OutputUpdateAllEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("OutputUpdateAllEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateEvent.java index 41b5327e5f3..0d6655cd4d5 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/OutputUpdateEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-20") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class OutputUpdateEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("OutputUpdateEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ParagraphInfo.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ParagraphInfo.java index 0d08ce08167..6bf5d3c45d8 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ParagraphInfo.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ParagraphInfo.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class ParagraphInfo implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("ParagraphInfo"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RegisterInfo.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RegisterInfo.java index 18f7fcf7556..e843ad2f4ee 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RegisterInfo.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RegisterInfo.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RegisterInfo implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("RegisterInfo"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteApplicationResult.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteApplicationResult.java index d3d4114eeed..5ccddb05616 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteApplicationResult.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteApplicationResult.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RemoteApplicationResult implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("RemoteApplicationResult"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterContext.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterContext.java index bd7ee2d868d..bd9256ef74c 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterContext.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterContext.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RemoteInterpreterContext implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("RemoteInterpreterContext"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEvent.java index b4499f3c918..05fba259c5c 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RemoteInterpreterEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("RemoteInterpreterEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventService.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventService.java index 0ccd4f54227..28bb8c815ee 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventService.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventService.java @@ -24,14 +24,14 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RemoteInterpreterEventService { public interface Iface { public void registerInterpreterProcess(RegisterInfo registerInfo) throws org.apache.zeppelin.interpreter.thrift.InterpreterRPCException, org.apache.thrift.TException; - public void unRegisterInterpreterProcess(java.lang.String intpGroupId) throws org.apache.zeppelin.interpreter.thrift.InterpreterRPCException, org.apache.thrift.TException; + public void unRegisterInterpreterProcess(java.lang.String intpGroupId, RegisterInfo registerInfo) throws org.apache.zeppelin.interpreter.thrift.InterpreterRPCException, org.apache.thrift.TException; public void appendOutput(OutputAppendEvent event) throws org.apache.zeppelin.interpreter.thrift.InterpreterRPCException, org.apache.thrift.TException; @@ -79,7 +79,7 @@ public interface AsyncIface { public void registerInterpreterProcess(RegisterInfo registerInfo, org.apache.thrift.async.AsyncMethodCallback resultHandler) throws org.apache.thrift.TException; - public void unRegisterInterpreterProcess(java.lang.String intpGroupId, org.apache.thrift.async.AsyncMethodCallback resultHandler) throws org.apache.thrift.TException; + public void unRegisterInterpreterProcess(java.lang.String intpGroupId, RegisterInfo registerInfo, org.apache.thrift.async.AsyncMethodCallback resultHandler) throws org.apache.thrift.TException; public void appendOutput(OutputAppendEvent event, org.apache.thrift.async.AsyncMethodCallback resultHandler) throws org.apache.thrift.TException; @@ -166,16 +166,17 @@ public void recv_registerInterpreterProcess() throws org.apache.zeppelin.interpr return; } - public void unRegisterInterpreterProcess(java.lang.String intpGroupId) throws org.apache.zeppelin.interpreter.thrift.InterpreterRPCException, org.apache.thrift.TException + public void unRegisterInterpreterProcess(java.lang.String intpGroupId, RegisterInfo registerInfo) throws org.apache.zeppelin.interpreter.thrift.InterpreterRPCException, org.apache.thrift.TException { - send_unRegisterInterpreterProcess(intpGroupId); + send_unRegisterInterpreterProcess(intpGroupId, registerInfo); recv_unRegisterInterpreterProcess(); } - public void send_unRegisterInterpreterProcess(java.lang.String intpGroupId) throws org.apache.thrift.TException + public void send_unRegisterInterpreterProcess(java.lang.String intpGroupId, RegisterInfo registerInfo) throws org.apache.thrift.TException { unRegisterInterpreterProcess_args args = new unRegisterInterpreterProcess_args(); args.setIntpGroupId(intpGroupId); + args.setRegisterInfo(registerInfo); sendBase("unRegisterInterpreterProcess", args); } @@ -723,24 +724,27 @@ public Void getResult() throws org.apache.zeppelin.interpreter.thrift.Interprete } } - public void unRegisterInterpreterProcess(java.lang.String intpGroupId, org.apache.thrift.async.AsyncMethodCallback resultHandler) throws org.apache.thrift.TException { + public void unRegisterInterpreterProcess(java.lang.String intpGroupId, RegisterInfo registerInfo, org.apache.thrift.async.AsyncMethodCallback resultHandler) throws org.apache.thrift.TException { checkReady(); - unRegisterInterpreterProcess_call method_call = new unRegisterInterpreterProcess_call(intpGroupId, resultHandler, this, ___protocolFactory, ___transport); + unRegisterInterpreterProcess_call method_call = new unRegisterInterpreterProcess_call(intpGroupId, registerInfo, resultHandler, this, ___protocolFactory, ___transport); this.___currentMethod = method_call; ___manager.call(method_call); } public static class unRegisterInterpreterProcess_call extends org.apache.thrift.async.TAsyncMethodCall { private java.lang.String intpGroupId; - public unRegisterInterpreterProcess_call(java.lang.String intpGroupId, org.apache.thrift.async.AsyncMethodCallback resultHandler, org.apache.thrift.async.TAsyncClient client, org.apache.thrift.protocol.TProtocolFactory protocolFactory, org.apache.thrift.transport.TNonblockingTransport transport) throws org.apache.thrift.TException { + private RegisterInfo registerInfo; + public unRegisterInterpreterProcess_call(java.lang.String intpGroupId, RegisterInfo registerInfo, org.apache.thrift.async.AsyncMethodCallback resultHandler, org.apache.thrift.async.TAsyncClient client, org.apache.thrift.protocol.TProtocolFactory protocolFactory, org.apache.thrift.transport.TNonblockingTransport transport) throws org.apache.thrift.TException { super(client, protocolFactory, transport, resultHandler, false); this.intpGroupId = intpGroupId; + this.registerInfo = registerInfo; } public void write_args(org.apache.thrift.protocol.TProtocol prot) throws org.apache.thrift.TException { prot.writeMessageBegin(new org.apache.thrift.protocol.TMessage("unRegisterInterpreterProcess", org.apache.thrift.protocol.TMessageType.CALL, 0)); unRegisterInterpreterProcess_args args = new unRegisterInterpreterProcess_args(); args.setIntpGroupId(intpGroupId); + args.setRegisterInfo(registerInfo); args.write(prot); prot.writeMessageEnd(); } @@ -1519,7 +1523,7 @@ protected boolean rethrowUnhandledExceptions() { public unRegisterInterpreterProcess_result getResult(I iface, unRegisterInterpreterProcess_args args) throws org.apache.thrift.TException { unRegisterInterpreterProcess_result result = new unRegisterInterpreterProcess_result(); try { - iface.unRegisterInterpreterProcess(args.intpGroupId); + iface.unRegisterInterpreterProcess(args.intpGroupId, args.registerInfo); } catch (org.apache.zeppelin.interpreter.thrift.InterpreterRPCException ex) { result.ex = ex; } @@ -2261,7 +2265,7 @@ protected boolean isOneway() { } public void start(I iface, unRegisterInterpreterProcess_args args, org.apache.thrift.async.AsyncMethodCallback resultHandler) throws org.apache.thrift.TException { - iface.unRegisterInterpreterProcess(args.intpGroupId,resultHandler); + iface.unRegisterInterpreterProcess(args.intpGroupId, args.registerInfo,resultHandler); } } @@ -4290,15 +4294,18 @@ public static class unRegisterInterpreterProcess_args implements org.apache.thri private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("unRegisterInterpreterProcess_args"); private static final org.apache.thrift.protocol.TField INTP_GROUP_ID_FIELD_DESC = new org.apache.thrift.protocol.TField("intpGroupId", org.apache.thrift.protocol.TType.STRING, (short)1); + private static final org.apache.thrift.protocol.TField REGISTER_INFO_FIELD_DESC = new org.apache.thrift.protocol.TField("registerInfo", org.apache.thrift.protocol.TType.STRUCT, (short)2); private static final org.apache.thrift.scheme.SchemeFactory STANDARD_SCHEME_FACTORY = new unRegisterInterpreterProcess_argsStandardSchemeFactory(); private static final org.apache.thrift.scheme.SchemeFactory TUPLE_SCHEME_FACTORY = new unRegisterInterpreterProcess_argsTupleSchemeFactory(); public @org.apache.thrift.annotation.Nullable java.lang.String intpGroupId; // required + public @org.apache.thrift.annotation.Nullable RegisterInfo registerInfo; // required /** The set of fields this struct contains, along with convenience methods for finding and manipulating them. */ public enum _Fields implements org.apache.thrift.TFieldIdEnum { - INTP_GROUP_ID((short)1, "intpGroupId"); + INTP_GROUP_ID((short)1, "intpGroupId"), + REGISTER_INFO((short)2, "registerInfo"); private static final java.util.Map byName = new java.util.HashMap(); @@ -4316,6 +4323,8 @@ public static _Fields findByThriftId(int fieldId) { switch(fieldId) { case 1: // INTP_GROUP_ID return INTP_GROUP_ID; + case 2: // REGISTER_INFO + return REGISTER_INFO; default: return null; } @@ -4362,6 +4371,8 @@ public java.lang.String getFieldName() { java.util.Map<_Fields, org.apache.thrift.meta_data.FieldMetaData> tmpMap = new java.util.EnumMap<_Fields, org.apache.thrift.meta_data.FieldMetaData>(_Fields.class); tmpMap.put(_Fields.INTP_GROUP_ID, new org.apache.thrift.meta_data.FieldMetaData("intpGroupId", org.apache.thrift.TFieldRequirementType.DEFAULT, new org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING))); + tmpMap.put(_Fields.REGISTER_INFO, new org.apache.thrift.meta_data.FieldMetaData("registerInfo", org.apache.thrift.TFieldRequirementType.DEFAULT, + new org.apache.thrift.meta_data.StructMetaData(org.apache.thrift.protocol.TType.STRUCT, RegisterInfo.class))); metaDataMap = java.util.Collections.unmodifiableMap(tmpMap); org.apache.thrift.meta_data.FieldMetaData.addStructMetaDataMap(unRegisterInterpreterProcess_args.class, metaDataMap); } @@ -4370,10 +4381,12 @@ public unRegisterInterpreterProcess_args() { } public unRegisterInterpreterProcess_args( - java.lang.String intpGroupId) + java.lang.String intpGroupId, + RegisterInfo registerInfo) { this(); this.intpGroupId = intpGroupId; + this.registerInfo = registerInfo; } /** @@ -4383,6 +4396,9 @@ public unRegisterInterpreterProcess_args(unRegisterInterpreterProcess_args other if (other.isSetIntpGroupId()) { this.intpGroupId = other.intpGroupId; } + if (other.isSetRegisterInfo()) { + this.registerInfo = new RegisterInfo(other.registerInfo); + } } public unRegisterInterpreterProcess_args deepCopy() { @@ -4392,6 +4408,7 @@ public unRegisterInterpreterProcess_args deepCopy() { @Override public void clear() { this.intpGroupId = null; + this.registerInfo = null; } @org.apache.thrift.annotation.Nullable @@ -4419,6 +4436,31 @@ public void setIntpGroupIdIsSet(boolean value) { } } + @org.apache.thrift.annotation.Nullable + public RegisterInfo getRegisterInfo() { + return this.registerInfo; + } + + public unRegisterInterpreterProcess_args setRegisterInfo(@org.apache.thrift.annotation.Nullable RegisterInfo registerInfo) { + this.registerInfo = registerInfo; + return this; + } + + public void unsetRegisterInfo() { + this.registerInfo = null; + } + + /** Returns true if field registerInfo is set (has been assigned a value) and false otherwise */ + public boolean isSetRegisterInfo() { + return this.registerInfo != null; + } + + public void setRegisterInfoIsSet(boolean value) { + if (!value) { + this.registerInfo = null; + } + } + public void setFieldValue(_Fields field, @org.apache.thrift.annotation.Nullable java.lang.Object value) { switch (field) { case INTP_GROUP_ID: @@ -4429,6 +4471,14 @@ public void setFieldValue(_Fields field, @org.apache.thrift.annotation.Nullable } break; + case REGISTER_INFO: + if (value == null) { + unsetRegisterInfo(); + } else { + setRegisterInfo((RegisterInfo)value); + } + break; + } } @@ -4438,6 +4488,9 @@ public java.lang.Object getFieldValue(_Fields field) { case INTP_GROUP_ID: return getIntpGroupId(); + case REGISTER_INFO: + return getRegisterInfo(); + } throw new java.lang.IllegalStateException(); } @@ -4451,6 +4504,8 @@ public boolean isSet(_Fields field) { switch (field) { case INTP_GROUP_ID: return isSetIntpGroupId(); + case REGISTER_INFO: + return isSetRegisterInfo(); } throw new java.lang.IllegalStateException(); } @@ -4479,6 +4534,15 @@ public boolean equals(unRegisterInterpreterProcess_args that) { return false; } + boolean this_present_registerInfo = true && this.isSetRegisterInfo(); + boolean that_present_registerInfo = true && that.isSetRegisterInfo(); + if (this_present_registerInfo || that_present_registerInfo) { + if (!(this_present_registerInfo && that_present_registerInfo)) + return false; + if (!this.registerInfo.equals(that.registerInfo)) + return false; + } + return true; } @@ -4490,6 +4554,10 @@ public int hashCode() { if (isSetIntpGroupId()) hashCode = hashCode * 8191 + intpGroupId.hashCode(); + hashCode = hashCode * 8191 + ((isSetRegisterInfo()) ? 131071 : 524287); + if (isSetRegisterInfo()) + hashCode = hashCode * 8191 + registerInfo.hashCode(); + return hashCode; } @@ -4511,6 +4579,16 @@ public int compareTo(unRegisterInterpreterProcess_args other) { return lastComparison; } } + lastComparison = java.lang.Boolean.valueOf(isSetRegisterInfo()).compareTo(other.isSetRegisterInfo()); + if (lastComparison != 0) { + return lastComparison; + } + if (isSetRegisterInfo()) { + lastComparison = org.apache.thrift.TBaseHelper.compareTo(this.registerInfo, other.registerInfo); + if (lastComparison != 0) { + return lastComparison; + } + } return 0; } @@ -4539,6 +4617,14 @@ public java.lang.String toString() { sb.append(this.intpGroupId); } first = false; + if (!first) sb.append(", "); + sb.append("registerInfo:"); + if (this.registerInfo == null) { + sb.append("null"); + } else { + sb.append(this.registerInfo); + } + first = false; sb.append(")"); return sb.toString(); } @@ -4546,6 +4632,9 @@ public java.lang.String toString() { public void validate() throws org.apache.thrift.TException { // check for required fields // check for sub-struct validity + if (registerInfo != null) { + registerInfo.validate(); + } } private void writeObject(java.io.ObjectOutputStream out) throws java.io.IOException { @@ -4590,6 +4679,15 @@ public void read(org.apache.thrift.protocol.TProtocol iprot, unRegisterInterpret org.apache.thrift.protocol.TProtocolUtil.skip(iprot, schemeField.type); } break; + case 2: // REGISTER_INFO + if (schemeField.type == org.apache.thrift.protocol.TType.STRUCT) { + struct.registerInfo = new RegisterInfo(); + struct.registerInfo.read(iprot); + struct.setRegisterInfoIsSet(true); + } else { + org.apache.thrift.protocol.TProtocolUtil.skip(iprot, schemeField.type); + } + break; default: org.apache.thrift.protocol.TProtocolUtil.skip(iprot, schemeField.type); } @@ -4610,6 +4708,11 @@ public void write(org.apache.thrift.protocol.TProtocol oprot, unRegisterInterpre oprot.writeString(struct.intpGroupId); oprot.writeFieldEnd(); } + if (struct.registerInfo != null) { + oprot.writeFieldBegin(REGISTER_INFO_FIELD_DESC); + struct.registerInfo.write(oprot); + oprot.writeFieldEnd(); + } oprot.writeFieldStop(); oprot.writeStructEnd(); } @@ -4631,20 +4734,31 @@ public void write(org.apache.thrift.protocol.TProtocol prot, unRegisterInterpret if (struct.isSetIntpGroupId()) { optionals.set(0); } - oprot.writeBitSet(optionals, 1); + if (struct.isSetRegisterInfo()) { + optionals.set(1); + } + oprot.writeBitSet(optionals, 2); if (struct.isSetIntpGroupId()) { oprot.writeString(struct.intpGroupId); } + if (struct.isSetRegisterInfo()) { + struct.registerInfo.write(oprot); + } } @Override public void read(org.apache.thrift.protocol.TProtocol prot, unRegisterInterpreterProcess_args struct) throws org.apache.thrift.TException { org.apache.thrift.protocol.TTupleProtocol iprot = (org.apache.thrift.protocol.TTupleProtocol) prot; - java.util.BitSet incoming = iprot.readBitSet(1); + java.util.BitSet incoming = iprot.readBitSet(2); if (incoming.get(0)) { struct.intpGroupId = iprot.readString(); struct.setIntpGroupIdIsSet(true); } + if (incoming.get(1)) { + struct.registerInfo = new RegisterInfo(); + struct.registerInfo.read(iprot); + struct.setRegisterInfoIsSet(true); + } } } diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java index 5d8244ae899..508b20dd8bd 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterEventType.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public enum RemoteInterpreterEventType implements org.apache.thrift.TEnum { NO_OP(1), ANGULAR_OBJECT_ADD(2), diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResult.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResult.java index 00ceeb62e5f..93024c3bc9d 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResult.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResult.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RemoteInterpreterResult implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("RemoteInterpreterResult"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResultMessage.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResultMessage.java index 3a00fe14f67..ed9b2f8de8a 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResultMessage.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterResultMessage.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RemoteInterpreterResultMessage implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("RemoteInterpreterResultMessage"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterService.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterService.java index a6ad9d5f976..842aca095a7 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterService.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RemoteInterpreterService.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RemoteInterpreterService { public interface Iface { diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RunParagraphsEvent.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RunParagraphsEvent.java index d9272dd4cc4..6343f865427 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RunParagraphsEvent.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/RunParagraphsEvent.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class RunParagraphsEvent implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("RunParagraphsEvent"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ServiceException.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ServiceException.java index 469b9d8ffea..05d31c3814f 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ServiceException.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/ServiceException.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class ServiceException extends org.apache.thrift.TException implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("ServiceException"); diff --git a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/WebUrlInfo.java b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/WebUrlInfo.java index 3dd97226515..5f69e9af76e 100644 --- a/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/WebUrlInfo.java +++ b/zeppelin-interpreter/src/main/java/org/apache/zeppelin/interpreter/thrift/WebUrlInfo.java @@ -24,7 +24,7 @@ package org.apache.zeppelin.interpreter.thrift; @SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"}) -@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-13") +@javax.annotation.Generated(value = "Autogenerated by Thrift Compiler (0.13.0)", date = "2026-09-30") public class WebUrlInfo implements org.apache.thrift.TBase, java.io.Serializable, Cloneable, Comparable { private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new org.apache.thrift.protocol.TStruct("WebUrlInfo"); diff --git a/zeppelin-interpreter/src/main/thrift/RemoteInterpreterEventService.thrift b/zeppelin-interpreter/src/main/thrift/RemoteInterpreterEventService.thrift index c7c6e039f0c..30f04f31eb0 100644 --- a/zeppelin-interpreter/src/main/thrift/RemoteInterpreterEventService.thrift +++ b/zeppelin-interpreter/src/main/thrift/RemoteInterpreterEventService.thrift @@ -113,7 +113,7 @@ exception ServiceException{ service RemoteInterpreterEventService { void registerInterpreterProcess(1: RegisterInfo registerInfo) throws (1: RemoteInterpreterService.InterpreterRPCException ex); - void unRegisterInterpreterProcess(1: string intpGroupId) throws (1: RemoteInterpreterService.InterpreterRPCException ex); + void unRegisterInterpreterProcess(1: string intpGroupId, 2: RegisterInfo registerInfo) throws (1: RemoteInterpreterService.InterpreterRPCException ex); void appendOutput(1: OutputAppendEvent event) throws (1: RemoteInterpreterService.InterpreterRPCException ex); void updateOutput(1: OutputUpdateEvent event) throws (1: RemoteInterpreterService.InterpreterRPCException ex); diff --git a/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServerTest.java b/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServerTest.java index 58ef866a744..7f90d4c8640 100644 --- a/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServerTest.java +++ b/zeppelin-interpreter/src/test/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterServerTest.java @@ -23,9 +23,11 @@ import org.apache.zeppelin.interpreter.InterpreterException; import org.apache.zeppelin.interpreter.InterpreterResult; import org.apache.zeppelin.interpreter.LazyOpenInterpreter; +import org.apache.zeppelin.interpreter.thrift.RegisterInfo; import org.apache.zeppelin.interpreter.thrift.RemoteInterpreterContext; import org.apache.zeppelin.interpreter.thrift.RemoteInterpreterResult; import org.junit.jupiter.api.Test; +import org.mockito.ArgumentCaptor; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -41,6 +43,7 @@ import static org.junit.jupiter.api.Assertions.assertNotNull; import static org.junit.jupiter.api.Assertions.assertTrue; import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.verify; class RemoteInterpreterServerTest { @@ -66,6 +69,23 @@ void testStartStopWithQueuedEvents() throws Exception { stopRemoteInterpreterServer(server, 10 * 10000); } + @Test + void testShutdownUnregistersWithItsHostAndPort() throws Exception { + RemoteInterpreterServer server = new RemoteInterpreterServer("localhost", + RemoteInterpreterUtils.findRandomAvailablePortOnAllLocalInterfaces(), ":", "groupId", true); + server.intpEventClient = mock(RemoteInterpreterEventClient.class); + startRemoteInterpreterServer(server, 10 * 1000); + + stopRemoteInterpreterServer(server, 10 * 10000); + + // The server compares these with what the process registered with. + ArgumentCaptor registerInfo = ArgumentCaptor.forClass(RegisterInfo.class); + verify(server.intpEventClient).unRegisterInterpreterProcess(registerInfo.capture()); + assertNotNull(registerInfo.getValue().getHost()); + assertEquals(server.getPort(), registerInfo.getValue().getPort()); + assertEquals("groupId", registerInfo.getValue().getInterpreterGroupId()); + } + private void startRemoteInterpreterServer(RemoteInterpreterServer server, int timeout) throws InterruptedException, TException { assertEquals(false, server.isRunning()); diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java index f26bc54f0e7..fe36a0bea7d 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/InterpreterSetting.java @@ -476,6 +476,22 @@ void removeInterpreterGroup(String groupId) { } } + /** + * Removes the given group only if it is still the one registered under its id. A group with the + * same id can be created while this one is closing, and that group must stay registered. + * Compares references, because {@link InterpreterGroup#equals} compares ids only. + */ + void removeInterpreterGroup(ManagedInterpreterGroup interpreterGroup) { + try { + interpreterGroupWriteLock.lock(); + if (this.interpreterGroups.get(interpreterGroup.getId()) == interpreterGroup) { + this.interpreterGroups.remove(interpreterGroup.getId()); + } + } finally { + interpreterGroupWriteLock.unlock(); + } + } + public ManagedInterpreterGroup getInterpreterGroup(String user, String noteId) { return getInterpreterGroup(getExecutionContext(user, noteId)); } 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 1a85f82d5d6..7f974912eb3 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 @@ -733,13 +733,6 @@ public List getInterpreterProcessStatuses() { 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()) { - interpreterSetting.removeInterpreterGroup(intpGroupId); - } - } - //TODO(zjffdu) move Resource related api to ResourceManager public ResourceSet getAllResources() { return getAllResourcesExcept(null); 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 3a8b14ee81e..f92de80e843 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 @@ -21,6 +21,7 @@ import org.apache.zeppelin.conf.ZeppelinConfiguration; import org.apache.zeppelin.interpreter.lifecycle.IdleInterpreterReclaimer; import org.apache.zeppelin.interpreter.remote.RemoteInterpreterProcess; +import org.apache.zeppelin.interpreter.thrift.RegisterInfo; import org.apache.zeppelin.scheduler.Job; import org.apache.zeppelin.scheduler.Scheduler; import org.apache.zeppelin.scheduler.SchedulerFactory; @@ -46,6 +47,8 @@ public class ManagedInterpreterGroup extends InterpreterGroup { private final ZeppelinConfiguration zConf; private volatile long lastUsedTimeInMillis = System.currentTimeMillis(); private volatile boolean launchingInterpreterProcess; + // The registration accepted for this group's process, null until the process registers. + private volatile RegisterInfo registerInfo; /** * Create InterpreterGroup with given id and interpreterSetting, used in ZeppelinServer @@ -123,6 +126,14 @@ public RemoteInterpreterProcess getRemoteInterpreterProcess() { return remoteInterpreterProcess; } + RegisterInfo getRegisterInfo() { + return registerInfo; + } + + void setRegisterInfo(RegisterInfo registerInfo) { + this.registerInfo = registerInfo; + } + /** * Close all interpreter instances in this group diff --git a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServer.java b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServer.java index 5504283308b..94cc00abb15 100644 --- a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServer.java +++ b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServer.java @@ -184,23 +184,54 @@ public void registerInterpreterProcess(RegisterInfo registerInfo) throws Interpr LOGGER.info("Register interpreter process: {}:{}, interpreterGroup: {}", registerInfo.getHost(), registerInfo.getPort(), registerInfo.getInterpreterGroupId()); interpreterProcess.processStarted(registerInfo.port, registerInfo.host); + ((ManagedInterpreterGroup) interpreterGroup).setRegisterInfo(registerInfo); } @Override - public void unRegisterInterpreterProcess(String intpGroupId) throws InterpreterRPCException, TException { + public void unRegisterInterpreterProcess(String intpGroupId, RegisterInfo registerInfo) + throws InterpreterRPCException, TException { LOGGER.info("Unregister interpreter process: {}", intpGroupId); - InterpreterGroup interpreterGroup = + ManagedInterpreterGroup interpreterGroup = interpreterSettingManager.getInterpreterGroupById(intpGroupId); if (interpreterGroup == null) { LOGGER.warn("Unable to unregister interpreter process because no such interpreterGroup: {}", intpGroupId); return; } + if (!isSentByProcessOf(interpreterGroup, registerInfo)) { + LOGGER.warn("Ignore the unregister of interpreter process {}:{}, because it is not the " + + "process of interpreterGroup: {}", registerInfo.getHost(), registerInfo.getPort(), + intpGroupId); + return; + } // Close RemoteInterpreter when RemoteInterpreterServer already timeout. // Otherwise the ProgressBar will be missing when rerun after the RemoteInterpreterServer timeout // and old RemoteInterpreterGroups will always alive after GC. interpreterGroup.close(); - interpreterSettingManager.removeInterpreterGroup(intpGroupId); + if (interpreterGroup.getInterpreterSetting() != null) { + interpreterGroup.getInterpreterSetting().removeInterpreterGroup(interpreterGroup); + } + } + + /** + * The unregister is resolved by group id, and a group with that id can belong to another process + * than the sender, e.g. when a new group is created while the old process is being stopped. + */ + private static boolean isSentByProcessOf(ManagedInterpreterGroup interpreterGroup, + RegisterInfo sender) { + if (sender == null) { + // An interpreter process from before this check does not say who it is. + return true; + } + RegisterInfo registered = interpreterGroup.getRegisterInfo(); + if (registered != null) { + return registered.equals(sender); + } + // No registration was accepted for this group. Without a process, or while its process is + // still launching, it has no process that could send this. A recovered or an externally + // running process is attached without a registration, so keep the id based behaviour for it. + return interpreterGroup.getInterpreterProcess() != null + && !interpreterGroup.isLaunchingInterpreterProcess(); } @Override diff --git a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerRegistrationTest.java b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerRegistrationTest.java new file mode 100644 index 00000000000..2656f22cbdd --- /dev/null +++ b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerRegistrationTest.java @@ -0,0 +1,163 @@ +/* + * 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 static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertNull; +import static org.junit.jupiter.api.Assertions.assertSame; +import static org.mockito.ArgumentMatchers.anyString; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import com.google.common.collect.Lists; +import java.lang.reflect.Field; +import java.util.HashMap; +import org.apache.zeppelin.conf.ZeppelinConfiguration; +import org.apache.zeppelin.interpreter.remote.RemoteInterpreterProcess; +import org.apache.zeppelin.interpreter.thrift.RegisterInfo; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +/** + * An unregister is resolved by interpreter group id, and the group with that id can belong to + * another process than the sender. + */ +class RemoteInterpreterEventServerRegistrationTest { + + private InterpreterSetting interpreterSetting; + private RemoteInterpreterEventServer eventServer; + private ManagedInterpreterGroup group; + // Two processes of the same interpreter group id, e.g. one being stopped and its replacement. + private RegisterInfo oldProcess; + private RegisterInfo newProcess; + + @BeforeEach + void setUp() { + ZeppelinConfiguration zConf = ZeppelinConfiguration.load(); + InterpreterOption option = new InterpreterOption(); + option.setPerUser(InterpreterOption.SHARED); + InterpreterInfo interpreterInfo = new InterpreterInfo(EchoInterpreter.class.getName(), + "echo", true, new HashMap(), new HashMap()); + interpreterSetting = new InterpreterSetting.Builder() + .setId("id") + .setName("test") + .setGroup("test") + .setInterpreterInfos(Lists.newArrayList(interpreterInfo)) + .setOption(option) + .setConf(zConf) + .create(); + + InterpreterSettingManager interpreterSettingManager = mock(InterpreterSettingManager.class); + when(interpreterSettingManager.getInterpreterGroupById(anyString())) + .thenAnswer(invocation -> interpreterSetting.getInterpreterGroup( + (String) invocation.getArgument(0))); + eventServer = new RemoteInterpreterEventServer(zConf, interpreterSettingManager); + + group = interpreterSetting.getOrCreateInterpreterGroup(new ExecutionContext("user1", "note1", + "test")); + group.getOrCreateSession("user1", "shared_session"); + oldProcess = new RegisterInfo("10.0.0.1", 30001, group.getId()); + newProcess = new RegisterInfo("10.0.0.1", 30002, group.getId()); + } + + @Test + void unregisterFromTheRegisteredProcessClosesItsGroup() throws Exception { + setInterpreterProcess(group, mock(RemoteInterpreterProcess.class)); + eventServer.registerInterpreterProcess(newProcess); + + eventServer.unRegisterInterpreterProcess(group.getId(), newProcess); + + assertClosed(); + } + + @Test + void unregisterFromAnotherProcessKeepsARegisteredGroup() throws Exception { + // The old process was stopped or left behind, and the group with its id is now a new one. + setInterpreterProcess(group, mock(RemoteInterpreterProcess.class)); + eventServer.registerInterpreterProcess(newProcess); + + eventServer.unRegisterInterpreterProcess(group.getId(), oldProcess); + + assertKept(); + } + + @Test + void unregisterKeepsAGroupWhoseProcessIsStillLaunching() throws Exception { + setInterpreterProcess(group, mock(RemoteInterpreterProcess.class)); + setLaunching(group, true); + + eventServer.unRegisterInterpreterProcess(group.getId(), oldProcess); + + assertKept(); + } + + @Test + void unregisterKeepsAGroupWithoutProcess() throws Exception { + eventServer.unRegisterInterpreterProcess(group.getId(), oldProcess); + + assertKept(); + } + + @Test + void unregisterWithoutRegisterInfoClosesTheGroup() throws Exception { + // An interpreter process from before the sender was sent along. + setInterpreterProcess(group, mock(RemoteInterpreterProcess.class)); + eventServer.registerInterpreterProcess(newProcess); + + eventServer.unRegisterInterpreterProcess(group.getId(), null); + + assertClosed(); + } + + @Test + void unregisterClosesAGroupWhoseProcessWasAttachedWithoutRegistration() throws Exception { + // A recovered or an externally running process does not register. + setInterpreterProcess(group, mock(RemoteInterpreterProcess.class)); + + eventServer.unRegisterInterpreterProcess(group.getId(), oldProcess); + + assertClosed(); + } + + private void assertClosed() { + assertEquals(0, group.getSessionNum()); + assertNull(interpreterSetting.getInterpreterGroup(group.getId())); + } + + private void assertKept() { + assertEquals(1, group.getSessionNum()); + assertSame(group, interpreterSetting.getInterpreterGroup(group.getId())); + } + + private static void setInterpreterProcess(ManagedInterpreterGroup interpreterGroup, + RemoteInterpreterProcess process) throws Exception { + setField(interpreterGroup, "remoteInterpreterProcess", process); + } + + private static void setLaunching(ManagedInterpreterGroup interpreterGroup, boolean launching) + throws Exception { + setField(interpreterGroup, "launchingInterpreterProcess", launching); + } + + private static void setField(ManagedInterpreterGroup interpreterGroup, String name, Object value) + throws Exception { + Field field = ManagedInterpreterGroup.class.getDeclaredField(name); + field.setAccessible(true); + field.set(interpreterGroup, value); + } +}