From 457935dabd52e58ca7ff4423ac05ad7bf1c57f41 Mon Sep 17 00:00:00 2001 From: gyowoo1113 Date: Thu, 3 Sep 2026 14:43:07 +0900 Subject: [PATCH 1/6] [ZEPPELIN-6646] Add websocket payload validation in Message.receive --- .../src/message-payload-guards/index.ts | 28 +++++++++++ .../src/message-payload-guards/job.ts | 47 +++++++++++++++++++ .../projects/zeppelin-sdk/src/message.ts | 11 +++++ 3 files changed, 86 insertions(+) create mode 100644 zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts create mode 100644 zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/job.ts diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts b/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts new file mode 100644 index 00000000000..b5b8a9a78ef --- /dev/null +++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts @@ -0,0 +1,28 @@ +/* + * 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 type { MessageReceiveDataTypeMap } from '../interfaces/message-data-type-map.interface'; +import { OP } from '../interfaces/message-operator.interface'; + +import { isListUpdateNoteJobsPayload } from './job'; + +export type MessagePayloadGuard = (value: unknown) => boolean; + +type ReceiveOP = keyof MessageReceiveDataTypeMap; + +const MESSAGE_PAYLOAD_GUARDS: Partial> = { + [OP.LIST_UPDATE_NOTE_JOBS]: isListUpdateNoteJobsPayload +}; + +export const getMessagePayloadGuard = (op: ReceiveOP): MessagePayloadGuard | undefined => { + return MESSAGE_PAYLOAD_GUARDS[op]; +}; diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/job.ts b/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/job.ts new file mode 100644 index 00000000000..f3edecd093b --- /dev/null +++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/job.ts @@ -0,0 +1,47 @@ +/* + * 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. + */ + +const isRecord = (value: unknown): value is Record => typeof value === 'object' && value !== null; + +const isJobUpdate = (value: unknown): boolean => { + if (!isRecord(value)) { + return false; + } + + if (typeof value.noteId !== 'string') { + return false; + } + + if (typeof value.isRemoved !== 'boolean') { + return false; + } + + if (value.isRemoved) { + return true; + } + + return typeof value.noteName === 'string'; +}; + +export const isListUpdateNoteJobsPayload = (value: unknown): boolean => { + if (!isRecord(value)) { + return false; + } + + const noteRunningJobs = value.noteRunningJobs; + + if (!isRecord(noteRunningJobs)) { + return false; + } + + return Array.isArray(noteRunningJobs.jobs) && noteRunningJobs.jobs.every(isJobUpdate); +}; diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts index 0f070c6354f..6274c97f5ba 100644 --- a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts +++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts @@ -30,6 +30,8 @@ import { } from './interfaces/message-paragraph.interface'; import { WebSocketMessage } from './interfaces/websocket-message.interface'; +import { getMessagePayloadGuard } from './message-payload-guards'; + export type ArgumentsType = T extends (...args: infer U) => void ? U : never; export type SendArgumentsType = MessageSendDataTypeMap[K] extends undefined @@ -175,6 +177,15 @@ export class Message { receive(op: K): Observable[K]> { return this.received$.pipe( filter(message => message.op === op), + filter(message => { + const guard = getMessagePayloadGuard(op); + + if (!guard) { + return true; + } + + return guard(message.data); + }), map(message => message.data) ) as Observable[K]>; } From 83029642ffbbac288eb4544ccac1cde37c959407 Mon Sep 17 00:00:00 2001 From: gyowoo1113 Date: Thu, 3 Sep 2026 14:55:07 +0900 Subject: [PATCH 2/6] [ZEPPELIN-6646] Add Message.receive payload validation tests --- .../projects/zeppelin-sdk/src/message.spec.ts | 136 ++++++++++++++++++ 1 file changed, 136 insertions(+) create mode 100644 zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts new file mode 100644 index 00000000000..256c30029e0 --- /dev/null +++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts @@ -0,0 +1,136 @@ +/* + * 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 { describe, expect, it, vi } from 'vitest'; + +import type { MessageReceiveDataTypeMap } from './interfaces/message-data-type-map.interface'; +import { OP } from './interfaces/message-operator.interface'; +import type { WebSocketMessage } from './interfaces/websocket-message.interface'; +import { Message } from './message'; + +const asReceivedMessage = (message: unknown): WebSocketMessage => + message as WebSocketMessage; + +describe('Message.receive', () => { + it('passes a non-removal job update with noteName', () => { + const message = new Message(); + const listener = vi.fn(); + const data = { + noteRunningJobs: { + jobs: [ + { + noteId: 'note-1', + noteName: 'Test Note', + isRemoved: false + } + ] + } + }; + + message.receive(OP.LIST_UPDATE_NOTE_JOBS).subscribe(listener); + + message.shortCircuit( + asReceivedMessage({ + op: OP.LIST_UPDATE_NOTE_JOBS, + data + }) + ); + + expect(listener).toHaveBeenCalledWith(data); + }); + + it('passes a partial removal payload without noteName', () => { + const message = new Message(); + const listener = vi.fn(); + const data = { + noteRunningJobs: { + jobs: [ + { + noteId: 'note-1', + isRemoved: true + } + ] + } + }; + + message.receive(OP.LIST_UPDATE_NOTE_JOBS).subscribe(listener); + + message.shortCircuit( + asReceivedMessage({ + op: OP.LIST_UPDATE_NOTE_JOBS, + data + }) + ); + + expect(listener).toHaveBeenCalledWith(data); + }); + + it('filters a non-removal job update without noteName', () => { + const message = new Message(); + const listener = vi.fn(); + + message.receive(OP.LIST_UPDATE_NOTE_JOBS).subscribe(listener); + + message.shortCircuit( + asReceivedMessage({ + op: OP.LIST_UPDATE_NOTE_JOBS, + data: { + noteRunningJobs: { + jobs: [ + { + noteId: 'note-1', + isRemoved: false + } + ] + } + } + }) + ); + + expect(listener).not.toHaveBeenCalled(); + }); + + it('filters a payload without a jobs array', () => { + const message = new Message(); + const listener = vi.fn(); + + message.receive(OP.LIST_UPDATE_NOTE_JOBS).subscribe(listener); + + message.shortCircuit( + asReceivedMessage({ + op: OP.LIST_UPDATE_NOTE_JOBS, + data: { + noteRunningJobs: {} + } + }) + ); + + expect(listener).not.toHaveBeenCalled(); + }); + + it('keeps existing behavior for an OP without a guard', () => { + const message = new Message(); + const listener = vi.fn(); + const data = {}; + + message.receive(OP.NOTE).subscribe(listener); + + message.shortCircuit( + asReceivedMessage({ + op: OP.NOTE, + data + }) + ); + + expect(listener).toHaveBeenCalledWith(data); + }); +}); From d7029a8a3e70a3d2f9dd98a832fb890ce6f4ea0e Mon Sep 17 00:00:00 2001 From: gyowoo1113 Date: Thu, 3 Sep 2026 15:41:07 +0900 Subject: [PATCH 3/6] [ZEPPELIN-6646] Log MessageListener handler errors with OP context --- .../src/app/core/message-listener/message-listener.ts | 8 ++++++-- 1 file changed, 6 insertions(+), 2 deletions(-) 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 6487124ecc7..c05a06f5d55 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 @@ -49,8 +49,12 @@ export function MessageListener(op: K this.__zeppelinMessageListeners$__.add( this.messageService.receive(op).subscribe(data => { - // @ts-ignore - oldValue.apply(this, [data]); + try { + // @ts-ignore + oldValue.apply(this, [data]); + } catch (error) { + console.error(`Failed to handle WebSocket OP ${String(op)}`, error); + } }) ); }; From af501838b731643a2941e3a89a81937f22e58716 Mon Sep 17 00:00:00 2001 From: gyowoo1113 Date: Thu, 3 Sep 2026 15:42:26 +0900 Subject: [PATCH 4/6] [ZEPPELIN-6646] Add tests for MessageListener error logging --- .../message-listener/message-listener.spec.ts | 85 +++++++++++++++++++ 1 file changed, 85 insertions(+) create mode 100644 zeppelin-web-angular/src/app/core/message-listener/message-listener.spec.ts 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 new file mode 100644 index 00000000000..b5dcf91563a --- /dev/null +++ b/zeppelin-web-angular/src/app/core/message-listener/message-listener.spec.ts @@ -0,0 +1,85 @@ +/* + * 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 { Subject } from 'rxjs'; +import { afterEach, describe, expect, it, vi } from 'vitest'; + +import { Message, OP, MessageReceiveDataTypeMap } from '@zeppelin/sdk'; + +import { MessageListener, MessageListenersManager } from './message-listener'; + +afterEach(() => { + vi.restoreAllMocks(); +}); + +describe('MessageListener', () => { + it('logs handler errors with the OP and keeps the subscription active', () => { + const received$ = new Subject(); + const messageService = { + receive: vi.fn(() => received$.asObservable()) + } as unknown as Message; + + const error = new Error('boom'); + const consoleError = vi.spyOn(console, 'error').mockImplementation(() => {}); + + class TestComponent extends MessageListenersManager { + calls = 0; + + handleNote(_data: MessageReceiveDataTypeMap[OP.NOTE]): void { + this.calls++; + + if (this.calls === 1) { + throw error; + } + } + } + + const descriptor = Object.getOwnPropertyDescriptor(TestComponent.prototype, 'handleNote')!; + + MessageListener(OP.NOTE)(TestComponent.prototype, 'handleNote', descriptor); + + const component = new TestComponent(messageService); + const data = {} as MessageReceiveDataTypeMap[OP.NOTE]; + + received$.next(data); + received$.next(data); + + expect(component.calls).toBe(2); + expect(consoleError).toHaveBeenCalledWith(`Failed to handle WebSocket OP ${String(OP.NOTE)}`, error); + }); + + it('passes received data to the handler', () => { + const received$ = new Subject(); + const messageService = { + receive: vi.fn(() => received$.asObservable()) + } as unknown as Message; + + class TestComponent extends MessageListenersManager { + receivedData?: MessageReceiveDataTypeMap[OP.NOTE]; + + handleNote(data: MessageReceiveDataTypeMap[OP.NOTE]): void { + this.receivedData = data; + } + } + + const descriptor = Object.getOwnPropertyDescriptor(TestComponent.prototype, 'handleNote')!; + + MessageListener(OP.NOTE)(TestComponent.prototype, 'handleNote', descriptor); + + const component = new TestComponent(messageService); + const data = {} as MessageReceiveDataTypeMap[OP.NOTE]; + + received$.next(data); + + expect(component.receivedData).toBe(data); + }); +}); From ce8440433d37cc2f8d77b7b37017b2ce3d3f18a3 Mon Sep 17 00:00:00 2001 From: gyowoo1113 Date: Thu, 3 Sep 2026 16:31:56 +0900 Subject: [PATCH 5/6] [ZEPPELIN-6646] Document payload guard registration convention --- .../zeppelin-sdk/src/message-payload-guards/index.ts | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts b/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts index b5b8a9a78ef..3ab759d6ce9 100644 --- a/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts +++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message-payload-guards/index.ts @@ -19,6 +19,11 @@ export type MessagePayloadGuard = (value: unknown) => boolean; type ReceiveOP = keyof MessageReceiveDataTypeMap; +/** + * Runtime payload guards are registered only for OPs with a demonstrated + * payload-shape failure. Add new guards when a concrete runtime failure + * shows that validation is needed. + */ const MESSAGE_PAYLOAD_GUARDS: Partial> = { [OP.LIST_UPDATE_NOTE_JOBS]: isListUpdateNoteJobsPayload }; From 19482b65c991690c589c413cb8b3cd65108856d5 Mon Sep 17 00:00:00 2001 From: gyowoo1113 Date: Sat, 5 Sep 2026 00:28:03 +0900 Subject: [PATCH 6/6] [ZEPPELIN-6646] Log dropped payloads and rethrow listener errors with OP context --- .../projects/zeppelin-sdk/src/message.spec.ts | 14 ++++++++++++-- .../projects/zeppelin-sdk/src/message.ts | 10 ++++++---- .../core/message-listener/message-listener.spec.ts | 7 ++++++- .../app/core/message-listener/message-listener.ts | 1 + 4 files changed, 25 insertions(+), 7 deletions(-) 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 256c30029e0..0bb02529cc7 100644 --- a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts +++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.spec.ts @@ -10,7 +10,7 @@ * limitations under the License. */ -import { describe, expect, it, vi } from 'vitest'; +import { afterEach, describe, expect, it, vi } from 'vitest'; import type { MessageReceiveDataTypeMap } from './interfaces/message-data-type-map.interface'; import { OP } from './interfaces/message-operator.interface'; @@ -20,6 +20,10 @@ import { Message } from './message'; const asReceivedMessage = (message: unknown): WebSocketMessage => message as WebSocketMessage; +afterEach(() => { + vi.restoreAllMocks(); +}); + describe('Message.receive', () => { it('passes a non-removal job update with noteName', () => { const message = new Message(); @@ -74,9 +78,10 @@ describe('Message.receive', () => { expect(listener).toHaveBeenCalledWith(data); }); - it('filters a non-removal job update without noteName', () => { + it('filters a non-removal job update without noteName and warns with the OP only', () => { const message = new Message(); const listener = vi.fn(); + const consoleWarn = vi.spyOn(console, 'warn').mockImplementation(() => {}); message.receive(OP.LIST_UPDATE_NOTE_JOBS).subscribe(listener); @@ -97,11 +102,16 @@ describe('Message.receive', () => { ); expect(listener).not.toHaveBeenCalled(); + expect(consoleWarn).toHaveBeenCalledTimes(1); + expect(consoleWarn).toHaveBeenCalledWith( + `Dropped WebSocket OP ${String(OP.LIST_UPDATE_NOTE_JOBS)}: payload failed validation` + ); }); it('filters a payload without a jobs array', () => { const message = new Message(); const listener = vi.fn(); + vi.spyOn(console, 'warn').mockImplementation(() => {}); message.receive(OP.LIST_UPDATE_NOTE_JOBS).subscribe(listener); diff --git a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts index 6274c97f5ba..4d559a86aa1 100644 --- a/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts +++ b/zeppelin-web-angular/projects/zeppelin-sdk/src/message.ts @@ -175,16 +175,18 @@ export class Message { } receive(op: K): Observable[K]> { + const guard = getMessagePayloadGuard(op); + return this.received$.pipe( filter(message => message.op === op), filter(message => { - const guard = getMessagePayloadGuard(op); - - if (!guard) { + if (!guard || guard(message.data)) { return true; } - return guard(message.data); + // 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]>; 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 b5dcf91563a..e12767691d6 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 @@ -19,10 +19,14 @@ import { MessageListener, MessageListenersManager } from './message-listener'; afterEach(() => { vi.restoreAllMocks(); + vi.useRealTimers(); }); describe('MessageListener', () => { - it('logs handler errors with the OP and keeps the subscription active', () => { + it('logs handler errors with the OP, rethrows them, and keeps the subscription active', () => { + // RxJS rethrows an error thrown inside `next` from a timer, so the subscription itself survives. + vi.useFakeTimers(); + const received$ = new Subject(); const messageService = { receive: vi.fn(() => received$.asObservable()) @@ -55,6 +59,7 @@ describe('MessageListener', () => { expect(component.calls).toBe(2); expect(consoleError).toHaveBeenCalledWith(`Failed to handle WebSocket OP ${String(OP.NOTE)}`, error); + expect(() => vi.runAllTimers()).toThrow(error); }); it('passes received data to the handler', () => { 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 c05a06f5d55..1b2f0209ae7 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 @@ -54,6 +54,7 @@ export function MessageListener(op: K oldValue.apply(this, [data]); } catch (error) { console.error(`Failed to handle WebSocket OP ${String(op)}`, error); + throw error; } }) );