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 3219ccc8ed [ZEPPELIN-6683] Prevent stale notebook WebSocket replies 
from mutating the active note
3219ccc8ed is described below

commit 3219ccc8eddd84396a1c2ae249b9c19439cc17a5
Author: gyowoo1113 <[email protected]>
AuthorDate: Thu Sep 24 23:04:54 2026 +0900

    [ZEPPELIN-6683] Prevent stale notebook WebSocket replies from mutating the 
active note
    
    ### What is this PR for?
    This PR prevents stale notebook WebSocket replies from mutating the 
currently active note.
    
    `GET_INTERPRETER_BINDINGS`, `SAVE_INTERPRETER_BINDINGS`, 
`LIST_REVISION_HISTORY`, checkpoint creation, and `SET_NOTE_REVISION` 
previously returned replies without preserving the originating request 
identity. The SDK `receive()` path also exposed only `message.data`, so 
notebook listeners could not determine which note a late reply belonged to.
    
    This PR preserves the originating `msgId` on the relevant server reply 
paths, exposes the full WebSocket envelope through the SDK, and records the 
note context for pending requests on the client. Notebook handlers now reject 
replies with missing or unknown request identities, as well as replies whose 
request note no longer matches the active route.
    
    The existing operation names are unchanged, and the change is limited to 
interpreter-binding and revision-related reply paths.
    
    
    ### What type of PR is it?
    Bug Fix
    
    ### Todos
    - [x] Preserve `msgId` on `INTERPRETER_BINDINGS` responses
    - [x] Preserve `msgId` on `LIST_REVISION_HISTORY` responses, including 
checkpoint replies
    - [x] Preserve `msgId` on `SET_NOTE_REVISION` responses
    - [x] Expose WebSocket reply envelopes through the SDK
    - [x] Track pending request note context and reject stale, missing, or 
unknown replies
    - [x] Add server tests for all affected response paths
    - [x] Add frontend tests for stale reply handling and valid same-note 
handling
    
    ### What is the Jira issue?
    [[ZEPPELIN-6683]](https://issues.apache.org/jira/browse/ZEPPELIN-6683)
    
    ### How should this be tested?
    Frontend focused tests (from `zeppelin-web-angular`):
    
    ```bash
    npm run test:shell -- \
      projects/zeppelin-sdk/src/message.spec.ts \
      src/app/core/message-listener/message-listener.spec.ts \
      src/app/pages/workspace/notebook/notebook.component.spec.ts
    ```
    WebSocket contract check:
    
    ```bash
    npm run check:websocket-contract
    ```
    
    Server focused test:
    ```bash
    ./mvnw -pl zeppelin-server -Dtest=NotebookServerTest test
    ```
    All checks above pass successfully.
    
    ### Screenshots (if appropriate)
    N/A
    
    ### Questions:
    * Does the license files need to update? No
    * Is there breaking changes for older versions? No
    * Does this needs documentation? No
    
    
    Closes #5487 from gyowoo1113/ZEPPELIN-6683-stale-notebook-replies.
    
    Signed-off-by: ChanHo Lee <[email protected]>
---
 .../org/apache/zeppelin/socket/NotebookServer.java |  20 ++-
 .../apache/zeppelin/socket/NotebookServerTest.java |  92 +++++++++-
 .../interfaces/message-interpreter.interface.ts    |   1 +
 .../src/interfaces/message-notebook.interface.ts   |   3 +-
 .../projects/zeppelin-sdk/src/message.spec.ts      |  18 ++
 .../projects/zeppelin-sdk/src/message.ts           |  39 ++--
 .../core/message-listener/message-listener.spec.ts |  38 +++-
 .../app/core/message-listener/message-listener.ts  |  27 ++-
 .../workspace/notebook/notebook.component.spec.ts  | 197 +++++++++++++++++++++
 .../pages/workspace/notebook/notebook.component.ts |  32 +++-
 .../src/app/services/message.service.ts            |   4 -
 11 files changed, 425 insertions(+), 46 deletions(-)

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 8e568e1217..41a65e3a4f 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
@@ -696,7 +696,9 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
                 setting.getInterpreterInfos(), true));
           }
         }
-        conn.send(serializeMessage(new 
Message(OP.INTERPRETER_BINDINGS).put("interpreterBindings", settingList)));
+        conn.send(serializeMessage(new Message(OP.INTERPRETER_BINDINGS)
+            .put("noteId", noteId)
+            .put("interpreterBindings", settingList)));
         return null;
       });
   }
