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 8916d7aee1 [ZEPPELIN-6704] Scope personalized streaming output to its
execution user
8916d7aee1 is described below
commit 8916d7aee1982d32b8f8538518f3ac9dbd3a5b7a
Author: dae won <[email protected]>
AuthorDate: Thu Sep 24 22:34:07 2026 +0900
[ZEPPELIN-6704] Scope personalized streaming output to its execution user
### What is this PR for?
In a personalized note, paragraph results are stored per user, but the
incremental output produced while a paragraph runs is not. The interpreter
knows who started the execution, yet the output events it sends back carry no
owner, so `NotebookServer` cannot tell who the output belongs to. Personalized
incremental output is therefore dropped, and a user sees nothing until the
paragraph finishes.
The owner is lost in two places. `OutputAppendEvent` and
`OutputUpdateEvent` have no user field, and `createInterpreterOutput` captures
only the note and paragraph ids from the interpreter context, even though the
authentication info is available there too. Recovering it on the server is not
an option either. The shared paragraph's `user` field is updated only when the
note owner runs it, so it can hold a stale user from an earlier run rather than
the owner of the current output.
This PR carries the execution owner end to end, in three commits:
1. Add an optional `user` field to the interpreter output events and thread
it from the interpreter to the server. The interpreter captures the owner when
the output is created, before any chunks are produced, so every event from that
execution carries the same owner. `AppendOutputRunner` now keys its buffer by
owner as well, because two users running the same paragraph would otherwise
have their chunks merged into the same buffer, leaving one chunk with no single
owner. Behaviour is [...]
2. Deliver personalized append and update events to the execution owner
through `multicastToUser`, write updates to that user's own paragraph copy
rather than the shared one, and make clear events discard only the owner's
output. An interpreter that predates the field reports no owner, and that
output is still dropped rather than broadcast.
3. Add two-user e2e coverage.
Shared mode is untouched. A shared paragraph has only one execution
producing output for a given index at a time, so adding the owner to the buffer
key does not change its batching behavior.
`AppendOutputRunnerTest.appendsFromOneExecutionAreStillBatchedIntoOneChunk`
verifies this.
Two notes on scope. The paragraph output callbacks on
`RemoteInterpreterProcessListener` are renamed to `onParagraphOutputAppend`,
`onParagraphOutputUpdated` and `onParagraphOutputClear`: once the owner is
added, they have the same erasure as the app output overloads that
`NotebookServer` also implements, so the class would not compile. The Helium
app output events carry no owner either, but I could not find any personalized
code path that uses them, so they are left out of this PR.
### What type of PR is it?
Bug Fix
### Todos
* [x] - Carry the execution owner from the interpreter to `NotebookServer`
* [x] - Key the append buffer by owner so concurrent runs are not merged
* [x] - Route personalized append, update and clear to the execution owner
* [x] - Add server tests and a two-user e2e test
### What is the Jira issue?
* [ZEPPELIN-6704](https://issues.apache.org/jira/browse/ZEPPELIN-6704)
### How should this be tested?
Server side:
```
./mvnw package -pl zeppelin-server --am \
-Dtest="NotebookServerStreamingScopeTest,AppendOutputRunnerTest,RemoteInterpreterEventServerTest"
\
-DfailIfNoTests=false -Dsurefire.failIfNoSpecifiedTests=false
```
27 tests pass. `NotebookServerStreamingScopeTest` covers the routing,
including a test that reproduces the original failure mode: the shared
paragraph is given a stale owner and a different user runs the paragraph, and
the output must still reach the user who started that execution.
The e2e test signs in as two shiro accounts in separate browser contexts,
opens one personalized note, and records the `PARAGRAPH_APPEND_OUTPUT` and
`PARAGRAPH_UPDATE_OUTPUT` frames each socket receives. The test asserts on
WebSocket frames rather than rendered UI because the bug is in server-side
connection routing. It runs on the authenticated leg:
```
cd zeppelin-web-angular && npx playwright test
e2e/tests/notebook/personalized --project=chromium
```
Reverting the routing change makes the new tests fail:
```
NotebookServerStreamingScopeTest.personalizedOutputReachesOnlyItsExecutionOwner
Wanted but not invoked: connectionManager.multicastToUser("owner", ...)
NotebookServerStreamingScopeTest.personalizedUpdateWritesToTheOwnersCopyNotTheSharedParagraph
the owner's copy did not receive its own output ==> expected: <1> but
was: <0>
NotebookServerStreamingScopeTest.personalizedClearOnlyDiscardsTheOwnersOutput
the owner's output survived its own clear ==> expected: <null> but was:
<%text owned>
personalized-streaming-scope.spec.ts
no marker in 0 streaming frames
Expected substring: "RUNNER-ONLY" / Received string: ""
```
### Screenshots (if appropriate)
N/A
### Questions:
* Does the license files need to update? No
* Is there breaking changes for older versions? No. The thrift field is
optional, so an older interpreter process simply reports no owner and its
personalized output is dropped instead of broadcast.
* Does this needs documentation? No
Closes #5479 from big-cir/ZEPPELIN-6704.
Signed-off-by: ChanHo Lee <[email protected]>
---
.../helium/ZeppelinApplicationDevServer.java | 7 +-
.../apache/zeppelin/helium/ZeppelinDevServer.java | 7 +-
.../remote/RemoteInterpreterEventClient.java | 19 +--
.../remote/RemoteInterpreterServer.java | 17 ++-
.../interpreter/thrift/AngularObjectId.java | 2 +-
.../interpreter/thrift/AppOutputAppendEvent.java | 2 +-
.../interpreter/thrift/AppOutputUpdateEvent.java | 2 +-
.../interpreter/thrift/AppStatusUpdateEvent.java | 2 +-
.../interpreter/thrift/InterpreterCompletion.java | 2 +-
.../thrift/InterpreterRPCException.java | 2 +-
.../interpreter/thrift/LibraryMetadata.java | 2 +-
.../interpreter/thrift/OutputAppendEvent.java | 115 ++++++++++++++-
.../interpreter/thrift/OutputUpdateAllEvent.java | 115 ++++++++++++++-
.../interpreter/thrift/OutputUpdateEvent.java | 115 ++++++++++++++-
.../zeppelin/interpreter/thrift/ParagraphInfo.java | 2 +-
.../zeppelin/interpreter/thrift/RegisterInfo.java | 2 +-
.../thrift/RemoteApplicationResult.java | 2 +-
.../thrift/RemoteInterpreterContext.java | 2 +-
.../interpreter/thrift/RemoteInterpreterEvent.java | 2 +-
.../thrift/RemoteInterpreterEventService.java | 2 +-
.../thrift/RemoteInterpreterEventType.java | 2 +-
.../thrift/RemoteInterpreterResult.java | 2 +-
.../thrift/RemoteInterpreterResultMessage.java | 2 +-
.../thrift/RemoteInterpreterService.java | 2 +-
.../interpreter/thrift/RunParagraphsEvent.java | 2 +-
.../interpreter/thrift/ServiceException.java | 2 +-
.../zeppelin/interpreter/thrift/WebUrlInfo.java | 2 +-
.../thrift/RemoteInterpreterEventService.thrift | 7 +-
.../interpreter/RemoteInterpreterEventServer.java | 15 +-
.../interpreter/remote/AppendOutputBuffer.java | 9 +-
.../interpreter/remote/AppendOutputRunner.java | 89 +++++++++---
.../remote/RemoteInterpreterProcessListener.java | 12 +-
.../interpreter/remote/UpdateOutputBuffer.java | 4 +-
.../org/apache/zeppelin/socket/NotebookServer.java | 53 +++++--
.../RemoteInterpreterEventServerTest.java | 47 +++---
.../interpreter/remote/AppendOutputRunnerTest.java | 141 ++++++++++++------
.../socket/NotebookServerStreamingScopeTest.java | 107 +++++++++++++-
.../personalized-streaming-scope.spec.ts | 158 +++++++++++++++++++++
zeppelin-web-angular/e2e/utils.ts | 37 ++++-
39 files changed, 936 insertions(+), 178 deletions(-)
diff --git
a/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinApplicationDevServer.java
b/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinApplicationDevServer.java
index fe1f191a01..1c6ebcb16d 100644
---
a/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinApplicationDevServer.java
+++
b/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinApplicationDevServer.java
@@ -136,7 +136,7 @@ public class ZeppelinApplicationDevServer extends
ZeppelinDevServer {
@Override
protected InterpreterOutput createInterpreterOutput(
- final String noteId, final String paragraphId) {
+ final String noteId, final String paragraphId, final String
executionOwner) {
if (out == null) {
final RemoteInterpreterEventClient eventClient = getIntpEventClient();
try {
@@ -148,14 +148,15 @@ public class ZeppelinApplicationDevServer extends
ZeppelinDevServer {
@Override
public void onAppend(int index, InterpreterResultMessageOutput out,
byte[] line) {
- eventClient.onInterpreterOutputAppend(noteId, paragraphId, index,
new String(line));
+ eventClient.onInterpreterOutputAppend(
+ noteId, paragraphId, index, executionOwner, new String(line));
}
@Override
public void onUpdate(int index, InterpreterResultMessageOutput out) {
try {
eventClient.onInterpreterOutputUpdate(noteId, paragraphId,
- index, out.getType(), new String(out.toByteArray()));
+ index, executionOwner, out.getType(), new
String(out.toByteArray()));
} catch (IOException e) {
LOGGER.error(e.getMessage(), e);
}
diff --git
a/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinDevServer.java
b/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinDevServer.java
index 81e0a61018..116f1adab1 100644
--- a/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinDevServer.java
+++ b/helium-dev/src/main/java/org/apache/zeppelin/helium/ZeppelinDevServer.java
@@ -66,7 +66,7 @@ public class ZeppelinDevServer extends
@Override
protected InterpreterOutput createInterpreterOutput(
- final String noteId, final String paragraphId) {
+ final String noteId, final String paragraphId, final String
executionOwner) {
if (out == null) {
final RemoteInterpreterEventClient eventClient = getIntpEventClient();
try {
@@ -78,14 +78,15 @@ public class ZeppelinDevServer extends
@Override
public void onAppend(int index, InterpreterResultMessageOutput out,
byte[] line) {
- eventClient.onInterpreterOutputAppend(noteId, paragraphId, index,
new String(line));
+ eventClient.onInterpreterOutputAppend(
+ noteId, paragraphId, index, executionOwner, new String(line));
}
@Override
public void onUpdate(int index, InterpreterResultMessageOutput out) {
try {
eventClient.onInterpreterOutputUpdate(noteId, paragraphId,
- index, out.getType(), new String(out.toByteArray()));
+ index, executionOwner, out.getType(), new
String(out.toByteArray()));
} catch (IOException e) {
LOGGER.error(e.getMessage(), e);
}
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 174de3bc19..df281c84d4 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
@@ -222,11 +222,11 @@ public class RemoteInterpreterEventClient implements
ResourcePoolConnector,
}
public void onInterpreterOutputAppend(
- String noteId, String paragraphId, int outputIndex, String output) {
+ String noteId, String paragraphId, int outputIndex, String
executionOwner, String output) {
try {
callRemoteFunction(client -> {
- client.appendOutput(
- new OutputAppendEvent(noteId, paragraphId, outputIndex,
output, null));
+ client.appendOutput(new OutputAppendEvent(
+ noteId, paragraphId, outputIndex, output, null,
executionOwner));
return null;
});
} catch (Exception e) {
@@ -235,12 +235,12 @@ public class RemoteInterpreterEventClient implements
ResourcePoolConnector,
}
public void onInterpreterOutputUpdate(
- String noteId, String paragraphId, int outputIndex,
+ String noteId, String paragraphId, int outputIndex, String
executionOwner,
InterpreterResult.Type type, String output) {
try {
callRemoteFunction(client -> {
- client.updateOutput(
- new OutputUpdateEvent(noteId, paragraphId, outputIndex,
type.name(), output, null));
+ client.updateOutput(new OutputUpdateEvent(
+ noteId, paragraphId, outputIndex, type.name(), output, null,
executionOwner));
return null;
});
@@ -250,11 +250,12 @@ public class RemoteInterpreterEventClient implements
ResourcePoolConnector,
}
public void onInterpreterOutputUpdateAll(
- String noteId, String paragraphId, List<InterpreterResultMessage>
messages) {
+ String noteId, String paragraphId, String executionOwner,
+ List<InterpreterResultMessage> messages) {
try {
callRemoteFunction(client -> {
- client.updateAllOutput(
- new OutputUpdateAllEvent(noteId, paragraphId,
convertToThrift(messages)));
+ client.updateAllOutput(new OutputUpdateAllEvent(
+ noteId, paragraphId, convertToThrift(messages),
executionOwner));
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 e733fd57f8..346a8641ea 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
@@ -955,7 +955,14 @@ public class RemoteInterpreterServer extends Thread
}
private InterpreterContext convert(RemoteInterpreterContext ric) {
- return convert(ric, createInterpreterOutput(ric.getNoteId(),
ric.getParagraphId()));
+ return convert(ric, createInterpreterOutput(
+ ric.getNoteId(), ric.getParagraphId(), executionOwnerOf(ric)));
+ }
+
+ private static String executionOwnerOf(RemoteInterpreterContext ric) {
+ AuthenticationInfo authenticationInfo =
+ AuthenticationInfo.fromJson(ric.getAuthenticationInfo());
+ return authenticationInfo == null ? null : authenticationInfo.getUser();
}
private InterpreterContext convert(RemoteInterpreterContext ric,
InterpreterOutput output) {
@@ -982,13 +989,13 @@ public class RemoteInterpreterServer extends Thread
protected InterpreterOutput createInterpreterOutput(final String noteId,
final String
- paragraphId) {
+ paragraphId, final String executionOwner) {
return new InterpreterOutput(new InterpreterOutputListener() {
@Override
public void onUpdateAll(InterpreterOutput out) {
try {
intpEventClient.onInterpreterOutputUpdateAll(
- noteId, paragraphId, out.toInterpreterResultMessage());
+ noteId, paragraphId, executionOwner,
out.toInterpreterResultMessage());
} catch (IOException e) {
LOGGER.error(e.getMessage(), e);
}
@@ -999,7 +1006,7 @@ public class RemoteInterpreterServer extends Thread
String output = new String(line);
LOGGER.debug("Output Append: {}", output);
intpEventClient.onInterpreterOutputAppend(
- noteId, paragraphId, index, output);
+ noteId, paragraphId, index, executionOwner, output);
}
@Override
@@ -1009,7 +1016,7 @@ public class RemoteInterpreterServer extends Thread
output = new String(out.toByteArray());
LOGGER.debug("Output Update for index {}: {}", index, output);
intpEventClient.onInterpreterOutputUpdate(
- noteId, paragraphId, index, out.getType(), output);
+ noteId, paragraphId, index, executionOwner, out.getType(),
output);
} catch (IOException e) {
LOGGER.error(e.getMessage(), 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 8f93f2ed44..7a62542b34 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class AngularObjectId implements
org.apache.thrift.TBase<AngularObjectId, AngularObjectId._Fields>,
java.io.Serializable, Cloneable, Comparable<AngularObjectId> {
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 676e46919c..5e19d8d315 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class AppOutputAppendEvent implements
org.apache.thrift.TBase<AppOutputAppendEvent, AppOutputAppendEvent._Fields>,
java.io.Serializable, Cloneable, Comparable<AppOutputAppendEvent> {
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 46b0ac6eec..6890bffa66 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class AppOutputUpdateEvent implements
org.apache.thrift.TBase<AppOutputUpdateEvent, AppOutputUpdateEvent._Fields>,
java.io.Serializable, Cloneable, Comparable<AppOutputUpdateEvent> {
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 57a18c645c..af9c829272 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class AppStatusUpdateEvent implements
org.apache.thrift.TBase<AppStatusUpdateEvent, AppStatusUpdateEvent._Fields>,
java.io.Serializable, Cloneable, Comparable<AppStatusUpdateEvent> {
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 7a92c079fa..af96cedca9 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class InterpreterCompletion implements
org.apache.thrift.TBase<InterpreterCompletion, InterpreterCompletion._Fields>,
java.io.Serializable, Cloneable, Comparable<InterpreterCompletion> {
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 f86ac4de5e..90f1d943a4 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class InterpreterRPCException extends org.apache.thrift.TException
implements org.apache.thrift.TBase<InterpreterRPCException,
InterpreterRPCException._Fields>, java.io.Serializable, Cloneable,
Comparable<InterpreterRPCException> {
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 4299eb6711..381b889975 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class LibraryMetadata implements
org.apache.thrift.TBase<LibraryMetadata, LibraryMetadata._Fields>,
java.io.Serializable, Cloneable, Comparable<LibraryMetadata> {
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 67a3a8f054..fbcf43da98 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-20")
public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEvent, OutputAppendEvent._Fields>,
java.io.Serializable, Cloneable, Comparable<OutputAppendEvent> {
private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new
org.apache.thrift.protocol.TStruct("OutputAppendEvent");
@@ -33,6 +33,7 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
private static final org.apache.thrift.protocol.TField INDEX_FIELD_DESC =
new org.apache.thrift.protocol.TField("index",
org.apache.thrift.protocol.TType.I32, (short)3);
private static final org.apache.thrift.protocol.TField DATA_FIELD_DESC = new
org.apache.thrift.protocol.TField("data",
org.apache.thrift.protocol.TType.STRING, (short)4);
private static final org.apache.thrift.protocol.TField APP_ID_FIELD_DESC =
new org.apache.thrift.protocol.TField("appId",
org.apache.thrift.protocol.TType.STRING, (short)5);
+ private static final org.apache.thrift.protocol.TField
EXECUTION_OWNER_FIELD_DESC = new
org.apache.thrift.protocol.TField("executionOwner",
org.apache.thrift.protocol.TType.STRING, (short)6);
private static final org.apache.thrift.scheme.SchemeFactory
STANDARD_SCHEME_FACTORY = new OutputAppendEventStandardSchemeFactory();
private static final org.apache.thrift.scheme.SchemeFactory
TUPLE_SCHEME_FACTORY = new OutputAppendEventTupleSchemeFactory();
@@ -42,6 +43,7 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
public int index; // required
public @org.apache.thrift.annotation.Nullable java.lang.String data; //
required
public @org.apache.thrift.annotation.Nullable java.lang.String appId; //
required
+ public @org.apache.thrift.annotation.Nullable java.lang.String
executionOwner; // 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 {
@@ -49,7 +51,8 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
PARAGRAPH_ID((short)2, "paragraphId"),
INDEX((short)3, "index"),
DATA((short)4, "data"),
- APP_ID((short)5, "appId");
+ APP_ID((short)5, "appId"),
+ EXECUTION_OWNER((short)6, "executionOwner");
private static final java.util.Map<java.lang.String, _Fields> byName = new
java.util.HashMap<java.lang.String, _Fields>();
@@ -75,6 +78,8 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
return DATA;
case 5: // APP_ID
return APP_ID;
+ case 6: // EXECUTION_OWNER
+ return EXECUTION_OWNER;
default:
return null;
}
@@ -131,6 +136,8 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
new
org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
tmpMap.put(_Fields.APP_ID, new
org.apache.thrift.meta_data.FieldMetaData("appId",
org.apache.thrift.TFieldRequirementType.DEFAULT,
new
org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
+ tmpMap.put(_Fields.EXECUTION_OWNER, new
org.apache.thrift.meta_data.FieldMetaData("executionOwner",
org.apache.thrift.TFieldRequirementType.DEFAULT,
+ new
org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
metaDataMap = java.util.Collections.unmodifiableMap(tmpMap);
org.apache.thrift.meta_data.FieldMetaData.addStructMetaDataMap(OutputAppendEvent.class,
metaDataMap);
}
@@ -143,7 +150,8 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
java.lang.String paragraphId,
int index,
java.lang.String data,
- java.lang.String appId)
+ java.lang.String appId,
+ java.lang.String executionOwner)
{
this();
this.noteId = noteId;
@@ -152,6 +160,7 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
setIndexIsSet(true);
this.data = data;
this.appId = appId;
+ this.executionOwner = executionOwner;
}
/**
@@ -172,6 +181,9 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
if (other.isSetAppId()) {
this.appId = other.appId;
}
+ if (other.isSetExecutionOwner()) {
+ this.executionOwner = other.executionOwner;
+ }
}
public OutputAppendEvent deepCopy() {
@@ -186,6 +198,7 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
this.index = 0;
this.data = null;
this.appId = null;
+ this.executionOwner = null;
}
@org.apache.thrift.annotation.Nullable
@@ -311,6 +324,31 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
}
}
+ @org.apache.thrift.annotation.Nullable
+ public java.lang.String getExecutionOwner() {
+ return this.executionOwner;
+ }
+
+ public OutputAppendEvent
setExecutionOwner(@org.apache.thrift.annotation.Nullable java.lang.String
executionOwner) {
+ this.executionOwner = executionOwner;
+ return this;
+ }
+
+ public void unsetExecutionOwner() {
+ this.executionOwner = null;
+ }
+
+ /** Returns true if field executionOwner is set (has been assigned a value)
and false otherwise */
+ public boolean isSetExecutionOwner() {
+ return this.executionOwner != null;
+ }
+
+ public void setExecutionOwnerIsSet(boolean value) {
+ if (!value) {
+ this.executionOwner = null;
+ }
+ }
+
public void setFieldValue(_Fields field,
@org.apache.thrift.annotation.Nullable java.lang.Object value) {
switch (field) {
case NOTE_ID:
@@ -353,6 +391,14 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
}
break;
+ case EXECUTION_OWNER:
+ if (value == null) {
+ unsetExecutionOwner();
+ } else {
+ setExecutionOwner((java.lang.String)value);
+ }
+ break;
+
}
}
@@ -374,6 +420,9 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
case APP_ID:
return getAppId();
+ case EXECUTION_OWNER:
+ return getExecutionOwner();
+
}
throw new java.lang.IllegalStateException();
}
@@ -395,6 +444,8 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
return isSetData();
case APP_ID:
return isSetAppId();
+ case EXECUTION_OWNER:
+ return isSetExecutionOwner();
}
throw new java.lang.IllegalStateException();
}
@@ -459,6 +510,15 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
return false;
}
+ boolean this_present_executionOwner = true && this.isSetExecutionOwner();
+ boolean that_present_executionOwner = true && that.isSetExecutionOwner();
+ if (this_present_executionOwner || that_present_executionOwner) {
+ if (!(this_present_executionOwner && that_present_executionOwner))
+ return false;
+ if (!this.executionOwner.equals(that.executionOwner))
+ return false;
+ }
+
return true;
}
@@ -484,6 +544,10 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
if (isSetAppId())
hashCode = hashCode * 8191 + appId.hashCode();
+ hashCode = hashCode * 8191 + ((isSetExecutionOwner()) ? 131071 : 524287);
+ if (isSetExecutionOwner())
+ hashCode = hashCode * 8191 + executionOwner.hashCode();
+
return hashCode;
}
@@ -545,6 +609,16 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
return lastComparison;
}
}
+ lastComparison =
java.lang.Boolean.valueOf(isSetExecutionOwner()).compareTo(other.isSetExecutionOwner());
+ if (lastComparison != 0) {
+ return lastComparison;
+ }
+ if (isSetExecutionOwner()) {
+ lastComparison =
org.apache.thrift.TBaseHelper.compareTo(this.executionOwner,
other.executionOwner);
+ if (lastComparison != 0) {
+ return lastComparison;
+ }
+ }
return 0;
}
@@ -601,6 +675,14 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
sb.append(this.appId);
}
first = false;
+ if (!first) sb.append(", ");
+ sb.append("executionOwner:");
+ if (this.executionOwner == null) {
+ sb.append("null");
+ } else {
+ sb.append(this.executionOwner);
+ }
+ first = false;
sb.append(")");
return sb.toString();
}
@@ -686,6 +768,14 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
}
break;
+ case 6: // EXECUTION_OWNER
+ if (schemeField.type == org.apache.thrift.protocol.TType.STRING) {
+ struct.executionOwner = iprot.readString();
+ struct.setExecutionOwnerIsSet(true);
+ } else {
+ org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
+ }
+ break;
default:
org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
}
@@ -724,6 +814,11 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
oprot.writeString(struct.appId);
oprot.writeFieldEnd();
}
+ if (struct.executionOwner != null) {
+ oprot.writeFieldBegin(EXECUTION_OWNER_FIELD_DESC);
+ oprot.writeString(struct.executionOwner);
+ oprot.writeFieldEnd();
+ }
oprot.writeFieldStop();
oprot.writeStructEnd();
}
@@ -757,7 +852,10 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
if (struct.isSetAppId()) {
optionals.set(4);
}
- oprot.writeBitSet(optionals, 5);
+ if (struct.isSetExecutionOwner()) {
+ optionals.set(5);
+ }
+ oprot.writeBitSet(optionals, 6);
if (struct.isSetNoteId()) {
oprot.writeString(struct.noteId);
}
@@ -773,12 +871,15 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
if (struct.isSetAppId()) {
oprot.writeString(struct.appId);
}
+ if (struct.isSetExecutionOwner()) {
+ oprot.writeString(struct.executionOwner);
+ }
}
@Override
public void read(org.apache.thrift.protocol.TProtocol prot,
OutputAppendEvent struct) throws org.apache.thrift.TException {
org.apache.thrift.protocol.TTupleProtocol iprot =
(org.apache.thrift.protocol.TTupleProtocol) prot;
- java.util.BitSet incoming = iprot.readBitSet(5);
+ java.util.BitSet incoming = iprot.readBitSet(6);
if (incoming.get(0)) {
struct.noteId = iprot.readString();
struct.setNoteIdIsSet(true);
@@ -799,6 +900,10 @@ public class OutputAppendEvent implements
org.apache.thrift.TBase<OutputAppendEv
struct.appId = iprot.readString();
struct.setAppIdIsSet(true);
}
+ if (incoming.get(5)) {
+ struct.executionOwner = iprot.readString();
+ struct.setExecutionOwnerIsSet(true);
+ }
}
}
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 587359fffd..7856a361e3 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,13 +24,14 @@
package org.apache.zeppelin.interpreter.thrift;
@SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-20")
public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdateAllEvent, OutputUpdateAllEvent._Fields>,
java.io.Serializable, Cloneable, Comparable<OutputUpdateAllEvent> {
private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new
org.apache.thrift.protocol.TStruct("OutputUpdateAllEvent");
private static final org.apache.thrift.protocol.TField NOTE_ID_FIELD_DESC =
new org.apache.thrift.protocol.TField("noteId",
org.apache.thrift.protocol.TType.STRING, (short)1);
private static final org.apache.thrift.protocol.TField
PARAGRAPH_ID_FIELD_DESC = new org.apache.thrift.protocol.TField("paragraphId",
org.apache.thrift.protocol.TType.STRING, (short)2);
private static final org.apache.thrift.protocol.TField MSG_FIELD_DESC = new
org.apache.thrift.protocol.TField("msg", org.apache.thrift.protocol.TType.LIST,
(short)3);
+ private static final org.apache.thrift.protocol.TField
EXECUTION_OWNER_FIELD_DESC = new
org.apache.thrift.protocol.TField("executionOwner",
org.apache.thrift.protocol.TType.STRING, (short)4);
private static final org.apache.thrift.scheme.SchemeFactory
STANDARD_SCHEME_FACTORY = new OutputUpdateAllEventStandardSchemeFactory();
private static final org.apache.thrift.scheme.SchemeFactory
TUPLE_SCHEME_FACTORY = new OutputUpdateAllEventTupleSchemeFactory();
@@ -38,12 +39,14 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
public @org.apache.thrift.annotation.Nullable java.lang.String noteId; //
required
public @org.apache.thrift.annotation.Nullable java.lang.String paragraphId;
// required
public @org.apache.thrift.annotation.Nullable
java.util.List<org.apache.zeppelin.interpreter.thrift.RemoteInterpreterResultMessage>
msg; // required
+ public @org.apache.thrift.annotation.Nullable java.lang.String
executionOwner; // 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 {
NOTE_ID((short)1, "noteId"),
PARAGRAPH_ID((short)2, "paragraphId"),
- MSG((short)3, "msg");
+ MSG((short)3, "msg"),
+ EXECUTION_OWNER((short)4, "executionOwner");
private static final java.util.Map<java.lang.String, _Fields> byName = new
java.util.HashMap<java.lang.String, _Fields>();
@@ -65,6 +68,8 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
return PARAGRAPH_ID;
case 3: // MSG
return MSG;
+ case 4: // EXECUTION_OWNER
+ return EXECUTION_OWNER;
default:
return null;
}
@@ -116,6 +121,8 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
tmpMap.put(_Fields.MSG, new
org.apache.thrift.meta_data.FieldMetaData("msg",
org.apache.thrift.TFieldRequirementType.DEFAULT,
new
org.apache.thrift.meta_data.ListMetaData(org.apache.thrift.protocol.TType.LIST,
new
org.apache.thrift.meta_data.StructMetaData(org.apache.thrift.protocol.TType.STRUCT,
org.apache.zeppelin.interpreter.thrift.RemoteInterpreterResultMessage.class))));
+ tmpMap.put(_Fields.EXECUTION_OWNER, new
org.apache.thrift.meta_data.FieldMetaData("executionOwner",
org.apache.thrift.TFieldRequirementType.DEFAULT,
+ new
org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
metaDataMap = java.util.Collections.unmodifiableMap(tmpMap);
org.apache.thrift.meta_data.FieldMetaData.addStructMetaDataMap(OutputUpdateAllEvent.class,
metaDataMap);
}
@@ -126,12 +133,14 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
public OutputUpdateAllEvent(
java.lang.String noteId,
java.lang.String paragraphId,
-
java.util.List<org.apache.zeppelin.interpreter.thrift.RemoteInterpreterResultMessage>
msg)
+
java.util.List<org.apache.zeppelin.interpreter.thrift.RemoteInterpreterResultMessage>
msg,
+ java.lang.String executionOwner)
{
this();
this.noteId = noteId;
this.paragraphId = paragraphId;
this.msg = msg;
+ this.executionOwner = executionOwner;
}
/**
@@ -151,6 +160,9 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
}
this.msg = __this__msg;
}
+ if (other.isSetExecutionOwner()) {
+ this.executionOwner = other.executionOwner;
+ }
}
public OutputUpdateAllEvent deepCopy() {
@@ -162,6 +174,7 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
this.noteId = null;
this.paragraphId = null;
this.msg = null;
+ this.executionOwner = null;
}
@org.apache.thrift.annotation.Nullable
@@ -255,6 +268,31 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
}
}
+ @org.apache.thrift.annotation.Nullable
+ public java.lang.String getExecutionOwner() {
+ return this.executionOwner;
+ }
+
+ public OutputUpdateAllEvent
setExecutionOwner(@org.apache.thrift.annotation.Nullable java.lang.String
executionOwner) {
+ this.executionOwner = executionOwner;
+ return this;
+ }
+
+ public void unsetExecutionOwner() {
+ this.executionOwner = null;
+ }
+
+ /** Returns true if field executionOwner is set (has been assigned a value)
and false otherwise */
+ public boolean isSetExecutionOwner() {
+ return this.executionOwner != null;
+ }
+
+ public void setExecutionOwnerIsSet(boolean value) {
+ if (!value) {
+ this.executionOwner = null;
+ }
+ }
+
public void setFieldValue(_Fields field,
@org.apache.thrift.annotation.Nullable java.lang.Object value) {
switch (field) {
case NOTE_ID:
@@ -281,6 +319,14 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
}
break;
+ case EXECUTION_OWNER:
+ if (value == null) {
+ unsetExecutionOwner();
+ } else {
+ setExecutionOwner((java.lang.String)value);
+ }
+ break;
+
}
}
@@ -296,6 +342,9 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
case MSG:
return getMsg();
+ case EXECUTION_OWNER:
+ return getExecutionOwner();
+
}
throw new java.lang.IllegalStateException();
}
@@ -313,6 +362,8 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
return isSetParagraphId();
case MSG:
return isSetMsg();
+ case EXECUTION_OWNER:
+ return isSetExecutionOwner();
}
throw new java.lang.IllegalStateException();
}
@@ -359,6 +410,15 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
return false;
}
+ boolean this_present_executionOwner = true && this.isSetExecutionOwner();
+ boolean that_present_executionOwner = true && that.isSetExecutionOwner();
+ if (this_present_executionOwner || that_present_executionOwner) {
+ if (!(this_present_executionOwner && that_present_executionOwner))
+ return false;
+ if (!this.executionOwner.equals(that.executionOwner))
+ return false;
+ }
+
return true;
}
@@ -378,6 +438,10 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
if (isSetMsg())
hashCode = hashCode * 8191 + msg.hashCode();
+ hashCode = hashCode * 8191 + ((isSetExecutionOwner()) ? 131071 : 524287);
+ if (isSetExecutionOwner())
+ hashCode = hashCode * 8191 + executionOwner.hashCode();
+
return hashCode;
}
@@ -419,6 +483,16 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
return lastComparison;
}
}
+ lastComparison =
java.lang.Boolean.valueOf(isSetExecutionOwner()).compareTo(other.isSetExecutionOwner());
+ if (lastComparison != 0) {
+ return lastComparison;
+ }
+ if (isSetExecutionOwner()) {
+ lastComparison =
org.apache.thrift.TBaseHelper.compareTo(this.executionOwner,
other.executionOwner);
+ if (lastComparison != 0) {
+ return lastComparison;
+ }
+ }
return 0;
}
@@ -463,6 +537,14 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
sb.append(this.msg);
}
first = false;
+ if (!first) sb.append(", ");
+ sb.append("executionOwner:");
+ if (this.executionOwner == null) {
+ sb.append("null");
+ } else {
+ sb.append(this.executionOwner);
+ }
+ first = false;
sb.append(")");
return sb.toString();
}
@@ -541,6 +623,14 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
}
break;
+ case 4: // EXECUTION_OWNER
+ if (schemeField.type == org.apache.thrift.protocol.TType.STRING) {
+ struct.executionOwner = iprot.readString();
+ struct.setExecutionOwnerIsSet(true);
+ } else {
+ org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
+ }
+ break;
default:
org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
}
@@ -578,6 +668,11 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
}
oprot.writeFieldEnd();
}
+ if (struct.executionOwner != null) {
+ oprot.writeFieldBegin(EXECUTION_OWNER_FIELD_DESC);
+ oprot.writeString(struct.executionOwner);
+ oprot.writeFieldEnd();
+ }
oprot.writeFieldStop();
oprot.writeStructEnd();
}
@@ -605,7 +700,10 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
if (struct.isSetMsg()) {
optionals.set(2);
}
- oprot.writeBitSet(optionals, 3);
+ if (struct.isSetExecutionOwner()) {
+ optionals.set(3);
+ }
+ oprot.writeBitSet(optionals, 4);
if (struct.isSetNoteId()) {
oprot.writeString(struct.noteId);
}
@@ -621,12 +719,15 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
}
}
}
+ if (struct.isSetExecutionOwner()) {
+ oprot.writeString(struct.executionOwner);
+ }
}
@Override
public void read(org.apache.thrift.protocol.TProtocol prot,
OutputUpdateAllEvent struct) throws org.apache.thrift.TException {
org.apache.thrift.protocol.TTupleProtocol iprot =
(org.apache.thrift.protocol.TTupleProtocol) prot;
- java.util.BitSet incoming = iprot.readBitSet(3);
+ java.util.BitSet incoming = iprot.readBitSet(4);
if (incoming.get(0)) {
struct.noteId = iprot.readString();
struct.setNoteIdIsSet(true);
@@ -649,6 +750,10 @@ public class OutputUpdateAllEvent implements
org.apache.thrift.TBase<OutputUpdat
}
struct.setMsgIsSet(true);
}
+ if (incoming.get(3)) {
+ struct.executionOwner = iprot.readString();
+ struct.setExecutionOwnerIsSet(true);
+ }
}
}
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 ef5d2f0d3b..41b5327e5f 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-20")
public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEvent, OutputUpdateEvent._Fields>,
java.io.Serializable, Cloneable, Comparable<OutputUpdateEvent> {
private static final org.apache.thrift.protocol.TStruct STRUCT_DESC = new
org.apache.thrift.protocol.TStruct("OutputUpdateEvent");
@@ -34,6 +34,7 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
private static final org.apache.thrift.protocol.TField TYPE_FIELD_DESC = new
org.apache.thrift.protocol.TField("type",
org.apache.thrift.protocol.TType.STRING, (short)4);
private static final org.apache.thrift.protocol.TField DATA_FIELD_DESC = new
org.apache.thrift.protocol.TField("data",
org.apache.thrift.protocol.TType.STRING, (short)5);
private static final org.apache.thrift.protocol.TField APP_ID_FIELD_DESC =
new org.apache.thrift.protocol.TField("appId",
org.apache.thrift.protocol.TType.STRING, (short)6);
+ private static final org.apache.thrift.protocol.TField
EXECUTION_OWNER_FIELD_DESC = new
org.apache.thrift.protocol.TField("executionOwner",
org.apache.thrift.protocol.TType.STRING, (short)7);
private static final org.apache.thrift.scheme.SchemeFactory
STANDARD_SCHEME_FACTORY = new OutputUpdateEventStandardSchemeFactory();
private static final org.apache.thrift.scheme.SchemeFactory
TUPLE_SCHEME_FACTORY = new OutputUpdateEventTupleSchemeFactory();
@@ -44,6 +45,7 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
public @org.apache.thrift.annotation.Nullable java.lang.String type; //
required
public @org.apache.thrift.annotation.Nullable java.lang.String data; //
required
public @org.apache.thrift.annotation.Nullable java.lang.String appId; //
required
+ public @org.apache.thrift.annotation.Nullable java.lang.String
executionOwner; // 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 {
@@ -52,7 +54,8 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
INDEX((short)3, "index"),
TYPE((short)4, "type"),
DATA((short)5, "data"),
- APP_ID((short)6, "appId");
+ APP_ID((short)6, "appId"),
+ EXECUTION_OWNER((short)7, "executionOwner");
private static final java.util.Map<java.lang.String, _Fields> byName = new
java.util.HashMap<java.lang.String, _Fields>();
@@ -80,6 +83,8 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
return DATA;
case 6: // APP_ID
return APP_ID;
+ case 7: // EXECUTION_OWNER
+ return EXECUTION_OWNER;
default:
return null;
}
@@ -138,6 +143,8 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
new
org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
tmpMap.put(_Fields.APP_ID, new
org.apache.thrift.meta_data.FieldMetaData("appId",
org.apache.thrift.TFieldRequirementType.DEFAULT,
new
org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
+ tmpMap.put(_Fields.EXECUTION_OWNER, new
org.apache.thrift.meta_data.FieldMetaData("executionOwner",
org.apache.thrift.TFieldRequirementType.DEFAULT,
+ new
org.apache.thrift.meta_data.FieldValueMetaData(org.apache.thrift.protocol.TType.STRING)));
metaDataMap = java.util.Collections.unmodifiableMap(tmpMap);
org.apache.thrift.meta_data.FieldMetaData.addStructMetaDataMap(OutputUpdateEvent.class,
metaDataMap);
}
@@ -151,7 +158,8 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
int index,
java.lang.String type,
java.lang.String data,
- java.lang.String appId)
+ java.lang.String appId,
+ java.lang.String executionOwner)
{
this();
this.noteId = noteId;
@@ -161,6 +169,7 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
this.type = type;
this.data = data;
this.appId = appId;
+ this.executionOwner = executionOwner;
}
/**
@@ -184,6 +193,9 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
if (other.isSetAppId()) {
this.appId = other.appId;
}
+ if (other.isSetExecutionOwner()) {
+ this.executionOwner = other.executionOwner;
+ }
}
public OutputUpdateEvent deepCopy() {
@@ -199,6 +211,7 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
this.type = null;
this.data = null;
this.appId = null;
+ this.executionOwner = null;
}
@org.apache.thrift.annotation.Nullable
@@ -349,6 +362,31 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
}
}
+ @org.apache.thrift.annotation.Nullable
+ public java.lang.String getExecutionOwner() {
+ return this.executionOwner;
+ }
+
+ public OutputUpdateEvent
setExecutionOwner(@org.apache.thrift.annotation.Nullable java.lang.String
executionOwner) {
+ this.executionOwner = executionOwner;
+ return this;
+ }
+
+ public void unsetExecutionOwner() {
+ this.executionOwner = null;
+ }
+
+ /** Returns true if field executionOwner is set (has been assigned a value)
and false otherwise */
+ public boolean isSetExecutionOwner() {
+ return this.executionOwner != null;
+ }
+
+ public void setExecutionOwnerIsSet(boolean value) {
+ if (!value) {
+ this.executionOwner = null;
+ }
+ }
+
public void setFieldValue(_Fields field,
@org.apache.thrift.annotation.Nullable java.lang.Object value) {
switch (field) {
case NOTE_ID:
@@ -399,6 +437,14 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
}
break;
+ case EXECUTION_OWNER:
+ if (value == null) {
+ unsetExecutionOwner();
+ } else {
+ setExecutionOwner((java.lang.String)value);
+ }
+ break;
+
}
}
@@ -423,6 +469,9 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
case APP_ID:
return getAppId();
+ case EXECUTION_OWNER:
+ return getExecutionOwner();
+
}
throw new java.lang.IllegalStateException();
}
@@ -446,6 +495,8 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
return isSetData();
case APP_ID:
return isSetAppId();
+ case EXECUTION_OWNER:
+ return isSetExecutionOwner();
}
throw new java.lang.IllegalStateException();
}
@@ -519,6 +570,15 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
return false;
}
+ boolean this_present_executionOwner = true && this.isSetExecutionOwner();
+ boolean that_present_executionOwner = true && that.isSetExecutionOwner();
+ if (this_present_executionOwner || that_present_executionOwner) {
+ if (!(this_present_executionOwner && that_present_executionOwner))
+ return false;
+ if (!this.executionOwner.equals(that.executionOwner))
+ return false;
+ }
+
return true;
}
@@ -548,6 +608,10 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
if (isSetAppId())
hashCode = hashCode * 8191 + appId.hashCode();
+ hashCode = hashCode * 8191 + ((isSetExecutionOwner()) ? 131071 : 524287);
+ if (isSetExecutionOwner())
+ hashCode = hashCode * 8191 + executionOwner.hashCode();
+
return hashCode;
}
@@ -619,6 +683,16 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
return lastComparison;
}
}
+ lastComparison =
java.lang.Boolean.valueOf(isSetExecutionOwner()).compareTo(other.isSetExecutionOwner());
+ if (lastComparison != 0) {
+ return lastComparison;
+ }
+ if (isSetExecutionOwner()) {
+ lastComparison =
org.apache.thrift.TBaseHelper.compareTo(this.executionOwner,
other.executionOwner);
+ if (lastComparison != 0) {
+ return lastComparison;
+ }
+ }
return 0;
}
@@ -683,6 +757,14 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
sb.append(this.appId);
}
first = false;
+ if (!first) sb.append(", ");
+ sb.append("executionOwner:");
+ if (this.executionOwner == null) {
+ sb.append("null");
+ } else {
+ sb.append(this.executionOwner);
+ }
+ first = false;
sb.append(")");
return sb.toString();
}
@@ -776,6 +858,14 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
}
break;
+ case 7: // EXECUTION_OWNER
+ if (schemeField.type == org.apache.thrift.protocol.TType.STRING) {
+ struct.executionOwner = iprot.readString();
+ struct.setExecutionOwnerIsSet(true);
+ } else {
+ org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
+ }
+ break;
default:
org.apache.thrift.protocol.TProtocolUtil.skip(iprot,
schemeField.type);
}
@@ -819,6 +909,11 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
oprot.writeString(struct.appId);
oprot.writeFieldEnd();
}
+ if (struct.executionOwner != null) {
+ oprot.writeFieldBegin(EXECUTION_OWNER_FIELD_DESC);
+ oprot.writeString(struct.executionOwner);
+ oprot.writeFieldEnd();
+ }
oprot.writeFieldStop();
oprot.writeStructEnd();
}
@@ -855,7 +950,10 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
if (struct.isSetAppId()) {
optionals.set(5);
}
- oprot.writeBitSet(optionals, 6);
+ if (struct.isSetExecutionOwner()) {
+ optionals.set(6);
+ }
+ oprot.writeBitSet(optionals, 7);
if (struct.isSetNoteId()) {
oprot.writeString(struct.noteId);
}
@@ -874,12 +972,15 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
if (struct.isSetAppId()) {
oprot.writeString(struct.appId);
}
+ if (struct.isSetExecutionOwner()) {
+ oprot.writeString(struct.executionOwner);
+ }
}
@Override
public void read(org.apache.thrift.protocol.TProtocol prot,
OutputUpdateEvent struct) throws org.apache.thrift.TException {
org.apache.thrift.protocol.TTupleProtocol iprot =
(org.apache.thrift.protocol.TTupleProtocol) prot;
- java.util.BitSet incoming = iprot.readBitSet(6);
+ java.util.BitSet incoming = iprot.readBitSet(7);
if (incoming.get(0)) {
struct.noteId = iprot.readString();
struct.setNoteIdIsSet(true);
@@ -904,6 +1005,10 @@ public class OutputUpdateEvent implements
org.apache.thrift.TBase<OutputUpdateEv
struct.appId = iprot.readString();
struct.setAppIdIsSet(true);
}
+ if (incoming.get(6)) {
+ struct.executionOwner = iprot.readString();
+ struct.setExecutionOwnerIsSet(true);
+ }
}
}
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 083d800351..0d08ce0816 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class ParagraphInfo implements org.apache.thrift.TBase<ParagraphInfo,
ParagraphInfo._Fields>, java.io.Serializable, Cloneable,
Comparable<ParagraphInfo> {
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 2c3667eb8f..18f7fcf755 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RegisterInfo implements org.apache.thrift.TBase<RegisterInfo,
RegisterInfo._Fields>, java.io.Serializable, Cloneable,
Comparable<RegisterInfo> {
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 8aa5bc2e6b..d3d4114eee 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RemoteApplicationResult implements
org.apache.thrift.TBase<RemoteApplicationResult,
RemoteApplicationResult._Fields>, java.io.Serializable, Cloneable,
Comparable<RemoteApplicationResult> {
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 a6243c83b9..bd7ee2d868 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RemoteInterpreterContext implements
org.apache.thrift.TBase<RemoteInterpreterContext,
RemoteInterpreterContext._Fields>, java.io.Serializable, Cloneable,
Comparable<RemoteInterpreterContext> {
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 eb22200e70..b4499f3c91 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RemoteInterpreterEvent implements
org.apache.thrift.TBase<RemoteInterpreterEvent,
RemoteInterpreterEvent._Fields>, java.io.Serializable, Cloneable,
Comparable<RemoteInterpreterEvent> {
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 9a91dbd059..0ccd4f5422 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,7 +24,7 @@
package org.apache.zeppelin.interpreter.thrift;
@SuppressWarnings({"cast", "rawtypes", "serial", "unchecked", "unused"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RemoteInterpreterEventService {
public interface Iface {
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 292c519253..5d8244ae89 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;
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
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 4fce5b0b9f..00ceeb62e5 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RemoteInterpreterResult implements
org.apache.thrift.TBase<RemoteInterpreterResult,
RemoteInterpreterResult._Fields>, java.io.Serializable, Cloneable,
Comparable<RemoteInterpreterResult> {
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 6ea6fa6aa7..3a00fe14f6 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RemoteInterpreterResultMessage implements
org.apache.thrift.TBase<RemoteInterpreterResultMessage,
RemoteInterpreterResultMessage._Fields>, java.io.Serializable, Cloneable,
Comparable<RemoteInterpreterResultMessage> {
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 4cde441cdc..a6ad9d5f97 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
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 47edc59883..d9272dd4cc 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class RunParagraphsEvent implements
org.apache.thrift.TBase<RunParagraphsEvent, RunParagraphsEvent._Fields>,
java.io.Serializable, Cloneable, Comparable<RunParagraphsEvent> {
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 1882f45f87..469b9d8ffe 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class ServiceException extends org.apache.thrift.TException implements
org.apache.thrift.TBase<ServiceException, ServiceException._Fields>,
java.io.Serializable, Cloneable, Comparable<ServiceException> {
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 a622c0d6f3..3dd9722651 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"})
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2021-03-09")
[email protected](value = "Autogenerated by Thrift Compiler
(0.13.0)", date = "2026-09-13")
public class WebUrlInfo implements org.apache.thrift.TBase<WebUrlInfo,
WebUrlInfo._Fields>, java.io.Serializable, Cloneable, Comparable<WebUrlInfo> {
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 ad08821f41..c7c6e039f0 100644
--- a/zeppelin-interpreter/src/main/thrift/RemoteInterpreterEventService.thrift
+++ b/zeppelin-interpreter/src/main/thrift/RemoteInterpreterEventService.thrift
@@ -36,7 +36,8 @@ struct OutputAppendEvent {
2: string paragraphId,
3: i32 index,
4: string data,
- 5: string appId
+ 5: string appId,
+ 6: string executionOwner
}
struct OutputUpdateEvent {
@@ -45,13 +46,15 @@ struct OutputUpdateEvent {
3: i32 index,
4: string type,
5: string data,
- 6: string appId
+ 6: string appId,
+ 7: string executionOwner
}
struct OutputUpdateAllEvent {
1: string noteId,
2: string paragraphId,
3: list<RemoteInterpreterService.RemoteInterpreterResultMessage> msg,
+ 4: string executionOwner
}
struct RunParagraphsEvent {
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 c8f8886e6a..5504283308 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
@@ -218,8 +218,8 @@ public class RemoteInterpreterEventServer implements
RemoteInterpreterEventServi
@Override
public void appendOutput(OutputAppendEvent event) throws
InterpreterRPCException, TException {
if (event.getAppId() == null) {
- runner.appendBuffer(
- event.getNoteId(), event.getParagraphId(), event.getIndex(),
event.getData());
+ runner.appendBuffer(event.getNoteId(), event.getParagraphId(),
event.getIndex(),
+ event.getExecutionOwner(), event.getData());
} else {
appListener.onOutputAppend(event.getNoteId(), event.getParagraphId(),
event.getIndex(),
event.getAppId(), event.getData());
@@ -230,7 +230,8 @@ public class RemoteInterpreterEventServer implements
RemoteInterpreterEventServi
public void updateOutput(OutputUpdateEvent event) throws
InterpreterRPCException, TException {
if (event.getAppId() == null) {
runner.updateBuffer(event.getNoteId(), event.getParagraphId(),
event.getIndex(),
- InterpreterResult.Type.valueOf(event.getType()), event.getData());
+ event.getExecutionOwner(),
InterpreterResult.Type.valueOf(event.getType()),
+ event.getData());
// Complete replacements before the interpreter can publish its terminal
result.
runner.run();
} else {
@@ -244,11 +245,13 @@ public class RemoteInterpreterEventServer implements
RemoteInterpreterEventServi
synchronized (runner) {
// Finish earlier output before the clear; keep replacements ahead of
the next drain.
runner.run();
- listener.onOutputClear(event.getNoteId(), event.getParagraphId());
+ listener.onParagraphOutputClear(
+ event.getNoteId(), event.getParagraphId(),
event.getExecutionOwner());
for (int i = 0; i < event.getMsg().size(); i++) {
RemoteInterpreterResultMessage msg = event.getMsg().get(i);
- listener.onOutputUpdated(event.getNoteId(), event.getParagraphId(), i,
- InterpreterResult.Type.valueOf(msg.getType()), msg.getData());
+ listener.onParagraphOutputUpdated(event.getNoteId(),
event.getParagraphId(), i,
+ event.getExecutionOwner(),
InterpreterResult.Type.valueOf(msg.getType()),
+ msg.getData());
}
}
}
diff --git
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputBuffer.java
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputBuffer.java
index b13940496f..97a13338b2 100644
---
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputBuffer.java
+++
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputBuffer.java
@@ -26,12 +26,15 @@ public class AppendOutputBuffer {
private String noteId;
private String paragraphId;
private int index;
+ private String executionOwner;
private String data;
- public AppendOutputBuffer(String noteId, String paragraphId, int index,
String data) {
+ public AppendOutputBuffer(String noteId, String paragraphId, int index,
String executionOwner,
+ String data) {
this.noteId = noteId;
this.paragraphId = paragraphId;
this.index = index;
+ this.executionOwner = executionOwner;
this.data = data;
}
@@ -47,6 +50,10 @@ public class AppendOutputBuffer {
return index;
}
+ public String getExecutionOwner() {
+ return executionOwner;
+ }
+
public String getData() {
return data;
}
diff --git
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunner.java
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunner.java
index 8bc063af10..4865b2464d 100644
---
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunner.java
+++
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunner.java
@@ -21,11 +21,12 @@ import org.apache.zeppelin.interpreter.InterpreterResult;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.util.HashMap;
+import java.util.LinkedHashMap;
import java.util.LinkedList;
import java.util.List;
import java.util.Map;
import java.util.Map.Entry;
+import java.util.Objects;
import java.util.concurrent.BlockingQueue;
import java.util.concurrent.LinkedBlockingQueue;
@@ -52,7 +53,7 @@ public class AppendOutputRunner implements Runnable {
@Override
public synchronized void run() {
- Map<String, StringBuilder> stringBufferMap = new HashMap<>();
+ Map<AppendKey, StringBuilder> stringBufferMap = new LinkedHashMap<>();
List<AppendOutputBuffer> list = new LinkedList<>();
queue.drainTo(list);
@@ -67,8 +68,8 @@ public class AppendOutputRunner implements Runnable {
sizeProcessed += flushAppendBuffers(stringBufferMap);
UpdateOutputBuffer update = (UpdateOutputBuffer) buffer;
try {
- listener.onOutputUpdated(update.getNoteId(),
update.getParagraphId(), update.getIndex(),
- update.getType(), update.getData());
+ listener.onParagraphOutputUpdated(update.getNoteId(),
update.getParagraphId(),
+ update.getIndex(), update.getExecutionOwner(), update.getType(),
update.getData());
} catch (RuntimeException e) {
// A stale callback must not abort another paragraph's synchronous
drain.
LOGGER.warn("Failed to update output for note {} paragraph {}",
@@ -77,16 +78,11 @@ public class AppendOutputRunner implements Runnable {
continue;
}
- String noteId = buffer.getNoteId();
- String paragraphId = buffer.getParagraphId();
- int index = buffer.getIndex();
- String stringBufferKey = noteId + ":" + paragraphId + ":" + index;
+ AppendKey key = new AppendKey(buffer.getNoteId(),
buffer.getParagraphId(),
+ buffer.getIndex(), buffer.getExecutionOwner());
- StringBuilder builder = stringBufferMap.containsKey(stringBufferKey) ?
- stringBufferMap.get(stringBufferKey) : new StringBuilder();
-
- builder.append(buffer.getData());
- stringBufferMap.put(stringBufferKey, builder);
+ stringBufferMap.computeIfAbsent(key, unused -> new StringBuilder())
+ .append(buffer.getData());
}
sizeProcessed += flushAppendBuffers(stringBufferMap);
Long processingTime = System.currentTimeMillis() - processingStartTime;
@@ -104,31 +100,78 @@ public class AppendOutputRunner implements Runnable {
}
}
- private long flushAppendBuffers(Map<String, StringBuilder> stringBufferMap) {
+ private long flushAppendBuffers(Map<AppendKey, StringBuilder>
stringBufferMap) {
long sizeProcessed = 0;
- for (Entry<String, StringBuilder> stringBufferMapEntry :
stringBufferMap.entrySet()) {
- String stringBufferKey = stringBufferMapEntry.getKey();
+ for (Entry<AppendKey, StringBuilder> stringBufferMapEntry :
stringBufferMap.entrySet()) {
+ AppendKey key = stringBufferMapEntry.getKey();
StringBuilder buffer = stringBufferMapEntry.getValue();
sizeProcessed += buffer.length();
try {
- String[] keys = stringBufferKey.split(":");
- listener.onOutputAppend(keys[0], keys[1], Integer.parseInt(keys[2]),
buffer.toString());
+ listener.onParagraphOutputAppend(key.noteId, key.paragraphId,
key.index,
+ key.executionOwner, buffer.toString());
} catch (RuntimeException e) {
// One stale append must not abort another paragraph's synchronous
drain.
- LOGGER.warn("Failed to append output for {}", stringBufferKey, e);
+ LOGGER.warn("Failed to append output for {}", key, e);
}
}
stringBufferMap.clear();
return sizeProcessed;
}
- public void appendBuffer(String noteId, String paragraphId, int index,
String outputToAppend) {
- queue.offer(new AppendOutputBuffer(noteId, paragraphId, index,
outputToAppend));
+ /**
+ * Identifies one stream of appended output. An owner name can contain any
character, so the
+ * parts are kept separate instead of being joined into a delimited string.
+ */
+ private static final class AppendKey {
+ private final String noteId;
+ private final String paragraphId;
+ private final int index;
+ private final String executionOwner;
+
+ private AppendKey(String noteId, String paragraphId, int index, String
executionOwner) {
+ this.noteId = noteId;
+ this.paragraphId = paragraphId;
+ this.index = index;
+ this.executionOwner = executionOwner;
+ }
+
+ @Override
+ public boolean equals(Object o) {
+ if (this == o) {
+ return true;
+ }
+ if (!(o instanceof AppendKey)) {
+ return false;
+ }
+ AppendKey other = (AppendKey) o;
+ return index == other.index
+ && Objects.equals(noteId, other.noteId)
+ && Objects.equals(paragraphId, other.paragraphId)
+ && Objects.equals(executionOwner, other.executionOwner);
+ }
+
+ @Override
+ public int hashCode() {
+ return Objects.hash(noteId, paragraphId, index, executionOwner);
+ }
+
+ @Override
+ public String toString() {
+ return "note " + noteId + " paragraph " + paragraphId + " index " + index
+ + " executionOwner " + executionOwner;
+ }
+ }
+
+ public void appendBuffer(String noteId, String paragraphId, int index,
String executionOwner,
+ String outputToAppend) {
+ queue.offer(
+ new AppendOutputBuffer(noteId, paragraphId, index, executionOwner,
outputToAppend));
}
/** Enqueues a replacement; callers needing completion must also invoke
run(). */
- public void updateBuffer(String noteId, String paragraphId, int index,
+ public void updateBuffer(String noteId, String paragraphId, int index,
String executionOwner,
InterpreterResult.Type type, String output) {
- queue.offer(new UpdateOutputBuffer(noteId, paragraphId, index, type,
output));
+ queue.offer(
+ new UpdateOutputBuffer(noteId, paragraphId, index, executionOwner,
type, output));
}
}
diff --git
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java
index e6876c72ea..7d32938078 100644
---
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java
+++
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/RemoteInterpreterProcessListener.java
@@ -33,27 +33,31 @@ public interface RemoteInterpreterProcessListener {
* @param noteId
* @param paragraphId
* @param index
+ * @param executionOwner null when the interpreter does not report one
* @param output
*/
- void onOutputAppend(String noteId, String paragraphId, int index, String
output);
+ void onParagraphOutputAppend(String noteId, String paragraphId, int index,
+ String executionOwner, String output);
/**
* Invoked when the whole output is updated
* @param noteId
* @param paragraphId
* @param index
+ * @param executionOwner null when the interpreter does not report one
* @param type
* @param output
*/
- void onOutputUpdated(
- String noteId, String paragraphId, int index, InterpreterResult.Type
type, String output);
+ void onParagraphOutputUpdated(String noteId, String paragraphId, int index,
+ String executionOwner, InterpreterResult.Type type, String output);
/**
* Invoked when output is cleared.
* @param noteId
* @param paragraphId
+ * @param executionOwner null when the interpreter does not report one
*/
- void onOutputClear(String noteId, String paragraphId);
+ void onParagraphOutputClear(String noteId, String paragraphId, String
executionOwner);
/**
* Run paragraphs, paragraphs can be specified via indices(paragraphIndices)
or ids(paragraphIds)
diff --git
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/UpdateOutputBuffer.java
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/UpdateOutputBuffer.java
index 15de2d5092..1243461d17 100644
---
a/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/UpdateOutputBuffer.java
+++
b/zeppelin-server/src/main/java/org/apache/zeppelin/interpreter/remote/UpdateOutputBuffer.java
@@ -28,9 +28,9 @@ public class UpdateOutputBuffer extends AppendOutputBuffer {
private final InterpreterResult.Type type;
- public UpdateOutputBuffer(String noteId, String paragraphId, int index,
+ public UpdateOutputBuffer(String noteId, String paragraphId, int index,
String executionOwner,
InterpreterResult.Type type, String data) {
- super(noteId, paragraphId, index, data);
+ super(noteId, paragraphId, index, executionOwner, data);
this.type = type;
}
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 4d31a06558..8e568e1217 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
@@ -1758,7 +1758,8 @@ public class NotebookServer implements
AngularObjectRegistryListener,
* @param output output to append
*/
@Override
- public void onOutputAppend(String noteId, String paragraphId, int index,
String output) {
+ public void onParagraphOutputAppend(String noteId, String paragraphId, int
index,
+ String executionOwner, String output) {
if (!sendParagraphStatusToFrontend()) {
return;
}
@@ -1772,13 +1773,17 @@ public class NotebookServer implements
AngularObjectRegistryListener,
if (note == null) {
LOGGER.warn("Note {} not found", noteId);
} else if (!note.isPersonalizedMode()) {
- // Streaming events do not identify the user that owns the execution.
connectionManager.broadcast(noteId, msg);
+ } else if (executionOwner != null) {
+ connectionManager.multicastToUser(executionOwner, msg);
+ } else {
+ LOGGER.debug("Dropping ownerless personalized output for note {}
paragraph {}",
+ noteId, paragraphId);
}
return null;
});
} catch (IOException e) {
- LOGGER.warn("Fail to call onOutputAppend", e);
+ LOGGER.warn("Fail to call onParagraphOutputAppend", e);
}
}
@@ -1788,8 +1793,9 @@ public class NotebookServer implements
AngularObjectRegistryListener,
* @param output output to update (replace)
*/
@Override
- public void onOutputUpdated(String noteId, String paragraphId, int index,
- InterpreterResult.Type type, String output) {
+ public void onParagraphOutputUpdated(String noteId, String paragraphId, int
index,
+ String executionOwner,
InterpreterResult.Type type,
+ String output) {
if (!sendParagraphStatusToFrontend()) {
return;
}
@@ -1807,10 +1813,19 @@ public class NotebookServer implements
AngularObjectRegistryListener,
return null;
}
if (note.isPersonalizedMode()) {
- // Streaming events carry no owner. The shared outputBuffer is
what checkpointOutput
- // saves as the shared result and what other users' paragraphs are
cloned from, so
- // one user's output must not be written there. Personalized
clients get their
- // user-specific terminal snapshot instead.
+ if (executionOwner == null) {
+ LOGGER.debug("Dropping ownerless personalized output for note {}
paragraph {}",
+ noteId, paragraphId);
+ return null;
+ }
+ // The shared outputBuffer is both what checkpointOutput saves and
what new users'
+ // copies are cloned from, so one user's output goes to that
user's own copy.
+ Paragraph userParagraph =
+
note.getParagraph(paragraphId).getUserParagraphMap().get(executionOwner);
+ if (userParagraph != null) {
+ userParagraph.updateOutputBuffer(index, type, output);
+ }
+ connectionManager.multicastToUser(executionOwner, msg);
return null;
}
note.getParagraph(paragraphId).updateOutputBuffer(index, type,
output);
@@ -1818,7 +1833,7 @@ public class NotebookServer implements
AngularObjectRegistryListener,
return null;
});
} catch (IOException e) {
- LOGGER.warn("Fail to call onOutputUpdated", e);
+ LOGGER.warn("Fail to call onParagraphOutputUpdated", e);
}
}
@@ -1826,7 +1841,7 @@ public class NotebookServer implements
AngularObjectRegistryListener,
* This callback is for the paragraph that runs on ZeppelinServer.
*/
@Override
- public void onOutputClear(String noteId, String paragraphId) {
+ public void onParagraphOutputClear(String noteId, String paragraphId, String
executionOwner) {
if (!sendParagraphStatusToFrontend()) {
return;
}
@@ -1840,7 +1855,19 @@ public class NotebookServer implements
AngularObjectRegistryListener,
return null;
}
if (note.isPersonalizedMode()) {
- // Streaming events carry no owner, so they must not mutate shared
paragraph state.
+ if (executionOwner == null) {
+ LOGGER.debug("Dropping ownerless personalized clear for note {}
paragraph {}",
+ noteId, paragraphId);
+ return null;
+ }
+ // Clearing the shared paragraph would discard output the other
users still own.
+ if (note.getParagraph(paragraphId).getUserParagraphMap()
+ .containsKey(executionOwner)) {
+ Paragraph userParagraph =
+ note.clearPersonalizedParagraphOutput(paragraphId,
executionOwner);
+ connectionManager.multicastToUser(executionOwner, new
Message(OP.PARAGRAPH)
+ .withMsgId(MSG_ID_NOT_DEFINED).put("paragraph",
userParagraph));
+ }
return null;
}
note.clearParagraphOutput(paragraphId);
@@ -1850,7 +1877,7 @@ public class NotebookServer implements
AngularObjectRegistryListener,
});
} catch (IOException e) {
- LOGGER.warn("Fail to call onOutputClear", e);
+ LOGGER.warn("Fail to call onParagraphOutputClear", e);
}
}
diff --git
a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerTest.java
b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerTest.java
index db73b10a42..8fac268b9e 100644
---
a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerTest.java
+++
b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/RemoteInterpreterEventServerTest.java
@@ -63,9 +63,10 @@ public class RemoteInterpreterEventServerTest {
RemoteInterpreterEventServer server =
serverWithRunner(listener, new AppendOutputRunner(listener));
try {
- server.updateOutput(new OutputUpdateEvent("note", "para", 0, "TEXT",
"final", null));
+ server.updateOutput(new OutputUpdateEvent("note", "para", 0, "TEXT",
"final", null, null));
// A caller may publish terminal status as soon as the RPC returns.
- verify(listener).onOutputUpdated("note", "para", 0,
InterpreterResult.Type.TEXT, "final");
+ verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.TEXT, "final");
} finally {
server.stop();
}
@@ -77,10 +78,10 @@ public class RemoteInterpreterEventServerTest {
RemoteInterpreterEventServer server =
serverWithRunner(listener, new AppendOutputRunner(listener));
try {
- server.appendOutput(new OutputAppendEvent("note", "para", 0, "pending",
null));
+ server.appendOutput(new OutputAppendEvent("note", "para", 0, "pending",
null, null));
server.checkpointOutput("note", "para");
InOrder order = inOrder(listener);
- order.verify(listener).onOutputAppend("note", "para", 0, "pending");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, null,
"pending");
order.verify(listener).checkpointOutput("note", "para");
} finally {
server.stop();
@@ -93,19 +94,19 @@ public class RemoteInterpreterEventServerTest {
AppendOutputRunner runner = new AppendOutputRunner(listener);
RemoteInterpreterEventServer server = serverWithRunner(listener, runner);
try {
- runner.appendBuffer("note", "para", 0, "old");
+ runner.appendBuffer("note", "para", 0, null, "old");
server.updateAllOutput(new OutputUpdateAllEvent("note", "para",
Collections.singletonList(
- new RemoteInterpreterResultMessage("HTML", "replacement"))));
- verify(listener).onOutputUpdated("note", "para", 0,
- InterpreterResult.Type.HTML, "replacement");
- runner.appendBuffer("note", "para", 0, "new");
+ new RemoteInterpreterResultMessage("HTML", "replacement")), null));
+ verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.HTML, "replacement");
+ runner.appendBuffer("note", "para", 0, null, "new");
runner.run();
InOrder order = inOrder(listener);
- order.verify(listener).onOutputAppend("note", "para", 0, "old");
- order.verify(listener).onOutputClear("note", "para");
- order.verify(listener).onOutputUpdated("note", "para", 0,
- InterpreterResult.Type.HTML, "replacement");
- order.verify(listener).onOutputAppend("note", "para", 0, "new");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, null,
"old");
+ order.verify(listener).onParagraphOutputClear("note", "para", null);
+ order.verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.HTML, "replacement");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, null,
"new");
} finally {
server.stop();
}
@@ -123,18 +124,18 @@ public class RemoteInterpreterEventServerTest {
entered.countDown();
assertTrue(release.await(5, TimeUnit.SECONDS));
return null;
- }).when(listener).onOutputAppend("note", "para", 0, "old");
+ }).when(listener).onParagraphOutputAppend("note", "para", 0, null, "old");
ExecutorService executor = Executors.newFixedThreadPool(2);
try {
- runner.appendBuffer("note", "para", 0, "old");
+ runner.appendBuffer("note", "para", 0, null, "old");
Future<?> first = executor.submit(runner);
assertTrue(entered.await(5, TimeUnit.SECONDS));
Future<?> update = executor.submit(() -> {
updateStarted.countDown();
server.updateAllOutput(new OutputUpdateAllEvent("note", "para",
Collections.singletonList(
- new RemoteInterpreterResultMessage("HTML", "replacement"))));
- verify(listener).onOutputUpdated("note", "para", 0,
- InterpreterResult.Type.HTML, "replacement");
+ new RemoteInterpreterResultMessage("HTML", "replacement")), null));
+ verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.HTML, "replacement");
return null;
});
assertTrue(updateStarted.await(5, TimeUnit.SECONDS));
@@ -143,10 +144,10 @@ public class RemoteInterpreterEventServerTest {
first.get(5, TimeUnit.SECONDS);
update.get(5, TimeUnit.SECONDS);
InOrder order = inOrder(listener);
- order.verify(listener).onOutputAppend("note", "para", 0, "old");
- order.verify(listener).onOutputClear("note", "para");
- order.verify(listener).onOutputUpdated("note", "para", 0,
- InterpreterResult.Type.HTML, "replacement");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, null,
"old");
+ order.verify(listener).onParagraphOutputClear("note", "para", null);
+ order.verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.HTML, "replacement");
} finally {
release.countDown();
executor.shutdownNow();
diff --git
a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunnerTest.java
b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunnerTest.java
index 2d6e08beba..786c70e23c 100644
---
a/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunnerTest.java
+++
b/zeppelin-server/src/test/java/org/apache/zeppelin/interpreter/remote/AppendOutputRunnerTest.java
@@ -47,6 +47,7 @@ import static org.mockito.Mockito.doThrow;
import static org.junit.jupiter.api.Assertions.fail;
import static org.mockito.ArgumentMatchers.any;
import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.ArgumentMatchers.isNull;
import static org.mockito.Mockito.atMost;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.inOrder;
@@ -79,8 +80,9 @@ class AppendOutputRunnerTest {
String[][] buffer = {{"note", "para", "data\n"}};
loopForCompletingEvents(listener, 1, buffer);
- verify(listener, times(1)).onOutputAppend(any(String.class),
any(String.class), anyInt(), any(String.class));
- verify(listener, times(1)).onOutputAppend("note", "para", 0, "data\n");
+ verify(listener, times(1)).onParagraphOutputAppend(
+ any(String.class), any(String.class), anyInt(), isNull(),
any(String.class));
+ verify(listener, times(1)).onParagraphOutputAppend("note", "para", 0,
null, "data\n");
}
@Test
@@ -95,26 +97,78 @@ class AppendOutputRunnerTest {
};
loopForCompletingEvents(listener, 1, buffer);
- verify(listener, times(1)).onOutputAppend(any(String.class),
any(String.class), anyInt(), any(String.class));
- verify(listener, times(1)).onOutputAppend(note1, para1, 0,
"data1\ndata2\ndata3\n");
+ verify(listener, times(1)).onParagraphOutputAppend(
+ any(String.class), any(String.class), anyInt(), isNull(),
any(String.class));
+ verify(listener, times(1)).onParagraphOutputAppend(
+ note1, para1, 0, null, "data1\ndata2\ndata3\n");
+ }
+
+ // A paragraph in shared mode is one Job, so one execution -- one user --
produces the output
+ // for a given index at a time, and keying by user must leave that batching
alone.
+ @Test
+ void appendsFromOneExecutionAreStillBatchedIntoOneChunk() {
+ RemoteInterpreterProcessListener listener =
mock(RemoteInterpreterProcessListener.class);
+ AppendOutputRunner runner = new AppendOutputRunner(listener);
+ runner.appendBuffer("note", "para", 0, "user1", "line-1\n");
+ runner.appendBuffer("note", "para", 0, "user1", "line-2\n");
+
+ runner.run();
+
+ verify(listener, times(1)).onParagraphOutputAppend(
+ "note", "para", 0, "user1", "line-1\nline-2\n");
+ }
+
+ // Two executions only overlap in personalized mode, where each user runs
their own copy.
+ // A merged chunk would have no single owner and could not be routed to
either user.
+ @Test
+ void appendsFromDifferentExecutionsAreNotBatchedTogether() {
+ RemoteInterpreterProcessListener listener =
mock(RemoteInterpreterProcessListener.class);
+ AppendOutputRunner runner = new AppendOutputRunner(listener);
+ runner.appendBuffer("note", "para", 0, "user1", "mine\n");
+ runner.appendBuffer("note", "para", 0, "user2", "theirs\n");
+
+ runner.run();
+
+ verify(listener, times(1)).onParagraphOutputAppend("note", "para", 0,
"user1", "mine\n");
+ verify(listener, times(1)).onParagraphOutputAppend("note", "para", 0,
"user2", "theirs\n");
+ }
+
+ // Keying by user splits the buffer, so the per-owner ordering the shared
queue guarantees
+ // must still hold once another user's output is interleaved.
+ @Test
+ void updatesDoNotOvertakeQueuedAppendsOfTheSameOwner() {
+ RemoteInterpreterProcessListener listener =
mock(RemoteInterpreterProcessListener.class);
+ AppendOutputRunner runner = new AppendOutputRunner(listener);
+ runner.appendBuffer("note", "para", 0, "owner", "before\n");
+ runner.appendBuffer("note", "para", 0, "other", "theirs\n");
+ runner.updateBuffer("note", "para", 0, "owner",
InterpreterResult.Type.TEXT, "replacement\n");
+ runner.appendBuffer("note", "para", 0, "owner", "after\n");
+
+ runner.run();
+
+ InOrder order = inOrder(listener);
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, "owner",
"before\n");
+ order.verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, "owner", InterpreterResult.Type.TEXT,
"replacement\n");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, "owner",
"after\n");
}
@Test
void testUpdateDoesNotOvertakeQueuedAppend() {
RemoteInterpreterProcessListener listener =
mock(RemoteInterpreterProcessListener.class);
AppendOutputRunner runner = new AppendOutputRunner(listener);
- runner.appendBuffer("note", "para", 0, "before-1\n");
- runner.appendBuffer("note", "para", 0, "before-2\n");
- runner.updateBuffer("note", "para", 0, InterpreterResult.Type.TEXT,
"replacement\n");
- runner.appendBuffer("note", "para", 0, "after\n");
+ runner.appendBuffer("note", "para", 0, null, "before-1\n");
+ runner.appendBuffer("note", "para", 0, null, "before-2\n");
+ runner.updateBuffer("note", "para", 0, null, InterpreterResult.Type.TEXT,
"replacement\n");
+ runner.appendBuffer("note", "para", 0, null, "after\n");
runner.run();
InOrder order = inOrder(listener);
- order.verify(listener).onOutputAppend("note", "para", 0,
"before-1\nbefore-2\n");
- order.verify(listener).onOutputUpdated(
- "note", "para", 0, InterpreterResult.Type.TEXT, "replacement\n");
- order.verify(listener).onOutputAppend("note", "para", 0, "after\n");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, null,
"before-1\nbefore-2\n");
+ order.verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.TEXT, "replacement\n");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, null,
"after\n");
}
@Test
@@ -132,11 +186,12 @@ class AppendOutputRunnerTest {
};
loopForCompletingEvents(listener, 4, buffer);
- verify(listener, times(4)).onOutputAppend(any(String.class),
any(String.class), anyInt(), any(String.class));
- verify(listener, times(1)).onOutputAppend(note1, para1, 0, "data1\n");
- verify(listener, times(1)).onOutputAppend(note1, para2, 0, "data2\n");
- verify(listener, times(1)).onOutputAppend(note2, para1, 0, "data3\n");
- verify(listener, times(1)).onOutputAppend(note2, para2, 0, "data4\n");
+ verify(listener, times(4)).onParagraphOutputAppend(
+ any(String.class), any(String.class), anyInt(), isNull(),
any(String.class));
+ verify(listener, times(1)).onParagraphOutputAppend(note1, para1, 0, null,
"data1\n");
+ verify(listener, times(1)).onParagraphOutputAppend(note1, para2, 0, null,
"data2\n");
+ verify(listener, times(1)).onParagraphOutputAppend(note2, para1, 0, null,
"data3\n");
+ verify(listener, times(1)).onParagraphOutputAppend(note2, para2, 0, null,
"data4\n");
}
@Test
@@ -155,7 +210,8 @@ class AppendOutputRunnerTest {
* calls, 30-40 Web-socket calls are made. Keeping
* the unit-test to a pessimistic 100 web-socket calls.
*/
- verify(listener,
atMost(NUM_CLUBBED_EVENTS)).onOutputAppend(any(String.class),
any(String.class), anyInt(), any(String.class));
+ verify(listener, atMost(NUM_CLUBBED_EVENTS)).onParagraphOutputAppend(
+ any(String.class), any(String.class), anyInt(), isNull(),
any(String.class));
}
@Test
@@ -166,7 +222,7 @@ class AppendOutputRunnerTest {
int numEvents = 100000;
for (int i=0; i<numEvents; i++) {
- runner.appendBuffer("noteId", "paraId", 0, data);
+ runner.appendBuffer("noteId", "paraId", 0, null, data);
}
TestAppender appender = new TestAppender();
@@ -196,15 +252,15 @@ class AppendOutputRunnerTest {
RemoteInterpreterProcessListener listener =
mock(RemoteInterpreterProcessListener.class);
AppendOutputRunner runner = new AppendOutputRunner(listener);
doThrow(new IllegalStateException("removed")).when(listener)
- .onOutputUpdated("note", "gone", 0, InterpreterResult.Type.TEXT,
"bad");
- runner.appendBuffer("note", "gone", 0, "bad");
- runner.updateBuffer("note", "gone", 0, InterpreterResult.Type.TEXT, "bad");
- runner.appendBuffer("note", "present", 0, "good");
+ .onParagraphOutputUpdated("note", "gone", 0, null,
InterpreterResult.Type.TEXT, "bad");
+ runner.appendBuffer("note", "gone", 0, null, "bad");
+ runner.updateBuffer("note", "gone", 0, null, InterpreterResult.Type.TEXT,
"bad");
+ runner.appendBuffer("note", "present", 0, null, "good");
runner.run();
- runner.appendBuffer("note", "present", 0, "later");
+ runner.appendBuffer("note", "present", 0, null, "later");
runner.run();
- verify(listener).onOutputAppend("note", "present", 0, "good");
- verify(listener).onOutputAppend("note", "present", 0, "later");
+ verify(listener).onParagraphOutputAppend("note", "present", 0, null,
"good");
+ verify(listener).onParagraphOutputAppend("note", "present", 0, null,
"later");
}
@Test
@@ -212,15 +268,15 @@ class AppendOutputRunnerTest {
RemoteInterpreterProcessListener listener =
mock(RemoteInterpreterProcessListener.class);
AppendOutputRunner runner = new AppendOutputRunner(listener);
doThrow(new IllegalStateException("removed")).when(listener)
- .onOutputAppend("note", "gone", 0, "bad");
- runner.appendBuffer("note", "gone", 0, "bad");
- runner.updateBuffer("note", "present", 0, InterpreterResult.Type.TEXT,
"current");
+ .onParagraphOutputAppend("note", "gone", 0, null, "bad");
+ runner.appendBuffer("note", "gone", 0, null, "bad");
+ runner.updateBuffer("note", "present", 0, null,
InterpreterResult.Type.TEXT, "current");
runner.run();
- runner.appendBuffer("note", "present", 0, "later");
+ runner.appendBuffer("note", "present", 0, null, "later");
runner.run();
- verify(listener).onOutputUpdated("note", "present", 0,
- InterpreterResult.Type.TEXT, "current");
- verify(listener).onOutputAppend("note", "present", 0, "later");
+ verify(listener).onParagraphOutputUpdated(
+ "note", "present", 0, null, InterpreterResult.Type.TEXT, "current");
+ verify(listener).onParagraphOutputAppend("note", "present", 0, null,
"later");
}
@Test
@@ -233,22 +289,22 @@ class AppendOutputRunnerTest {
entered.countDown();
assertTrue(release.await(5, TimeUnit.SECONDS));
return null;
- }).when(listener).onOutputAppend("note", "para", 0, "old");
+ }).when(listener).onParagraphOutputAppend("note", "para", 0, null, "old");
ExecutorService executor = Executors.newFixedThreadPool(2);
try {
- runner.appendBuffer("note", "para", 0, "old");
+ runner.appendBuffer("note", "para", 0, null, "old");
Future<?> first = executor.submit(runner);
assertTrue(entered.await(5, TimeUnit.SECONDS));
- runner.updateBuffer("note", "para", 0, InterpreterResult.Type.TEXT,
"new");
+ runner.updateBuffer("note", "para", 0, null,
InterpreterResult.Type.TEXT, "new");
Future<?> second = executor.submit(runner);
assertThrows(TimeoutException.class, () -> second.get(100,
TimeUnit.MILLISECONDS));
release.countDown();
first.get(5, TimeUnit.SECONDS);
second.get(5, TimeUnit.SECONDS);
InOrder order = inOrder(listener);
- order.verify(listener).onOutputAppend("note", "para", 0, "old");
- order.verify(listener).onOutputUpdated("note", "para", 0,
- InterpreterResult.Type.TEXT, "new");
+ order.verify(listener).onParagraphOutputAppend("note", "para", 0, null,
"old");
+ order.verify(listener).onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.TEXT, "new");
} finally {
release.countDown();
executor.shutdownNow();
@@ -268,7 +324,7 @@ class AppendOutputRunnerTest {
String noteId = "noteId";
String paraId = "paraId";
for (int i=0; i<NUM_EVENTS; i++) {
- runner.appendBuffer(noteId, paraId, 0, "data\n");
+ runner.appendBuffer(noteId, paraId, 0, null, "data\n");
}
}
}
@@ -302,7 +358,8 @@ class AppendOutputRunnerTest {
numInvocations += 1;
return null;
}
- }).when(listener).onOutputAppend(any(String.class), any(String.class),
anyInt(), any(String.class));
+ }).when(listener).onParagraphOutputAppend(
+ any(String.class), any(String.class), anyInt(), isNull(),
any(String.class));
}
private void loopForCompletingEvents(RemoteInterpreterProcessListener
listener,
@@ -311,7 +368,7 @@ class AppendOutputRunnerTest {
prepareInvocationCounts(listener);
AppendOutputRunner runner = new AppendOutputRunner(listener);
for (String[] bufferElement: buffer) {
- runner.appendBuffer(bufferElement[0], bufferElement[1], 0,
bufferElement[2]);
+ runner.appendBuffer(bufferElement[0], bufferElement[1], 0, null,
bufferElement[2]);
}
future = service.scheduleWithFixedDelay(runner, 0,
AppendOutputRunner.BUFFER_TIME_MS, TimeUnit.MILLISECONDS);
diff --git
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerStreamingScopeTest.java
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerStreamingScopeTest.java
index cdc4d5c7bf..697faf808b 100644
---
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerStreamingScopeTest.java
+++
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerStreamingScopeTest.java
@@ -18,6 +18,7 @@ package org.apache.zeppelin.socket;
import static org.junit.jupiter.api.Assertions.assertEquals;
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.eq;
@@ -36,6 +37,7 @@ import
org.apache.zeppelin.interpreter.InterpreterResultMessage;
import org.apache.zeppelin.notebook.Note;
import org.apache.zeppelin.notebook.Notebook;
import org.apache.zeppelin.notebook.Paragraph;
+import org.apache.zeppelin.user.AuthenticationInfo;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -68,8 +70,8 @@ class NotebookServerStreamingScopeTest {
@Test
void sharedNoteReceivesIncrementalOutput() {
- server.onOutputAppend("note", "para", 0, "append");
- server.onOutputUpdated("note", "para", 0, InterpreterResult.Type.TEXT,
"update");
+ server.onParagraphOutputAppend("note", "para", 0, null, "append");
+ server.onParagraphOutputUpdated("note", "para", 0, null,
InterpreterResult.Type.TEXT, "update");
verify(connections, times(2)).broadcast(eq("note"), any());
}
@@ -78,18 +80,111 @@ class NotebookServerStreamingScopeTest {
void personalizedNoteDoesNotReceiveUnownedIncrementalOutput() {
note.setPersonalizedMode(true);
- server.onOutputAppend("note", "para", 0, "private append");
- server.onOutputUpdated("note", "para", 0, InterpreterResult.Type.TEXT,
"private update");
+ server.onParagraphOutputAppend("note", "para", 0, null, "private append");
+ server.onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.TEXT, "private
update");
verify(connections, never()).broadcast(eq("note"), any());
verify(connections, never()).multicastToUser(any(), any());
}
+ @Test
+ void personalizedOutputReachesOnlyItsExecutionOwner() {
+ note.setPersonalizedMode(true);
+ note.getParagraph("para").getUserParagraph("owner");
+
+ server.onParagraphOutputAppend("note", "para", 0, "owner", "mine");
+ server.onParagraphOutputUpdated(
+ "note", "para", 0, "owner", InterpreterResult.Type.TEXT, "mine too");
+
+ verify(connections, never()).broadcast(eq("note"), any());
+ verify(connections, times(2)).multicastToUser(eq("owner"), any());
+ }
+
+ // The shared user field only tracks the note owner's runs, so reading it
misaddressed
+ // everyone else's output.
+ @Test
+ void personalizedOutputIgnoresTheSharedParagraphUserField() {
+ note.setPersonalizedMode(true);
+ Paragraph sharedParagraph = note.getParagraph("para");
+ sharedParagraph.setAuthenticationInfo(new
AuthenticationInfo("stale-owner"));
+ sharedParagraph.getUserParagraph("runner");
+
+ server.onParagraphOutputAppend("note", "para", 0, "runner", "mine");
+
+ verify(connections).multicastToUser(eq("runner"), any());
+ verify(connections, never()).multicastToUser(eq("stale-owner"), any());
+ }
+
+ @Test
+ void personalizedUpdateWritesToTheOwnersCopyNotTheSharedParagraph() {
+ note.setPersonalizedMode(true);
+ Paragraph sharedParagraph = note.getParagraph("para");
+ Paragraph ownerParagraph = sharedParagraph.getUserParagraph("owner");
+
+ server.onParagraphOutputUpdated(
+ "note", "para", 0, "owner", InterpreterResult.Type.TEXT, "owned
update");
+
+ // checkpointOutput always produces a result, so assert on what that
result carries.
+ sharedParagraph.checkpointOutput();
+ assertTrue(sharedParagraph.getReturn().message().isEmpty(),
+ "the shared paragraph absorbed one user's streaming output: "
+ + sharedParagraph.getReturn().message());
+
+ ownerParagraph.checkpointOutput();
+ assertEquals(1, ownerParagraph.getReturn().message().size(),
+ "the owner's copy did not receive its own output");
+ assertEquals("owned update",
ownerParagraph.getReturn().message().get(0).getData());
+ }
+
+ @Test
+ void personalizedClearOnlyDiscardsTheOwnersOutput() {
+ note.setPersonalizedMode(true);
+ Paragraph sharedParagraph = note.getParagraph("para");
+ Paragraph ownerParagraph = sharedParagraph.getUserParagraph("owner");
+ Paragraph otherParagraph = sharedParagraph.getUserParagraph("other");
+ ownerParagraph.setResult(new
InterpreterResult(InterpreterResult.Code.SUCCESS, "owned"));
+ otherParagraph.setResult(new
InterpreterResult(InterpreterResult.Code.SUCCESS, "untouched"));
+
+ server.onParagraphOutputClear("note", "para", "owner");
+
+ assertNull(ownerParagraph.getReturn(), "the owner's output survived its
own clear");
+ assertNotNull(otherParagraph.getReturn(), "another user's output was
cleared");
+ assertEquals("untouched",
otherParagraph.getReturn().message().get(0).getData());
+ verify(connections).multicastToUser(eq("owner"), any());
+ verify(connections, never()).multicastToUser(eq("other"), any());
+ }
+
+ // With paragraph status and progress events off, no streaming message is
sent at all and the
+ // terminal result the user already holds is left alone.
+ @Test
+ void disabledParagraphStatusEventsLeaveTheTerminalResultAlone() {
+ ZeppelinConfiguration disabled = mock(ZeppelinConfiguration.class);
+ when(disabled.getBoolean(
+
ZeppelinConfiguration.ConfVars.ZEPPELIN_WEBSOCKET_PARAGRAPH_STATUS_PROGRESS))
+ .thenReturn(false);
+ server.setZeppelinConfiguration(disabled);
+ note.setPersonalizedMode(true);
+ Paragraph ownerParagraph =
note.getParagraph("para").getUserParagraph("owner");
+ ownerParagraph.setResult(new
InterpreterResult(InterpreterResult.Code.SUCCESS, "terminal"));
+
+ server.onParagraphOutputAppend("note", "para", 0, "owner", "append");
+ server.onParagraphOutputUpdated(
+ "note", "para", 0, "owner", InterpreterResult.Type.TEXT, "update");
+ server.onParagraphOutputClear("note", "para", "owner");
+
+ verify(connections, never()).broadcast(eq("note"), any());
+ verify(connections, never()).multicastToUser(any(), any());
+ assertNotNull(ownerParagraph.getReturn(), "the terminal result was
discarded");
+ assertEquals("terminal",
ownerParagraph.getReturn().message().get(0).getData());
+ }
+
@Test
void personalizedCheckpointDoesNotExposeUnownedOutputToOtherUsers() {
note.setPersonalizedMode(true);
- server.onOutputUpdated("note", "para", 0, InterpreterResult.Type.TEXT,
"private update");
+ server.onParagraphOutputUpdated(
+ "note", "para", 0, null, InterpreterResult.Type.TEXT, "private
update");
server.checkpointOutput("note", "para");
InterpreterResult otherUserResult =
@@ -109,7 +204,7 @@ class NotebookServerStreamingScopeTest {
sharedParagraph.updateOutputBuffer(0, InterpreterResult.Type.TEXT,
"buffered result");
sharedParagraph.getUserParagraph("existing");
- server.onOutputClear("note", "para");
+ server.onParagraphOutputClear("note", "para", null);
InterpreterResult futureUserResult =
sharedParagraph.getUserParagraph("future").getReturn();
diff --git
a/zeppelin-web-angular/e2e/tests/notebook/personalized/personalized-streaming-scope.spec.ts
b/zeppelin-web-angular/e2e/tests/notebook/personalized/personalized-streaming-scope.spec.ts
new file mode 100644
index 0000000000..268e51cd6f
--- /dev/null
+++
b/zeppelin-web-angular/e2e/tests/notebook/personalized/personalized-streaming-scope.spec.ts
@@ -0,0 +1,158 @@
+/*
+ * Licensed 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.
+ */
+
+import { BrowserContext, expect, Page, test } from '@playwright/test';
+import { CollaborationPage } from 'e2e/models/collaboration-page';
+import { NotebookParagraphPage } from 'e2e/models/notebook-paragraph-page';
+import { addPageAnnotationBeforeEach, createTestNotebook, getTwoTestAccounts,
loginAs, PAGES } from '../../../utils';
+
+// browser.newContext() inherits the project's storageState, already signed in
as the first
+// account. These tests choose their own principal, so they start from nothing.
+const SIGNED_OUT = { cookies: [], origins: [] };
+
+/**
+ * Records the incremental-output frames the server sends to one page. The
defect is in which
+ * connection the server addresses, so the frames are the subject rather than
what the UI draws.
+ * Attach before the page navigates: the socket opens as the app boots.
+ */
+const recordStreamingFrames = (page: Page): string[] => {
+ const frames: string[] = [];
+ page.on('websocket', socket => {
+ socket.on('framereceived', frame => {
+ const payload = typeof frame.payload === 'string' ? frame.payload :
frame.payload.toString();
+ if (payload.includes('PARAGRAPH_APPEND_OUTPUT') ||
payload.includes('PARAGRAPH_UPDATE_OUTPUT')) {
+ frames.push(payload);
+ }
+ });
+ });
+ return frames;
+};
+
+/** The pid the shell printed, which differs for every execution of the same
code. */
+const markerIn = (frames: string[]): string => {
+ const marker = /MARK-\d+/.exec(frames.join('\n'));
+ expect(marker, `no marker in ${frames.length} streaming
frames`).not.toBeNull();
+ return marker![0];
+};
+
+test.describe('Personalized streaming scope', () => {
+ // JUSTIFIED: one shared notebook and two signed-in principals must stay
within one worker.
+ test.describe.configure({ mode: 'default' });
+ addPageAnnotationBeforeEach(PAGES.WORKSPACE.NOTEBOOK);
+
+ let ownerContext: BrowserContext;
+ let otherContext: BrowserContext;
+
+ test.afterEach(async () => {
+ await ownerContext?.close();
+ await otherContext?.close();
+ });
+
+ test('sends each user only the streaming output of their own run', async ({
browser }) => {
+ const accounts = await getTwoTestAccounts();
+ test.skip(!accounts, 'Two shiro accounts are required to tell personalized
users apart');
+ const [ownerAccount, otherAccount] = accounts!;
+
+ ownerContext = await browser.newContext({ storageState: SIGNED_OUT });
+ otherContext = await browser.newContext({ storageState: SIGNED_OUT });
+ const ownerPage = await ownerContext.newPage();
+ const otherPage = await otherContext.newPage();
+ const ownerFrames = recordStreamingFrames(ownerPage);
+ const otherFrames = recordStreamingFrames(otherPage);
+
+ const ownerNote = new CollaborationPage(ownerPage);
+ const otherNote = new CollaborationPage(otherPage);
+ const ownerParagraph = new NotebookParagraphPage(ownerPage);
+ const otherParagraph = new NotebookParagraphPage(otherPage);
+
+ await test.step('Given a personalized note two signed-in users share',
async () => {
+ await loginAs(ownerPage, ownerAccount);
+ await loginAs(otherPage, otherAccount);
+
+ const { noteId } = await createTestNotebook(ownerPage);
+ await ownerNote.openNotebook(noteId);
+ await ownerNote.switchToPersonalModeButton.click();
+ await ownerNote.confirmPersonalizedModeChange();
+ await expect(ownerNote.switchToCollaborationModeButton).toBeVisible({
timeout: 15000 });
+
+ await otherNote.openNotebook(noteId);
+ });
+
+ await test.step('When both run that paragraph', async () => {
+ // Personalized mode shares and live-syncs the paragraph text, so the
two users cannot hold
+ // different code. The shell pid differs per run, giving each execution
its own marker.
+ await ownerNote.typeInEditor('%sh\necho MARK-$$');
+ await expect(otherNote.editorText).toContainText('MARK-', { timeout:
15000 });
+
+ await ownerParagraph.runParagraph();
+ await otherParagraph.runParagraph();
+
+ await expect(ownerParagraph.status).toHaveText('FINISHED');
+ await expect(otherParagraph.status).toHaveText('FINISHED');
+ });
+
+ await test.step('Then neither socket carried the other run output', async
() => {
+ const ownerMarker = markerIn(ownerFrames);
+ const otherMarker = markerIn(otherFrames);
+ expect(ownerMarker).not.toEqual(otherMarker);
+
+ expect(ownerFrames.join('\n')).not.toContain(otherMarker);
+ expect(otherFrames.join('\n')).not.toContain(ownerMarker);
+ });
+
+ await test.step('And the terminal result each user sees matches their own
run', async () => {
+ const ownerMarker = markerIn(ownerFrames);
+ const otherMarker = markerIn(otherFrames);
+
+ await expect(ownerParagraph.resultDisplay).toContainText(ownerMarker);
+ await
expect(ownerParagraph.resultDisplay).not.toContainText(otherMarker);
+ await expect(otherParagraph.resultDisplay).toContainText(otherMarker);
+ await
expect(otherParagraph.resultDisplay).not.toContainText(ownerMarker);
+ });
+ });
+
+ test('sends no streaming output to a user who only opened the note', async
({ browser }) => {
+ const accounts = await getTwoTestAccounts();
+ test.skip(!accounts, 'Two shiro accounts are required to tell personalized
users apart');
+ const [ownerAccount, otherAccount] = accounts!;
+
+ ownerContext = await browser.newContext({ storageState: SIGNED_OUT });
+ otherContext = await browser.newContext({ storageState: SIGNED_OUT });
+ const ownerPage = await ownerContext.newPage();
+ const watcherPage = await otherContext.newPage();
+ const ownerFrames = recordStreamingFrames(ownerPage);
+ const watcherFrames = recordStreamingFrames(watcherPage);
+
+ await loginAs(ownerPage, ownerAccount);
+ await loginAs(watcherPage, otherAccount);
+
+ const { noteId } = await createTestNotebook(ownerPage);
+ const ownerNote = new CollaborationPage(ownerPage);
+ await ownerNote.openNotebook(noteId);
+ await ownerNote.switchToPersonalModeButton.click();
+ await ownerNote.confirmPersonalizedModeChange();
+ await expect(ownerNote.switchToCollaborationModeButton).toBeVisible({
timeout: 15000 });
+ await new CollaborationPage(watcherPage).openNotebook(noteId);
+
+ await ownerNote.typeInEditor('%sh\necho RUNNER-ONLY');
+ const ownerParagraph = new NotebookParagraphPage(ownerPage);
+ await ownerParagraph.runParagraph();
+ await expect(ownerParagraph.status).toHaveText('FINISHED');
+ await expect(ownerParagraph.resultDisplay).toContainText('RUNNER-ONLY');
+
+ // Both halves matter. On their own "the watcher got nothing" also holds
when the server
+ // sends the output to nobody, so it would pass against a build that
simply drops it.
+ expect(ownerFrames.join('\n')).toContain('RUNNER-ONLY');
+ expect(watcherFrames).toHaveLength(0);
+ await expect(new
NotebookParagraphPage(watcherPage).resultDisplay).toHaveCount(0);
+ });
+});
diff --git a/zeppelin-web-angular/e2e/utils.ts
b/zeppelin-web-angular/e2e/utils.ts
index dd1c626243..06cacc8ee0 100644
--- a/zeppelin-web-angular/e2e/utils.ts
+++ b/zeppelin-web-angular/e2e/utils.ts
@@ -13,7 +13,7 @@
import { globSync } from 'fs';
import { join, sep } from 'path';
import { test, expect, Page, TestInfo } from '@playwright/test';
-import { LoginTestUtil } from './models/login-page.util';
+import { LoginTestUtil, TestCredentials } from './models/login-page.util';
import { E2E_TEST_FOLDER } from './models/base-page';
import { LoginPage } from './models/login-page';
@@ -298,6 +298,41 @@ export const performLoginIfRequired = async (page: Page):
Promise<boolean> => {
return false;
};
+/**
+ * Signs in as one specific account. performLoginIfRequired always takes the
first configured
+ * user, which cannot express a test that needs two distinct principals at
once.
+ */
+export const loginAs = async (page: Page, credentials: TestCredentials):
Promise<void> => {
+ const loginPage = new LoginPage(page);
+ await loginPage.navigate();
+ // Wait for the form rather than probing visibility, which resolves before
Angular renders it.
+ await loginPage.userNameInput.waitFor({ state: 'visible', timeout: 30000 });
+ await loginPage.login(credentials.username, credentials.password);
+ await page.waitForSelector('zeppelin-login', { state: 'hidden', timeout:
30000 });
+ await page.evaluate(() => {
+ if (window.location.hash.includes('login')) {
+ window.location.hash = '#/';
+ }
+ });
+ // The note list only renders once the authenticated socket has delivered
it, so it is the
+ // signal that this principal can open a notebook -- not just that the form
was accepted.
+ await page.waitForSelector('zeppelin-node-list', { timeout: 30000 });
+ await waitForZeppelinReady(page);
+};
+
+/**
+ * Two distinct shiro accounts, or null when the deployment cannot provide
them.
+ */
+export const getTwoTestAccounts = async (): Promise<[TestCredentials,
TestCredentials] | null> => {
+ if (!(await LoginTestUtil.isShiroEnabled())) {
+ return null;
+ }
+ const accounts = Object.values(await
LoginTestUtil.getTestCredentials()).filter(
+ account => account.username && account.password
+ );
+ return accounts.length >= 2 ? [accounts[0], accounts[1]] : null;
+};
+
export const skipWhenAuthenticationIsStillRequired = async (page: Page):
Promise<void> => {
const loginStillVisible = await page
.locator('zeppelin-login')