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')

Reply via email to