@@ -734,7 +736,9 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
       });
     if (permitted) {
       conn.send(serializeMessage(
-          new Message(OP.INTERPRETER_BINDINGS).put("interpreterBindings", 
settingList)));
+          new Message(OP.INTERPRETER_BINDINGS)
+              .put("noteId", noteId)
+              .put("interpreterBindings", settingList)));
     }
   }
 
@@ -1669,7 +1673,9 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
 
               List<Revision> revisions = getNotebook().processNote(noteId,
                 note -> getNotebook().listRevisionHistory(noteId, 
note.getPath(), context.getAutheInfo()));
-              conn.send(serializeMessage(new 
Message(OP.LIST_REVISION_HISTORY).put("revisionList", revisions)));
+              conn.send(serializeMessage(new Message(OP.LIST_REVISION_HISTORY)
+                  .put("noteId", noteId)
+                  .put("revisionList", revisions)));
             } else {
               conn.send(serializeMessage(
                   new Message(OP.ERROR_INFO).put("info",
@@ -1689,7 +1695,9 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
           @Override
           public void onSuccess(List<Revision> revisions, ServiceContext 
context) throws IOException {
             super.onSuccess(revisions, context);
-            conn.send(serializeMessage(new 
Message(OP.LIST_REVISION_HISTORY).put("revisionList", revisions)));
+            conn.send(serializeMessage(new Message(OP.LIST_REVISION_HISTORY)
+                .put("noteId", noteId)
+                .put("revisionList", revisions)));
           }
         });
   }
@@ -1705,7 +1713,9 @@ public class NotebookServer implements 
AngularObjectRegistryListener,
           public void onSuccess(Note note, ServiceContext context) throws 
IOException {
             super.onSuccess(note, context);
             Note reloadedNote = getNotebook().loadNoteFromRepo(noteId, 
context.getAutheInfo());
-            conn.send(serializeMessage(new 
Message(OP.SET_NOTE_REVISION).put("status", true)));
+            conn.send(serializeMessage(new Message(OP.SET_NOTE_REVISION)
+                .put("noteId", noteId)
+                .put("status", true)));
             broadcastNote(reloadedNote);
           }
         });
diff --git 
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java
 
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java
index d288851fbf..5de5b1d10f 100644
--- 
a/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java
+++ 
b/zeppelin-server/src/test/java/org/apache/zeppelin/socket/NotebookServerTest.java
@@ -1067,15 +1067,17 @@ class NotebookServerTest extends AbstractTestRestApi {
 
       ArgumentCaptor<String> response = ArgumentCaptor.forClass(String.class);
       verify(socket).send(response.capture());
-      assertEquals(OP.AUTH_INFO, 
notebookServer.deserializeMessage(response.getValue()).op);
+      Message responseMessage = 
notebookServer.deserializeMessage(response.getValue());
+      assertEquals(OP.AUTH_INFO, responseMessage.op);
 
       reset(socket);
       setNotePermissions(noteId, "binding-owner", "binding-reader");
       notebookServer.getInterpreterBindings(socket, 
serviceContext("binding-reader"), message);
 
       verify(socket).send(response.capture());
-      assertEquals(OP.INTERPRETER_BINDINGS,
-          notebookServer.deserializeMessage(response.getValue()).op);
+      responseMessage = notebookServer.deserializeMessage(response.getValue());
+      assertEquals(OP.INTERPRETER_BINDINGS, responseMessage.op);
+      assertEquals(noteId, responseMessage.data.get("noteId"));
     } finally {
       notebook.removeNote(noteId, owner);
     }
@@ -1100,7 +1102,8 @@ class NotebookServerTest extends AbstractTestRestApi {
           notebook.processNote(noteId, Note::getDefaultInterpreterGroup));
       ArgumentCaptor<String> response = ArgumentCaptor.forClass(String.class);
       verify(socket).send(response.capture());
