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 4d31a06558c..8a2219d3219 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 void getInterpreterBindings(NotebookSocket conn, 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 void saveInterpreterBindings(NotebookSocket conn, ServiceContext context, }); 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 void onSuccess(Revision revision, ServiceContext context) throws IOExcept List 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 @@ private void listRevisionHistory(NotebookSocket conn, @Override public void onSuccess(List 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 @@ private void setNoteRevision(NotebookSocket conn, 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 d288851fbff..5de5b1d10f1 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 @@ void getInterpreterBindingsRequiresReaderPermission() throws IOException { ArgumentCaptor 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 @@ void saveInterpreterBindingsRequiresWriterPermission() throws IOException { notebook.processNote(noteId, Note::getDefaultInterpreterGroup)); ArgumentCaptor 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 @@ void saveInterpreterBindingsRequiresWriterPermission() throws IOException { 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 @@ void testNoteRevision() throws IOException { } } + @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 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 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 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 c59e459410e..f9a4ac4957e 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 665e8dfd71f..fca5c071989 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 0bb02529cc7..37f8d44c2bf 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 4d559a86aa1..f9dff2e804a 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(op: K): Observable[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] + >; + } - // 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[K]>; + receiveEnvelope( + op: K + ): Observable> { + return this.receiveMessage(op) as Observable>; } shortCircuit(message: WebSocketMessage) { @@ -565,4 +559,21 @@ export class Message { formName }); } + + private receiveMessage(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 e12767691d6..7fd1778ae6b 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>(); + const messageService = { + receiveEnvelope: vi.fn(() => received$.asObservable()) + } as unknown as Message; + + class TestComponent extends MessageListenersManager { + receivedMessage?: WebSocketMessage; + + handleNote(message: WebSocketMessage): 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 = { + 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 1b2f0209ae7..93c6294700c 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(op: K) { +type ListenerArgumentsType = T extends undefined ? () => void : (data: T) => void; + +const createMessageListener = ( + op: K, + receiver: (messageService: Message, op: K) => Observable +) => { return function ( target: MessageListenersManager, propertyKey: string, - descriptor: TypedPropertyDescriptor> + descriptor: TypedPropertyDescriptor> ) { - const oldValue = descriptor.value as ReceiveArgumentsType; + const oldValue = descriptor.value as ListenerArgumentsType; const fn = function (this: MessageListenersManager) { if (!this.__zeppelinMessageListeners$__) { @@ -48,7 +53,7 @@ export function MessageListener(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(op: K return descriptor; }; -} +}; + +export const MessageListener = (op: K) => { + return createMessageListener(op, (messageService, targetOp) => messageService.receive(targetOp)); +}; + +export const MessageEnvelopeListener = (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 00000000000..62186de41f4 --- /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 = { + 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 3dd6e739607..cf05e7c0a2a 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) { 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 050d78c44cc..937754d400a 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): WebSocketMessage { 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; }