-      assertEquals(OP.AUTH_INFO, 
notebookServer.deserializeMessage(response.getValue()).op);
+      Message responseMessage = 
notebookServer.deserializeMessage(response.getValue());
+      assertEquals(OP.AUTH_INFO, responseMessage.op);
 
       reset(socket);
       authorizationService.setWriters(noteId,
@@ -1110,8 +1113,9 @@ class NotebookServerTest extends AbstractTestRestApi {
       assertEquals(replacementGroup,
           notebook.processNote(noteId, Note::getDefaultInterpreterGroup));
       verify(socket).send(response.capture());
-      assertEquals(OP.INTERPRETER_BINDINGS,
-          notebookServer.deserializeMessage(response.getValue()).op);
+      responseMessage = notebookServer.deserializeMessage(response.getValue());
+      assertEquals(OP.INTERPRETER_BINDINGS, responseMessage.op);
+      assertEquals(noteId, responseMessage.data.get("noteId"));
     } finally {
       notebook.removeNote(noteId, owner);
     }
@@ -1170,6 +1174,82 @@ class NotebookServerTest extends AbstractTestRestApi {
     }
   }
 
+  @Test
+  void listRevisionHistoryIncludesNoteId() throws IOException {
+    String noteId = notebook.createNote("revision-list-note-id", anonymous);
+
+    try {
+      NotebookSocket socket = createWebSocket();
+      Message request = new Message(OP.LIST_REVISION_HISTORY)
+          .put("noteId", noteId);
+
+      notebookServer.onMessage(socket, request.toJson());
+
+      ArgumentCaptor<String> response = ArgumentCaptor.forClass(String.class);
+      verify(socket).send(response.capture());
+
+      Message responseMessage = 
notebookServer.deserializeMessage(response.getValue());
+      assertEquals(OP.LIST_REVISION_HISTORY, responseMessage.op);
+      assertEquals(noteId, responseMessage.data.get("noteId"));
+    } finally {
+      notebook.removeNote(noteId, anonymous);
+    }
+  }
+
+  @Test
+  void checkpointNoteIncludesNoteId() throws IOException {
+    String noteId = notebook.createNote("checkpoint-note-id", anonymous);
+
+    try {
+      NotebookSocket socket = createWebSocket();
+      Message request = new Message(OP.CHECKPOINT_NOTE)
+          .put("noteId", noteId)
+          .put("commitMessage", "checkpoint");
+
+      notebookServer.onMessage(socket, request.toJson());
+
+      ArgumentCaptor<String> response = ArgumentCaptor.forClass(String.class);
+      verify(socket).send(response.capture());
+
+      Message responseMessage = 
notebookServer.deserializeMessage(response.getValue());
+      assertEquals(OP.LIST_REVISION_HISTORY, responseMessage.op);
+      assertEquals(noteId, responseMessage.data.get("noteId"));
+    } finally {
+      notebook.removeNote(noteId, anonymous);
+    }
+  }
+
+  @Test
+  void setNoteRevisionIncludesNoteId() throws IOException {
+    String noteId = notebook.createNote("set-revision-note-id", anonymous);
+
+    try {
+      NotebookRepoWithVersionControl.Revision revision =
+          notebook.processNote(noteId,
+              note -> notebook.checkpointNote(
+                  note.getId(),
+                  note.getPath(),
+                  "revision",
+                  anonymous));
+
+      NotebookSocket socket = createWebSocket();
+      Message request = new Message(OP.SET_NOTE_REVISION)
+          .put("noteId", noteId)
+          .put("revisionId", revision.id);
+
+      notebookServer.onMessage(socket, request.toJson());
+
+      ArgumentCaptor<String> response = ArgumentCaptor.forClass(String.class);
+      verify(socket).send(response.capture());
+
+      Message responseMessage = 
notebookServer.deserializeMessage(response.getValue());
+      assertEquals(OP.SET_NOTE_REVISION, responseMessage.op);
+      assertEquals(noteId, responseMessage.data.get("noteId"));
+    } finally {
+      notebook.removeNote(noteId, anonymous);
+    }
+  }
+
   private NotebookSocket createWebSocket() {
     NotebookSocket sock = mock(NotebookSocket.class);
     return sock;
diff --git 
a/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-interpreter.interface.ts
 
b/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-interpreter.interface.ts
index c59e459410..f9a4ac4957 100644
--- 
a/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-interpreter.interface.ts
+++ 
b/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-interpreter.interface.ts
@@ -26,6 +26,7 @@ export interface InterpreterItem {
 }
 
 export interface InterpreterBindings {
+  noteId: string;
   interpreterBindings: InterpreterBindingItem[];
 }
 
diff --git 
a/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-notebook.interface.ts
 
b/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-notebook.interface.ts
index 665e8dfd71..fca5c07198 100644
--- 
a/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-notebook.interface.ts
+++ 
b/zeppelin-web-angular/projects/zeppelin-sdk/src/interfaces/message-notebook.interface.ts
@@ -167,15 +167,16 @@ export interface ImportNoteReceived {
 
 export interface ParagraphAdded {
   index: number;
-  msgId?: string;
   paragraph: ParagraphItem;
 }
 
 export interface SetNoteRevisionStatus {
+  noteId: string;
   status: boolean;
 }
 
 export interface ListRevision {
+  noteId: string;
   revisionList: RevisionListItem[];
 }
 
diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts 
b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts
index 0bb02529cc..37f8d44c2b 100644
--- a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts
+++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts
@@ -143,4 +143,22 @@ describe('Message.receive', () => {
 
     expect(listener).toHaveBeenCalledWith(data);
   });
+
+  it('passes the full message envelope with msgId', () => {
+    const message = new Message();
+    const listener = vi.fn();
+    const data = {};
+
+    message.receiveEnvelope(OP.NOTE).subscribe(listener);
+
+    const envelope = asReceivedMessage({
+      op: OP.NOTE,
+      msgId: 'note-request-1',
+      data
+    });
+
+    message.shortCircuit(envelope);
+
+    expect(listener).toHaveBeenCalledWith(envelope);
+  });
 });
diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts 
b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts
index 4d559a86aa..f9dff2e804 100644
--- a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts
+++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts
@@ -175,21 +175,15 @@ export class Message {
   }
 
   receive<K extends keyof MessageReceiveDataTypeMap>(op: K): 
Observable<Record<K, MessageReceiveDataTypeMap[K]>[K]> {
-    const guard = getMessagePayloadGuard(op);
-
-    return this.received$.pipe(
-      filter(message => message.op === op),
-      filter(message => {
-        if (!guard || guard(message.data)) {
-          return true;
-        }
+    return this.receiveMessage(op).pipe(map(message => message.data)) as 
Observable<
+      Record<K, MessageReceiveDataTypeMap[K]>[K]
+    >;
+  }
 
-        // The payload can be large and carries note names, so log the OP 
alone.
-        console.warn(`Dropped WebSocket OP ${String(op)}: payload failed 
validation`);
-        return false;
-      }),
-      map(message => message.data)
-    ) as Observable<Record<K, MessageReceiveDataTypeMap[K]>[K]>;
+  receiveEnvelope<K extends keyof MessageReceiveDataTypeMap>(
+    op: K
+  ): Observable<WebSocketMessage<MessageReceiveDataTypeMap, K>> {
+    return this.receiveMessage(op) as 
Observable<WebSocketMessage<MessageReceiveDataTypeMap, K>>;
   }
 
   shortCircuit(message: WebSocketMessage<MessageReceiveDataTypeMap>) {
@@ -565,4 +559,21 @@ export class Message {
       formName
     });
   }
+
+  private receiveMessage<K extends keyof MessageReceiveDataTypeMap>(op: K) {
+    const guard = getMessagePayloadGuard(op);
+
+    return this.received$.pipe(
+      filter(message => message.op === op),
+      filter(message => {
+        if (!guard || guard(message.data)) {
+          return true;
+        }
+
+        // The payload can be large and carries note names, so log the OP 
alone.
+        console.warn(`Dropped WebSocket OP ${String(op)}: payload failed 
validation`);
+        return false;
+      })
+    );
+  }
 }
diff --git 
a/zeppelin-web-angular/src/app/core/message-listener/message-listener.spec.ts 
b/zeppelin-web-angular/src/app/core/message-listener/message-listener.spec.ts
index e12767691d..7fd1778ae6 100644
--- 
a/zeppelin-web-angular/src/app/core/message-listener/message-listener.spec.ts
+++ 
b/zeppelin-web-angular/src/app/core/message-listener/message-listener.spec.ts
@@ -13,9 +13,8 @@
 import { Subject } from 'rxjs';
 import { afterEach, describe, expect, it, vi } from 'vitest';
 
-import { Message, OP, MessageReceiveDataTypeMap } from '@zeppelin/sdk';
-
-import { MessageListener, MessageListenersManager } from './message-listener';
+import { Message, OP, MessageReceiveDataTypeMap, type WebSocketMessage } from 
'@zeppelin/sdk';
+import { MessageEnvelopeListener, MessageListener, MessageListenersManager } 
from './message-listener';
 
 afterEach(() => {
   vi.restoreAllMocks();
@@ -88,3 +87,36 @@ describe('MessageListener', () => {
     expect(component.receivedData).toBe(data);
   });
 });
+
+describe('MessageEnvelopeListener', () => {
+  it('passes the received message envelope with msgId to the handler', () => {
+    const received$ = new Subject<WebSocketMessage<MessageReceiveDataTypeMap, 
OP.NOTE>>();
+    const messageService = {
+      receiveEnvelope: vi.fn(() => received$.asObservable())
+    } as unknown as Message;
+
+    class TestComponent extends MessageListenersManager {
+      receivedMessage?: WebSocketMessage<MessageReceiveDataTypeMap, OP.NOTE>;
+
+      handleNote(message: WebSocketMessage<MessageReceiveDataTypeMap, 
OP.NOTE>): void {
+        this.receivedMessage = message;
+      }
+    }
+
+    const descriptor = 
Object.getOwnPropertyDescriptor(TestComponent.prototype, 'handleNote')!;
+
+    MessageEnvelopeListener(OP.NOTE)(TestComponent.prototype, 'handleNote', 
descriptor);
+
+    const component = new TestComponent(messageService);
+    const envelope: WebSocketMessage<MessageReceiveDataTypeMap, OP.NOTE> = {
+      op: OP.NOTE,
+      msgId: 'note-request-1',
+      data: {} as MessageReceiveDataTypeMap[OP.NOTE]
+    };
+
+    received$.next(envelope);
+
+    expect(component.receivedMessage).toBe(envelope);
+    expect(messageService.receiveEnvelope).toHaveBeenCalledWith(OP.NOTE);
+  });
+});
diff --git 
a/zeppelin-web-angular/src/app/core/message-listener/message-listener.ts 
b/zeppelin-web-angular/src/app/core/message-listener/message-listener.ts
index 1b2f0209ae..93c6294700 100644
--- a/zeppelin-web-angular/src/app/core/message-listener/message-listener.ts
+++ b/zeppelin-web-angular/src/app/core/message-listener/message-listener.ts
@@ -11,9 +11,9 @@
  */
 
 import { Component, OnDestroy } from '@angular/core';
-import { Subscriber } from 'rxjs';
+import { Observable, Subscriber } from 'rxjs';
 
-import { Message, MessageReceiveDataTypeMap, ReceiveArgumentsType } from 
'@zeppelin/sdk';
+import { Message, MessageReceiveDataTypeMap } from '@zeppelin/sdk';
 
 @Component({
   template: '',
@@ -34,13 +34,18 @@ export class MessageListenersManager implements OnDestroy {
   }
 }
 
-export function MessageListener<K extends keyof MessageReceiveDataTypeMap>(op: 
K) {
+type ListenerArgumentsType<T> = T extends undefined ? () => void : (data: T) 
=> void;
+
+const createMessageListener = <K extends keyof MessageReceiveDataTypeMap, T>(
+  op: K,
+  receiver: (messageService: Message, op: K) => Observable<T>
+) => {
   return function (
     target: MessageListenersManager,
     propertyKey: string,
-    descriptor: TypedPropertyDescriptor<ReceiveArgumentsType<K>>
+    descriptor: TypedPropertyDescriptor<ListenerArgumentsType<T>>
   ) {
-    const oldValue = descriptor.value as ReceiveArgumentsType<K>;
+    const oldValue = descriptor.value as ListenerArgumentsType<T>;
 
     const fn = function (this: MessageListenersManager) {
       if (!this.__zeppelinMessageListeners$__) {
@@ -48,7 +53,7 @@ export function MessageListener<K extends keyof 
MessageReceiveDataTypeMap>(op: K
       }
 
       this.__zeppelinMessageListeners$__.add(
-        this.messageService.receive(op).subscribe(data => {
+        receiver(this.messageService, op).subscribe(data => {
           try {
             // @ts-ignore
             oldValue.apply(this, [data]);
@@ -68,4 +73,12 @@ export function MessageListener<K extends keyof 
MessageReceiveDataTypeMap>(op: K
 
     return descriptor;
   };
-}
+};
+
+export const MessageListener = <K extends keyof MessageReceiveDataTypeMap>(op: 
K) => {
+  return createMessageListener(op, (messageService, targetOp) => 
messageService.receive(targetOp));
+};
+
+export const MessageEnvelopeListener = <K extends keyof 
MessageReceiveDataTypeMap>(op: K) => {
+  return createMessageListener(op, (messageService, targetOp) => 
messageService.receiveEnvelope(targetOp));
+};
diff --git 
a/zeppelin-web-angular/src/app/pages/workspace/notebook/notebook.component.spec.ts
 
b/zeppelin-web-angular/src/app/pages/workspace/notebook/notebook.component.spec.ts
new file mode 100644
index 0000000000..62186de41f
--- /dev/null
+++ 
b/zeppelin-web-angular/src/app/pages/workspace/notebook/notebook.component.spec.ts
@@ -0,0 +1,197 @@
+/*
+ * 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 { NEVER } from 'rxjs';
+import { afterEach, describe, expect, it, vi } from 'vitest';
+
+import type { ChangeDetectorRef } from '@angular/core';
+import type { Title } from '@angular/platform-browser';
+import type { ActivatedRoute, Router } from '@angular/router';
+import { OP, type MessageReceiveDataTypeMap, type WebSocketMessage } from 
'@zeppelin/sdk';
+
+import type {
+  MessageService,
+  NgZService,
+  NoteStatusService,
+  NoteVarShareService,
+  ReactFeatureService,
+  SecurityService,
+  ThemeService,
+  TicketService
+} from '@zeppelin/services';
+
+vi.mock('./paragraph/paragraph.component', () => ({
+  NotebookParagraphComponent: class NotebookParagraphComponent {}
+}));
+
+import { NotebookComponent } from './notebook.component';
+
+const createComponent = () => {
+  const messageService = {
+    receive: vi.fn(() => NEVER),
+    receiveEnvelope: vi.fn(() => NEVER),
+    consumeLocalAddFocusMsgId: vi.fn(() => false)
+  };
+
+  const router = {
+    navigate: vi.fn(() => Promise.resolve(true))
+  } as unknown as Router;
+
+  const component = new NotebookComponent(
+    messageService as unknown as MessageService,
+    {} as NgZService,
+    {
+      snapshot: {
+        params: {
+          noteId: 'note-b'
+        }
+      }
+    } as unknown as ActivatedRoute,
+    {
+      markForCheck: vi.fn()
+    } as unknown as ChangeDetectorRef,
+    {} as NoteStatusService,
+    {} as NoteVarShareService,
+    {} as TicketService,
+    {} as SecurityService,
+    router,
+    {} as Title,
+    {} as ThemeService,
+    {} as ReactFeatureService
+  );
+
+  return {
+    component,
+    router,
+    messageService
+  };
+};
+
+afterEach(() => {
+  vi.restoreAllMocks();
+});
+
+describe('NotebookComponent', () => {
+  it('ignores stale interpreter bindings replies', () => {
+    const { component } = createComponent();
+
+    const originalBindings = [{ id: 'existing', selected: true }];
+    component.interpreterBindings = originalBindings as typeof 
component.interpreterBindings;
+
+    const data: MessageReceiveDataTypeMap[OP.INTERPRETER_BINDINGS] = {
+      noteId: 'note-a',
+      interpreterBindings: [{ id: 'stale', selected: false }]
+    };
+
+    component.loadInterpreterBindings(data);
+
+    expect(component.interpreterBindings).toBe(originalBindings);
+  });
+
+  it('ignores stale revision history replies', () => {
+    const { component } = createComponent();
+
+    const originalRevisions = [{ id: 'Head', message: 'Head' }];
+    component.noteRevisions = originalRevisions;
+    component.currentRevision = 'Head';
+
+    const data: MessageReceiveDataTypeMap[OP.LIST_REVISION_HISTORY] = {
+      noteId: 'note-a',
+      revisionList: [{ id: 'rev-1', message: 'stale revision' }]
+    };
+
+    component.listRevisionHistory(data);
+
+    expect(component.noteRevisions).toBe(originalRevisions);
+    expect(component.currentRevision).toBe('Head');
+  });
+
+  it('ignores stale set note revision replies', () => {
+    const { component, router } = createComponent();
+
+    const data: MessageReceiveDataTypeMap[OP.SET_NOTE_REVISION] = {
+      noteId: 'note-a',
+      status: true
+    };
+
+    component.setNoteRevision(data);
+
+    expect(router.navigate).not.toHaveBeenCalled();
+  });
+
+  it('applies interpreter bindings for the active note', () => {
+    const { component } = createComponent();
+
+    const data: MessageReceiveDataTypeMap[OP.INTERPRETER_BINDINGS] = {
+      noteId: 'note-b',
+      interpreterBindings: [{ id: 'current', selected: true }]
+    };
+
+    component.loadInterpreterBindings(data);
+
+    expect(component.interpreterBindings).toEqual([{ id: 'current', selected: 
true }]);
+  });
+
+  it('applies revision history for the active note', () => {
+    const { component } = createComponent();
+
+    const data: MessageReceiveDataTypeMap[OP.LIST_REVISION_HISTORY] = {
+      noteId: 'note-b',
+      revisionList: [{ id: 'rev-1', message: 'current revision' }]
+    };
+
+    component.listRevisionHistory(data);
+
+    expect(component.noteRevisions).toEqual([
+      { id: 'Head', message: 'Head' },
+      { id: 'rev-1', message: 'current revision' }
+    ]);
+
+    expect(component.currentRevision).toBe('Head');
+  });
+
+  it('navigates after setting a revision for the active note', () => {
+    const { component, router } = createComponent();
+
+    const data: MessageReceiveDataTypeMap[OP.SET_NOTE_REVISION] = {
+      noteId: 'note-b',
+      status: true
+    };
+
+    component.setNoteRevision(data);
+
+    expect(router.navigate).toHaveBeenCalledWith(['/notebook', 'note-b']);
+  });
+
+  it('uses the paragraph added envelope msgId for local focus', () => {
+    const { component, messageService } = createComponent();
+
+    component.note = {
+      paragraphs: []
+    } as typeof component.note;
+
+    const message: WebSocketMessage<MessageReceiveDataTypeMap, 
OP.PARAGRAPH_ADDED> = {
+      op: OP.PARAGRAPH_ADDED,
+      msgId: 'local-add-msg',
+      data: {
+        index: 0,
+        paragraph: {
+          id: 'paragraph-1'
+        } as MessageReceiveDataTypeMap[OP.PARAGRAPH_ADDED]['paragraph']
+      }
+    };
+
+    component.addParagraph(message);
+
+    
expect(messageService.consumeLocalAddFocusMsgId).toHaveBeenCalledWith('local-add-msg');
+  });
+});
diff --git 
a/zeppelin-web-angular/src/app/pages/workspace/notebook/notebook.component.ts 
b/zeppelin-web-angular/src/app/pages/workspace/notebook/notebook.component.ts
index 3dd6e73960..cf05e7c0a2 100644
--- 
a/zeppelin-web-angular/src/app/pages/workspace/notebook/notebook.component.ts
+++ 
b/zeppelin-web-angular/src/app/pages/workspace/notebook/notebook.component.ts
@@ -28,7 +28,7 @@ import { distinctUntilChanged, distinctUntilKeyChanged, 
startWith, takeUntil } f
 
 import { NzResizeEvent } from 'ng-zorro-antd/resizable';
 
-import { MessageListener, MessageListenersManager } from '@zeppelin/core';
+import { MessageEnvelopeListener, MessageListener, MessageListenersManager } 
from '@zeppelin/core';
 import { Permissions } from '@zeppelin/interfaces';
 import {
   DynamicFormParams,
@@ -36,7 +36,8 @@ import {
   MessageReceiveDataTypeMap,
   Note,
   OP,
-  RevisionListItem
+  RevisionListItem,
+  WebSocketMessage
 } from '@zeppelin/sdk';
 import {
   MessageService,
@@ -118,6 +119,11 @@ export class NotebookComponent extends 
MessageListenersManager implements OnInit
 
   @MessageListener(OP.INTERPRETER_BINDINGS)
   loadInterpreterBindings(data: 
MessageReceiveDataTypeMap[OP.INTERPRETER_BINDINGS]) {
+    const { noteId } = this.activatedRoute.snapshot.params;
+    if (data.noteId !== noteId) {
+      return;
+    }
+
     this.interpreterBindings = data.interpreterBindings;
     if (!this.interpreterBindings.some(item => item.selected)) {
       this.activatedExtension = 'interpreter';
@@ -146,8 +152,8 @@ export class NotebookComponent extends 
MessageListenersManager implements OnInit
     this.cdr.markForCheck();
   }
 
-  @MessageListener(OP.PARAGRAPH_ADDED)
-  addParagraph(data: MessageReceiveDataTypeMap[OP.PARAGRAPH_ADDED]) {
+  @MessageEnvelopeListener(OP.PARAGRAPH_ADDED)
+  addParagraph(message: WebSocketMessage<MessageReceiveDataTypeMap, 
OP.PARAGRAPH_ADDED>) {
     const { paragraphId } = this.activatedRoute.snapshot.params;
     if (paragraphId || this.revisionView) {
       return;
@@ -155,6 +161,12 @@ export class NotebookComponent extends 
MessageListenersManager implements OnInit
     if (!this.note) {
       return;
     }
+
+    const data = message.data;
+    if (data === undefined) {
+      return;
+    }
+
     const definedNote = this.note;
     definedNote.paragraphs.splice(data.index, 0, data.paragraph);
     const paragraphIndex = definedNote.paragraphs.findIndex(p => p.id === 
data.paragraph.id);
@@ -164,7 +176,7 @@ export class NotebookComponent extends 
MessageListenersManager implements OnInit
 
     // Focus the editor only for a clone/insert initiated by this client (not 
auto-append on run or remote inserts).
     // Defer a tick so the new paragraph's editor child exists, since `focus = 
true` alone misses it.
-    if (this.messageService.consumeLocalAddFocusMsgId(data.msgId)) {
+    if (this.messageService.consumeLocalAddFocusMsgId(message.msgId)) {
       const addedId = data.paragraph.id;
       setTimeout(() => {
         const added = this.listOfNotebookParagraphComponent?.find(e => 
e.paragraph.id === addedId);
@@ -198,8 +210,11 @@ export class NotebookComponent extends 
MessageListenersManager implements OnInit
   }
 
   @MessageListener(OP.SET_NOTE_REVISION)
-  setNoteRevision(_data: MessageReceiveDataTypeMap[OP.SET_NOTE_REVISION]) {
+  setNoteRevision(data: MessageReceiveDataTypeMap[OP.SET_NOTE_REVISION]) {
     const { noteId } = this.activatedRoute.snapshot.params;
+    if (data.noteId !== noteId) {
+      return;
+    }
     this.router.navigate(['/notebook', noteId]).then();
   }
 
@@ -257,6 +272,11 @@ export class NotebookComponent extends 
MessageListenersManager implements OnInit
 
   @MessageListener(OP.LIST_REVISION_HISTORY)
   listRevisionHistory(data: 
MessageReceiveDataTypeMap[OP.LIST_REVISION_HISTORY]) {
+    const { noteId } = this.activatedRoute.snapshot.params;
+    if (data.noteId !== noteId) {
+      return;
+    }
+
     this.noteRevisions = data.revisionList;
     if (this.noteRevisions) {
       if (this.noteRevisions.length === 0 || this.noteRevisions[0].id !== 
'Head') {
diff --git a/zeppelin-web-angular/src/app/services/message.service.ts 
b/zeppelin-web-angular/src/app/services/message.service.ts
index 050d78c44c..937754d400 100644
--- a/zeppelin-web-angular/src/app/services/message.service.ts
+++ b/zeppelin-web-angular/src/app/services/message.service.ts
@@ -23,7 +23,6 @@ import {
   MessageSendDataTypeMap,
   Note,
   NoteConfig,
-  OP,
   ParagraphConfig,
   ParagraphParams,
   PersonalizedMode,
@@ -52,9 +51,6 @@ export class MessageService extends Message implements 
OnDestroy {
 
   interceptReceived(data: WebSocketMessage<MessageReceiveDataTypeMap>): 
WebSocketMessage<MessageReceiveDataTypeMap> {
     const received = this.messageInterceptor ? 
this.messageInterceptor.received(data) : super.interceptReceived(data);
-    if (received.op === OP.PARAGRAPH_ADDED && received.data && received.msgId) 
{
-      (received.data as MessageReceiveDataTypeMap[OP.PARAGRAPH_ADDED]).msgId = 
received.msgId;
-    }
     return received;
   }
 

Reply via email to