This is an automated email from the ASF dual-hosted git repository. voidmatcha pushed a commit to branch ZEPPELIN-6666-transport-fixtures in repository https://gitbox.apache.org/repos/asf/zeppelin.git
commit 033d857848b4cd069984b89da8ca07d19fc20be0 Author: YONGJAE LEE <[email protected]> AuthorDate: Thu Sep 3 01:00:28 2026 +0900 [ZEPPELIN-6666] Add notebook transport fixtures --- .../e2e/core-contract/capture-server.sh | 181 +++++ .../e2e/core-contract/capture-server.test.mjs | 830 +++++++++++++++++++++ .../e2e/core-contract/capture-stub-zeppelin.mjs | 60 ++ .../core-contract/notebook-transport-fixture.mjs | 699 +++++++++++++++++ .../core-contract/capture-fixtures.spec.ts | 503 +++++++++++++ zeppelin-web-angular/package.json | 1 + zeppelin-web-angular/pom.xml | 12 + 7 files changed, 2286 insertions(+) diff --git a/zeppelin-web-angular/e2e/core-contract/capture-server.sh b/zeppelin-web-angular/e2e/core-contract/capture-server.sh new file mode 100755 index 0000000000..d8062a7c25 --- /dev/null +++ b/zeppelin-web-angular/e2e/core-contract/capture-server.sh @@ -0,0 +1,181 @@ +#!/usr/bin/env bash +# +# Licensed to the Apache Software Foundation (ASF) under one or more +# contributor license agreements. See the NOTICE file distributed with +# this work for additional information regarding copyright ownership. +# The ASF licenses this file to You 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. + +set -euo pipefail + +usage() { + echo "usage: $0 start|stop --root <dir> [--mode anonymous|auth] [--port <port>]" >&2 +} + +command="${1:-}" +shift || true +capture_root="" +capture_mode="anonymous" +zeppelin_port="8080" + +while [[ $# -gt 0 ]]; do + case "$1" in + --root) + capture_root="${2:-}" + shift 2 + ;; + --mode) + capture_mode="${2:-}" + shift 2 + ;; + --port) + zeppelin_port="${2:-}" + shift 2 + ;; + *) + usage + exit 2 + ;; + esac +done + +if [[ -z "${command}" || -z "${capture_root}" ]]; then + usage + exit 2 +fi + +repo_root="$(cd "$(dirname "$0")/../../.." && pwd)" +capture_root="$(mkdir -p "${capture_root}" && cd "${capture_root}" && pwd)" +marker_file="${capture_root}/.zeppelin-capture-root" +zeppelin_pid_file="${capture_root}/zeppelin.pid" + +port_in_use() { + lsof -nP -iTCP:"$1" -sTCP:LISTEN >/dev/null 2>&1 +} + +write_marker() { + { + echo "root=${capture_root}" + echo "repo=${repo_root}" + } > "${marker_file}" +} + +verify_root_marker() { + [[ -f "${marker_file}" ]] && grep -qx "root=${capture_root}" "${marker_file}" +} + +verify_pid_identity() { + local pid="$1" + local expected="$2" + [[ "${pid}" =~ ^[0-9]+$ ]] || return 1 + ps -p "${pid}" -o command= | grep -F -- "${expected}" >/dev/null 2>&1 +} + +stop_pid() { + local pid_file="$1" + local expected="$2" + [[ -f "${pid_file}" ]] || return 0 + local pid + pid="$(cat "${pid_file}")" + if ps -p "${pid}" >/dev/null 2>&1; then + if ! verify_pid_identity "${pid}" "${expected}"; then + echo "refusing to stop ${pid}: command does not match ${expected}" >&2 + exit 1 + fi + kill "${pid}" + for _ in {1..20}; do + ps -p "${pid}" >/dev/null 2>&1 || break + sleep 1 + done + fi + rm -f "${pid_file}" +} + +start_zeppelin() { + mkdir -p "${capture_root}/conf" "${capture_root}/notebook" "${capture_root}/index" \ + "${capture_root}/logs" "${capture_root}/run" "${capture_root}/recovery" "${capture_root}/webapps" + cp "${repo_root}/conf/log4j2.properties" "${capture_root}/conf/log4j2.properties" + cp "${repo_root}/conf/zeppelin-site.xml.template" "${capture_root}/conf/zeppelin-site.xml" + if [[ "${capture_mode}" == "auth" ]]; then + cp "${repo_root}/conf/shiro.ini.template" "${capture_root}/conf/shiro.ini" + else + rm -f "${capture_root}/conf/shiro.ini" + fi + + export ZEPPELIN_CONF_DIR="${capture_root}/conf" + export ZEPPELIN_NOTEBOOK_DIR="${capture_root}/notebook" + export ZEPPELIN_LOG_DIR="${capture_root}/logs" + export ZEPPELIN_PID_DIR="${capture_root}/run" + export ZEPPELIN_WAR_TEMPDIR="${capture_root}/webapps" + export ZEPPELIN_JAVA_OPTS="${ZEPPELIN_JAVA_OPTS:-} -Dzeppelin.server.port=${zeppelin_port} -Dzeppelin.notebook.dir=${capture_root}/notebook -Dzeppelin.search.index.path=${capture_root}/index -Dzeppelin.recovery.dir=${capture_root}/recovery -Dzeppelin.capture.root=${capture_root}" + export ZEPPELIN_CAPTURE_ROOT="${capture_root}" + export ZEPPELIN_PORT="${zeppelin_port}" + + if [[ -n "${CAPTURE_ZEPPELIN_COMMAND:-}" ]]; then + bash -c "${CAPTURE_ZEPPELIN_COMMAND}" >"${capture_root}/logs/zeppelin-stdout.log" 2>"${capture_root}/logs/zeppelin-stderr.log" & + echo "$!" > "${zeppelin_pid_file}" + else + "${repo_root}/bin/zeppelin.sh" >"${capture_root}/logs/zeppelin-stdout.log" 2>"${capture_root}/logs/zeppelin-stderr.log" & + echo "$!" > "${zeppelin_pid_file}" + fi +} + +wait_for_http() { + local url="$1" + for _ in {1..120}; do + if curl -fsS "${url}" >/dev/null 2>&1; then + return 0 + fi + sleep 1 + done + return 1 +} + +start_server() { + if [[ "${capture_mode}" != "anonymous" && "${capture_mode}" != "auth" ]]; then + echo "mode must be anonymous or auth" >&2 + exit 2 + fi + if port_in_use "${zeppelin_port}"; then + echo "port ${zeppelin_port} is already in use" >&2 + exit 1 + fi + + write_marker + start_zeppelin + if ! wait_for_http "http://127.0.0.1:${zeppelin_port}/api/version"; then + echo "zeppelin did not become ready on port ${zeppelin_port}" >&2 + stop_server + exit 1 + fi +} + +stop_server() { + if ! verify_root_marker; then + echo "refusing to stop without matching capture root marker: ${marker_file}" >&2 + exit 1 + fi + stop_pid "${zeppelin_pid_file}" "${capture_root}" +} + +case "${command}" in + start) + start_server + ;; + stop) + stop_server + ;; + *) + usage + exit 2 + ;; +esac diff --git a/zeppelin-web-angular/e2e/core-contract/capture-server.test.mjs b/zeppelin-web-angular/e2e/core-contract/capture-server.test.mjs new file mode 100644 index 0000000000..8d3f4171ab --- /dev/null +++ b/zeppelin-web-angular/e2e/core-contract/capture-server.test.mjs @@ -0,0 +1,830 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 assert from 'node:assert/strict'; +import { spawnSync } from 'node:child_process'; +import { EventEmitter } from 'node:events'; +import { existsSync, mkdtempSync, readFileSync } from 'node:fs'; +import http from 'node:http'; +import os from 'node:os'; +import path from 'node:path'; +import test from 'node:test'; + +import * as fixtureModule from './notebook-transport-fixture.mjs'; +import { + createNotebookTransportRecorder, + createPlaywrightFixtureAdapter, + fixtureVersion, + isNotebookRestUrl, + parseRestBody, + validateFixture, + webSocketPayloadMatches +} from './notebook-transport-fixture.mjs'; + +const script = path.resolve('e2e/core-contract/capture-server.sh'); +const stub = path.resolve('e2e/core-contract/capture-stub-zeppelin.mjs'); + +test('capture-server start and stop do not require a fixture arg', () => { + const root = createRoot(); + + start(root); + stop(root); + + assert.equal(existsSync(path.join(root.root, '.zeppelin-capture-root')), true); + assert.equal(existsSync(path.join(root.root, 'zeppelin.pid')), false); +}); + +test('capture-server writes anonymous and auth config in an isolated temp root', () => { + const anonymous = createRoot(); + const auth = createRoot(); + + start(anonymous); + stop(anonymous); + start(auth, { mode: 'auth' }); + stop(auth); + + assert.equal(existsSync(path.join(anonymous.root, 'conf/shiro.ini')), false); + assert.equal(existsSync(path.join(auth.root, 'conf/shiro.ini')), true); +}); + +test('capture-server reports explicit port conflicts', async () => { + const server = await listen(); + const root = createRoot(); + + try { + const result = run(['start', '--root', root.root, '--port', String(server.address().port)]); + + assert.equal(result.status, 1); + assert.match(result.stderr, /port .* is already in use/); + } finally { + await close(server); + } +}); + +test('recorder captures notebook REST and WebSocket browser events only', async () => { + const page = new EventEmitter(); + const socket = new EventEmitter(); + socket.url = () => 'http://127.0.0.1:8080/ws'; + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + page.emit('request', request('POST', 'http://127.0.0.1:8080/api/notebook/note-a', '{"msgId":"runtime"}')); + page.emit('response', response(request('GET', 'http://127.0.0.1:8080/api/notebook/note-a'), 200, '{"id":"note-a"}')); + page.emit('request', request('GET', 'http://127.0.0.1:8080/assets/app.js')); + page.emit('websocket', socket); + socket.emit('framesent', { payload: '{"op":"GET_NOTE","msgId":"runtime"}' }); + socket.emit('framereceived', { payload: '{"op":"NOTE","noteId":"note-a"}' }); + + await recorder.stop(); + const fixture = recorder.snapshot(); + + assert.deepEqual(validateFixture(fixture), []); + assert.deepEqual( + fixture.records.map(record => [ + record.sequence, + record.kind, + record.rest?.direction ?? record.websocket?.direction + ]), + [ + [1, 'rest', 'request'], + [2, 'rest', 'response'], + [3, 'websocket', 'send'], + [4, 'websocket', 'receive'] + ] + ); + assert.equal(fixture.records[0].rest.request.bodyJson.msgId, '<msgId>'); + assert.equal(fixture.records[1].rest.bodyJson.id, 'note-a'); +}); + +test('recorder redacts sensitive headers and fields before writing fixture files', async () => { + const root = createRoot(); + const page = new EventEmitter(); + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + page.emit( + 'request', + request('POST', 'http://127.0.0.1:8080/api/notebook', '{"ticket":"secret","id":"stable-id"}', { + accept: 'application/json', + authorization: 'Bearer secret', + cookie: 'ticket=secret', + 'content-type': 'application/json' + }) + ); + const written = await recorder.write(path.join(root.root, 'fixtures/notebook-transport.json')); + const text = readFileSync(path.join(root.root, 'fixtures/notebook-transport.json'), 'utf8'); + + assert.equal(text.includes('Bearer secret'), false); + assert.equal(text.includes('ticket=secret'), false); + assert.deepEqual(written.records[0].rest.request.headers, { + accept: 'application/json', + 'content-type': 'application/json' + }); + assert.deepEqual(written.records[0].rest.request.bodyJson, { id: 'stable-id', ticket: '<ticket>' }); +}); + +test('recorder redacts WebSocket JSON payload secrets before writing fixture files', async () => { + const root = createRoot(); + const page = new EventEmitter(); + const socket = new EventEmitter(); + socket.url = () => 'http://127.0.0.1:8080/ws'; + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + page.emit('websocket', socket); + socket.emit('framesent', { + payload: + '{"op":"GET_NOTE","id":"stable-id","noteId":"note-a","ticket":"secret-ticket","principal":"alice","msgId":"runtime"}' + }); + + await recorder.write(path.join(root.root, 'fixtures/notebook-transport.json')); + const text = readFileSync(path.join(root.root, 'fixtures/notebook-transport.json'), 'utf8'); + const fixture = JSON.parse(text); + + assert.equal(text.includes('secret-ticket'), false); + assert.equal(text.includes('alice'), false); + assert.equal(text.includes('runtime'), false); + assert.equal(fixture.records[0].websocket.payloadText.includes('"op":"GET_NOTE"'), true); + assert.equal(fixture.records[0].websocket.payloadText.includes('"id":"stable-id"'), true); + assert.equal(fixture.records[0].websocket.payloadText.includes('"noteId":"note-a"'), true); +}); + +test('recorder normalizes sensitive and volatile REST URL query values before writing', async () => { + const root = createRoot(); + const page = new EventEmitter(); + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + page.emit( + 'request', + request( + 'GET', + 'http://127.0.0.1:8080/api/notebook/note-a?ticket=secret-ticket&token=secret-token&msgId=runtime&view=stable' + ) + ); + + await recorder.write(path.join(root.root, 'fixtures/notebook-transport.json')); + const text = readFileSync(path.join(root.root, 'fixtures/notebook-transport.json'), 'utf8'); + const fixture = JSON.parse(text); + + assert.equal(text.includes('secret-ticket'), false); + assert.equal(text.includes('secret-token'), false); + assert.equal( + fixture.records[0].rest.request.url, + '/api/notebook/note-a?ticket=%3Cticket%3E&token=%3Ctoken%3E&msgId=%3CmsgId%3E&view=stable' + ); +}); + +test('recorder redacts credential-shaped body and query fields before writing', async () => { + const root = createRoot(); + const page = new EventEmitter(); + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + page.emit( + 'request', + request( + 'POST', + 'http://127.0.0.1:8080/api/notebook/note-a?apiKey=query-key&clientSecret=query-secret&view=stable', + '{"apiKey":"body-key","credential":"body-credential","secret":"body-secret","id":"stable-id"}', + { 'content-type': 'application/json' } + ) + ); + + await recorder.write(path.join(root.root, 'fixtures/notebook-transport.json')); + const text = readFileSync(path.join(root.root, 'fixtures/notebook-transport.json'), 'utf8'); + const fixture = JSON.parse(text); + + for (const value of ['query-key', 'query-secret', 'body-key', 'body-credential', 'body-secret']) { + assert.equal(text.includes(value), false); + } + assert.deepEqual(fixture.records[0].rest.request.bodyJson, { + apiKey: '<apiKey>', + credential: '<credential>', + id: 'stable-id', + secret: '<secret>' + }); + assert.equal( + fixture.records[0].rest.request.url, + '/api/notebook/note-a?apiKey=%3CapiKey%3E&clientSecret=%3CclientSecret%3E&view=stable' + ); +}); + +test('replay applies the same query normalization and ignores JSON property order', async () => { + const fulfilled = []; + const adapter = fixtureModule.createFixtureReplayAdapter({ + records: [ + { + kind: 'rest', + sequence: 1, + rest: { + bodyJson: { id: 'note-a' }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { + bodyRaw: '', + headers: { accept: 'application/json' }, + method: 'GET', + url: '/api/notebook/note-a?ticket=%3Cticket%3E&msgId=%3CmsgId%3E' + }, + status: 200 + } + } + ], + version: fixtureVersion + }); + + await adapter.route( + { fulfill: async value => fulfilled.push(value) }, + request('GET', 'http://127.0.0.1:8080/api/notebook/note-a?ticket=runtime-secret&msgId=runtime-id') + ); + + assert.equal(fulfilled.length, 1); + assert.equal( + webSocketPayloadMatches('{"op":"GET_NOTE","msgId":"<msgId>"}', '{"msgId":"runtime","op":"GET_NOTE"}'), + true + ); +}); + +test('replay rejects REST requests when the recorded request body or safe headers differ', async () => { + const adapter = fixtureModule.createFixtureReplayAdapter({ + records: [ + { + kind: 'rest', + sequence: 1, + rest: { + bodyJson: { id: 'note-a' }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { + bodyJson: { paragraphId: 'paragraph-a', text: 'print(1)' }, + headers: { accept: 'application/json', 'content-type': 'application/json' }, + method: 'POST', + url: '/api/notebook/job/note-a/paragraph-a' + }, + status: 200 + } + } + ], + version: fixtureVersion + }); + + await assert.rejects( + () => + adapter.route( + { fulfill: async () => undefined }, + request( + 'POST', + 'http://127.0.0.1:8080/api/notebook/job/note-a/paragraph-a', + '{"paragraphId":"paragraph-a","text":"print(2)"}', + { + accept: 'application/json', + 'content-type': 'application/json' + } + ) + ), + /REST fixture request mismatch/ + ); + + await assert.rejects( + () => + adapter.route( + { fulfill: async () => undefined }, + request( + 'POST', + 'http://127.0.0.1:8080/api/notebook/job/note-a/paragraph-a', + '{"text":"print(1)","paragraphId":"paragraph-a"}', + { + accept: 'text/plain', + 'content-type': 'application/json' + } + ) + ), + /REST fixture request mismatch/ + ); +}); + +test('recorder fails closed when response body capture fails', async () => { + const page = new EventEmitter(); + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + page.emit( + 'response', + response(request('GET', 'http://127.0.0.1:8080/api/notebook/note-a'), 200, async () => { + throw new Error('body unavailable'); + }) + ); + + await assert.rejects(() => recorder.stop(), /body unavailable/); +}); + +test('Playwright adapter replays WebSocket payloadBase64 as binary data', async () => { + const calls = []; + const page = { + route: async () => undefined, + routeWebSocket: async (_pattern, handler) => calls.push(handler) + }; + await createPlaywrightFixtureAdapter({ + records: [ + { + kind: 'websocket', + sequence: 1, + websocket: { direction: 'send', payloadText: 'client-ready' } + }, + { + kind: 'websocket', + sequence: 2, + websocket: { direction: 'receive', payloadBase64: Buffer.from([0, 255, 1, 2]).toString('base64') } + } + ], + version: fixtureVersion + }).install(page); + + const replies = []; + const handlers = []; + calls[0]({ + onMessage: handler => handlers.push(handler), + send: message => replies.push(message) + }); + handlers[0]('client-ready'); + + assert.equal(Buffer.isBuffer(replies[0]), true); + assert.deepEqual([...replies[0]], [0, 255, 1, 2]); +}); + +test('WebSocket binary payload matching compares bytes instead of UTF-8 replacement text', () => { + assert.equal(webSocketPayloadMatches(Buffer.from([0xff]), Buffer.from([0xff])), true); + assert.equal(webSocketPayloadMatches(Buffer.from([0xff]), Buffer.from([0xfe])), false); + assert.equal(webSocketPayloadMatches(Buffer.from([0xff]), '\ufffd'), false); +}); + +test('recorder write waits for pending response body capture', async () => { + const root = createRoot(); + const page = new EventEmitter(); + const recorder = createNotebookTransportRecorder(); + let resolveBody; + + recorder.install(page); + page.emit( + 'response', + response(request('GET', 'http://127.0.0.1:8080/api/notebook/note-a'), 200, () => { + return new Promise(resolve => { + resolveBody = resolve; + }); + }) + ); + const writePromise = recorder.write(path.join(root.root, 'fixtures/notebook-transport.json')); + resolveBody('{"id":"note-a"}'); + const written = await writePromise; + + assert.deepEqual(written.records[0].rest.bodyJson, { id: 'note-a' }); +}); + +test('notebook REST predicate excludes unrelated API traffic', async () => { + assert.equal(isNotebookRestUrl('http://127.0.0.1:8080/api/notebook'), true); + assert.equal(isNotebookRestUrl('http://127.0.0.1:8080/api/notebook/note-a'), true); + assert.equal(isNotebookRestUrl('http://127.0.0.1:8080/api/security/ticket'), false); + assert.equal(isNotebookRestUrl('http://127.0.0.1:8080/api/configurations/all'), false); +}); + +test('REST body parsing preserves raw non-JSON and parses JSON-looking bodies for normalization', () => { + assert.deepEqual(parseRestBody('plain text', { 'content-type': 'text/plain' }), { bodyRaw: 'plain text' }); + assert.deepEqual(parseRestBody('{"noteId":"note-a","stable":true}', { 'content-type': 'application/json' }), { + bodyJson: { noteId: 'note-a', stable: true } + }); +}); + +test('recorder redacts sensitive values from malformed JSON REST and WebSocket payloads', async () => { + const root = createRoot(); + const page = new EventEmitter(); + const socket = new EventEmitter(); + socket.url = () => 'http://127.0.0.1:8080/ws'; + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + page.emit('request', request('POST', 'http://127.0.0.1:8080/api/notebook', '{"token":"rest-secret')); + page.emit('websocket', socket); + socket.emit('framesent', { payload: '{"credential":"socket-secret' }); + + const fixture = await recorder.write(path.join(root.root, 'fixtures/notebook-transport.json')); + const serialized = JSON.stringify(fixture); + assert.equal(serialized.includes('rest-secret'), false); + assert.equal(serialized.includes('socket-secret'), false); +}); + +test('Playwright adapter replays WebSocket fixtures with cursors and no server forwarding by default', async () => { + const calls = []; + const page = { + route: async () => undefined, + routeWebSocket: async (_pattern, handler) => calls.push(handler) + }; + await createPlaywrightFixtureAdapter({ + records: [ + wsRecord(1, 'send', '{"op":"GET_NOTE","msgId":"<msgId>"}'), + wsRecord(2, 'receive', '{"op":"NOTE","noteId":"note-a"}'), + wsRecord(3, 'send', '{"op":"RUN_PARAGRAPH","paragraphId":"paragraph-a"}'), + wsRecord(4, 'receive', '{"op":"PARAGRAPH","paragraphId":"paragraph-a"}') + ], + version: fixtureVersion + }).install(page); + + const replies = []; + const forwarded = []; + const handlers = []; + calls[0]({ + connectToServer: () => ({ send: message => forwarded.push(message) }), + onMessage: handler => handlers.push(handler), + send: message => replies.push(message) + }); + + handlers[0]('{"op":"GET_NOTE","msgId":"runtime"}'); + handlers[0]('{"op":"RUN_PARAGRAPH","paragraphId":"paragraph-a"}'); + + assert.deepEqual(forwarded, []); + assert.deepEqual(replies, ['{"op":"NOTE","noteId":"note-a"}', '{"op":"PARAGRAPH","paragraphId":"paragraph-a"}']); + assert.throws(() => handlers[0]('{"op":"EXTRA"}'), /messages exhausted/); +}); + +test('Playwright adapter emits server-first messages and reports unconsumed messages', async () => { + const calls = []; + const page = { + route: async () => undefined, + routeWebSocket: async (_pattern, handler) => calls.push(handler) + }; + const fixtureAdapter = createPlaywrightFixtureAdapter({ + records: [wsRecord(1, 'receive', '{"op":"CONNECTED"}'), wsRecord(2, 'send', '{"op":"GET_NOTE"}')], + version: fixtureVersion + }); + await fixtureAdapter.install(page); + + const replies = []; + const handlers = []; + calls[0]({ + onMessage: handler => handlers.push(handler), + send: message => replies.push(message) + }); + + assert.deepEqual(replies, ['{"op":"CONNECTED"}']); + assert.throws(() => fixtureAdapter.assertComplete(), /1 unconsumed record/); + handlers[0]('{"op":"GET_NOTE"}'); + assert.doesNotThrow(() => fixtureAdapter.assertComplete()); +}); + +test('Playwright adapter requires every REST response and WebSocket connection to be consumed', async () => { + const restAdapter = createPlaywrightFixtureAdapter({ + records: [ + { + kind: 'rest', + sequence: 1, + rest: { + bodyJson: { id: 'note-a' }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url: '/api/notebook/note-a' }, + status: 200 + } + } + ], + version: fixtureVersion + }); + await restAdapter.install({ route: async () => undefined, routeWebSocket: async () => undefined }); + assert.throws(() => restAdapter.assertComplete(), /1 unconsumed record/); + + const webSocketAdapter = createPlaywrightFixtureAdapter({ + records: [wsRecord(1, 'send', '{"op":"GET_NOTE"}')], + version: fixtureVersion + }); + await webSocketAdapter.install({ route: async () => undefined, routeWebSocket: async () => undefined }); + assert.throws(() => webSocketAdapter.assertComplete(), /was never connected/); +}); + +test('replay rejects REST and WebSocket traffic that violates the captured transport order', async () => { + const adapter = fixtureModule.createFixtureReplayAdapter({ + records: [ + wsRecord(1, 'send', '{"op":"GET_NOTE"}'), + { + kind: 'rest', + sequence: 2, + rest: { + bodyJson: { id: 'note-a' }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url: '/api/notebook/note-a' }, + status: 200 + } + } + ], + version: fixtureVersion + }); + + await assert.rejects( + () => + adapter.route({ fulfill: async () => undefined }, request('GET', 'http://127.0.0.1:8080/api/notebook/note-a')), + /Transport fixture out of order: expected WebSocket send/ + ); +}); + +test('Playwright adapter preserves a REST request, server WebSocket frame, and REST response ordering', async () => { + const calls = []; + const page = { + route: async (_pattern, handler) => calls.push({ handler, kind: 'route' }), + routeWebSocket: async (_pattern, handler) => calls.push({ handler, kind: 'websocket' }) + }; + const fixtureAdapter = createPlaywrightFixtureAdapter({ + records: [ + { + kind: 'rest', + sequence: 1, + rest: { + direction: 'request', + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url: '/api/notebook/note-a' } + } + }, + wsRecord(2, 'receive', '{"op":"NOTE","noteId":"note-a"}'), + { + kind: 'rest', + sequence: 3, + rest: { + bodyJson: { id: 'note-a' }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url: '/api/notebook/note-a' }, + status: 200 + } + } + ], + version: fixtureVersion + }); + await fixtureAdapter.install(page); + + const events = []; + calls + .find(call => call.kind === 'websocket') + .handler({ + onMessage: () => undefined, + send: message => events.push(`websocket:${message}`) + }); + await calls + .find(call => call.kind === 'route') + .handler( + { fulfill: async value => events.push(`rest:${value.body}`) }, + request('GET', 'http://127.0.0.1:8080/api/notebook/note-a') + ); + + assert.deepEqual(events, ['websocket:{"op":"NOTE","noteId":"note-a"}', 'rest:{"id":"note-a"}']); + assert.doesNotThrow(() => fixtureAdapter.assertComplete()); +}); + +test('Playwright adapter rejects REST requests whose body does not match the captured request record', async () => { + const calls = []; + const page = { + route: async (_pattern, handler) => calls.push({ handler, kind: 'route' }), + routeWebSocket: async () => undefined + }; + const fixtureAdapter = createPlaywrightFixtureAdapter({ + records: [ + { + kind: 'rest', + sequence: 1, + rest: { + direction: 'request', + request: { + bodyJson: { paragraphId: 'paragraph-a' }, + headers: { accept: 'application/json', 'content-type': 'application/json' }, + method: 'POST', + url: '/api/notebook/note-a/paragraph' + } + } + }, + { + kind: 'rest', + sequence: 2, + rest: { + bodyJson: { status: 'ok' }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { + bodyJson: { paragraphId: 'paragraph-a' }, + headers: { accept: 'application/json', 'content-type': 'application/json' }, + method: 'POST', + url: '/api/notebook/note-a/paragraph' + }, + status: 200 + } + } + ], + version: fixtureVersion + }); + await fixtureAdapter.install(page); + + await assert.rejects( + () => + calls + .find(call => call.kind === 'route') + .handler( + { fulfill: async () => undefined }, + request( + 'POST', + 'http://127.0.0.1:8080/api/notebook/note-a/paragraph', + '{"paragraphId":"paragraph-b"}', + { accept: 'application/json', 'content-type': 'application/json' } + ) + ), + /REST fixture request mismatch/ + ); +}); + +test('Playwright adapter waits for an interleaved client WebSocket frame before fulfilling REST', async () => { + const calls = []; + const page = { + route: async (_pattern, handler) => calls.push({ handler, kind: 'route' }), + routeWebSocket: async (_pattern, handler) => calls.push({ handler, kind: 'websocket' }) + }; + const fixtureAdapter = createPlaywrightFixtureAdapter({ + records: [ + { + kind: 'rest', + sequence: 1, + rest: { + direction: 'request', + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url: '/api/notebook/note-a' } + } + }, + wsRecord(2, 'send', '{"op":"GET_NOTE"}'), + { + kind: 'rest', + sequence: 3, + rest: { + bodyJson: { id: 'note-a' }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url: '/api/notebook/note-a' }, + status: 200 + } + } + ], + version: fixtureVersion + }); + await fixtureAdapter.install(page); + + const handlers = []; + calls + .find(call => call.kind === 'websocket') + .handler({ + onMessage: handler => handlers.push(handler), + send: () => undefined + }); + const routePromise = calls + .find(call => call.kind === 'route') + .handler({ fulfill: async () => undefined }, request('GET', 'http://127.0.0.1:8080/api/notebook/note-a')); + handlers[0]('{"op":"GET_NOTE"}'); + await routePromise; + + assert.doesNotThrow(() => fixtureAdapter.assertComplete()); +}); + +test('Playwright adapter rejects a REST request that arrives before an expected client WebSocket frame', async () => { + const calls = []; + const page = { + route: async (_pattern, handler) => calls.push({ handler, kind: 'route' }), + routeWebSocket: async (_pattern, handler) => calls.push({ handler, kind: 'websocket' }) + }; + const fixtureAdapter = createPlaywrightFixtureAdapter({ + records: [wsRecord(1, 'send', '{"op":"GET_NOTE"}')], + version: fixtureVersion + }); + await fixtureAdapter.install(page); + + await assert.rejects( + calls + .find(call => call.kind === 'route') + .handler({ fulfill: async () => undefined }, request('GET', 'http://127.0.0.1:8080/api/notebook/note-a')), + /expected WebSocket send, got REST GET \/api\/notebook\/note-a/ + ); +}); + +test('fixture validation reports non-object records instead of throwing', () => { + assert.deepEqual(validateFixture({ records: [null], version: fixtureVersion }), ['records[0] must be an object']); +}); + +test('Playwright adapter forwards WebSocket messages only with explicit passthrough', async () => { + const calls = []; + const page = { + route: async () => undefined, + routeWebSocket: async (_pattern, handler) => calls.push(handler) + }; + await createPlaywrightFixtureAdapter( + { records: [wsRecord(1, 'send', '{"op":"GET_NOTE"}')], version: fixtureVersion }, + { passthrough: true } + ).install(page); + + const forwarded = []; + const handlers = []; + calls[0]({ + connectToServer: () => ({ send: message => forwarded.push(message) }), + onMessage: handler => handlers.push(handler), + send: () => undefined + }); + handlers[0]('{"op":"GET_NOTE"}'); + + assert.deepEqual(forwarded, ['{"op":"GET_NOTE"}']); +}); + +test('transport fixtures no longer expose a custom WebSocket frame parser', () => { + assert.equal('createWebSocketFrameParser' in fixtureModule, false); + assert.equal(existsSync(path.resolve('e2e/core-contract/capture-proxy.mjs')), false); +}); + +function createRoot() { + return { + root: mkdtempSync(path.join(os.tmpdir(), 'zeppelin-capture-')), + zeppelinPort: freePortSync() + }; +} + +function wsRecord(sequence, direction, payloadText) { + return { + kind: 'websocket', + sequence, + websocket: { direction, payloadText } + }; +} + +function start(root, options = {}) { + const result = run( + ['start', '--root', root.root, '--mode', options.mode ?? 'anonymous', '--port', String(root.zeppelinPort)], + { CAPTURE_ZEPPELIN_COMMAND: `node ${stub} ${root.root}` } + ); + assert.equal(result.status, 0, result.stderr); +} + +function stop(root) { + const result = run(['stop', '--root', root.root, '--port', String(root.zeppelinPort)]); + assert.equal(result.status, 0, result.stderr); +} + +function run(args, env = {}) { + return spawnSync('bash', [script, ...args], { + cwd: path.resolve('.'), + encoding: 'utf8', + env: { ...process.env, ...env } + }); +} + +function request(method, url, body = '', headers = { accept: 'application/json' }) { + return { + headers: () => headers, + method: () => method, + postData: () => body, + url: () => url + }; +} + +function response(sourceRequest, status, body, headers = { 'content-type': 'application/json' }) { + return { + headers: () => headers, + request: () => sourceRequest, + status: () => status, + text: async () => (typeof body === 'function' ? body() : body) + }; +} + +function listen() { + return new Promise(resolve => { + const server = http.createServer(); + server.listen(0, '127.0.0.1', () => resolve(server)); + }); +} + +function close(server) { + return new Promise(resolve => server.close(resolve)); +} + +function freePortSync() { + const result = spawnSync( + process.execPath, + [ + '-e', + "require('net').createServer().listen(0, '127.0.0.1', function () { console.log(this.address().port); this.close(); })" + ], + { + encoding: 'utf8' + } + ); + return Number(result.stdout.trim()); +} diff --git a/zeppelin-web-angular/e2e/core-contract/capture-stub-zeppelin.mjs b/zeppelin-web-angular/e2e/core-contract/capture-stub-zeppelin.mjs new file mode 100755 index 0000000000..2087e7483f --- /dev/null +++ b/zeppelin-web-angular/e2e/core-contract/capture-stub-zeppelin.mjs @@ -0,0 +1,60 @@ +#!/usr/bin/env node +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 crypto from 'node:crypto'; +import http from 'node:http'; + +const port = Number(process.env.ZEPPELIN_PORT); +const root = process.env.ZEPPELIN_CAPTURE_ROOT; + +if (!port || !root) { + process.stderr.write('ZEPPELIN_PORT and ZEPPELIN_CAPTURE_ROOT are required\n'); + process.exit(2); +} + +const server = http.createServer((request, response) => { + if (request.url === '/api/version') { + response.writeHead(200, { 'content-type': 'application/json' }); + response.end('{"version":"stub"}'); + return; + } + response.writeHead(200, { 'content-type': 'application/json' }); + response.end('{"id":"note-a","paragraphs":[{"id":"paragraph-a"}]}'); +}); + +server.on('upgrade', (request, socket) => { + const key = request.headers['sec-websocket-key']; + const accept = crypto.createHash('sha1').update(`${key}258EAFA5-E914-47DA-95CA-C5AB0DC85B11`).digest('base64'); + socket.write( + [ + 'HTTP/1.1 101 Switching Protocols', + 'Upgrade: websocket', + 'Connection: Upgrade', + `Sec-WebSocket-Accept: ${accept}`, + '', + '' + ].join('\r\n') + ); + socket.on('data', () => { + const payload = Buffer.from('{"op":"NOTE"}'); + socket.write(Buffer.concat([Buffer.from([0x81, payload.length]), payload])); + }); +}); + +server.listen(port, '127.0.0.1'); +process.on('SIGTERM', () => server.close(() => process.exit(0))); diff --git a/zeppelin-web-angular/e2e/core-contract/notebook-transport-fixture.mjs b/zeppelin-web-angular/e2e/core-contract/notebook-transport-fixture.mjs new file mode 100644 index 0000000000..783983570c --- /dev/null +++ b/zeppelin-web-angular/e2e/core-contract/notebook-transport-fixture.mjs @@ -0,0 +1,699 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 { mkdirSync, writeFileSync } from 'node:fs'; +import path from 'node:path'; + +export const fixtureVersion = 1; + +export const restDirections = new Set(['request', 'response']); +export const websocketDirections = new Set(['send', 'receive']); + +const safeHeaderNames = new Set(['accept', 'content-type']); +const volatileFieldNames = new Set(['dateCreated', 'dateFinished', 'dateStarted', 'lastUpdated', 'msgId', 'time']); +const sensitiveFieldNames = new Set([ + 'authorization', + 'api-key', + 'apikey', + 'client-secret', + 'clientsecret', + 'cookie', + 'credential', + 'credentials', + 'password', + 'principal', + 'secret', + 'set-cookie', + 'ticket', + 'token' +]); + +export function normalizeFixtureRecord(value) { + if (Array.isArray(value)) { + return value.map(item => normalizeFixtureRecord(item)); + } + if (!value || typeof value !== 'object') { + return value; + } + + return Object.fromEntries( + Object.entries(value).map(([key, entry]) => [ + key, + shouldRedactField(key) ? `<${key}>` : volatileFieldNames.has(key) ? `<${key}>` : normalizeFixtureRecord(entry) + ]) + ); +} + +export function normalizeFixture(fixture) { + return { + ...fixture, + records: fixture.records.map(record => normalizeFixtureRecord(record)) + }; +} + +export function sanitizeFixture(fixture) { + return normalizeFixture({ + ...fixture, + records: fixture.records.map(record => sanitizeRecord(record)) + }); +} + +export function validateFixture(fixture) { + const errors = []; + if (!fixture || typeof fixture !== 'object') { + return ['fixture must be an object']; + } + if (fixture.version !== fixtureVersion) { + errors.push(`Unsupported fixture version ${fixture.version}`); + } + if (!Array.isArray(fixture.records) || fixture.records.length === 0) { + errors.push('Fixture records must be a non-empty array'); + return errors; + } + + let previousSequence = 0; + for (const [index, record] of fixture.records.entries()) { + const prefix = `records[${index}]`; + if (!record || typeof record !== 'object' || Array.isArray(record)) { + errors.push(`${prefix} must be an object`); + continue; + } + if (!Number.isInteger(record.sequence) || record.sequence <= previousSequence) { + errors.push(`${prefix}.sequence must increase without reordering`); + } + previousSequence = record.sequence; + + if (record.kind === 'rest') { + validateRestRecord(errors, prefix, record); + } else if (record.kind === 'websocket') { + validateWebSocketRecord(errors, prefix, record); + } else { + errors.push(`${prefix}.kind must be rest or websocket`); + } + } + return errors; +} + +export function visitNormalizedFixtureRecords(fixture, visitor = () => undefined) { + const adapter = createFixtureReplayAdapter(fixture); + for (const record of adapter.records()) { + visitor(record); + } +} + +export const replayFixture = visitNormalizedFixtureRecords; + +export function createFixtureReplayAdapter(fixture) { + const errors = validateFixture(fixture); + if (errors.length > 0) { + throw new Error(errors.join('\n')); + } + + const records = sanitizeFixture(fixture).records; + const replayRecords = records.filter( + record => record.kind === 'websocket' || (record.kind === 'rest' && record.rest.direction === 'response') + ); + let replayCursor = 0; + + const nextRecord = () => replayRecords[replayCursor]; + const consumeServerMessages = () => { + const messages = []; + while (nextRecord()?.kind === 'websocket' && nextRecord().websocket.direction === 'receive') { + messages.push(deserializeWebSocketPayload(nextRecord().websocket)); + replayCursor += 1; + } + return messages; + }; + + return { + records: () => records[Symbol.iterator](), + route: async (route, request) => { + if (!isNotebookRestUrl(request.url())) { + await route.continue?.(); + return; + } + const method = request.method(); + const requestKey = `${method} ${urlPath(request.url())}`; + const response = nextRecord(); + if (!response || response.kind !== 'rest') { + throw new Error( + `Transport fixture out of order: expected ${describeReplayRecord(response)}, got REST ${requestKey}` + ); + } + const responseKey = `${response.rest.request.method} ${response.rest.request.url}`; + if (responseKey !== requestKey) { + throw new Error(`REST fixture request out of order: expected ${responseKey}, got ${requestKey}`); + } + assertRestRequestMatches(response.rest.request, summarizeRequest(request), requestKey); + replayCursor += 1; + await route.fulfill({ + body: serializeRestBody(response.rest), + contentType: response.rest.headers['content-type'] ?? 'application/json', + headers: response.rest.headers, + status: response.rest.status + }); + }, + receiveWebSocketMessages: consumeServerMessages, + sendWebSocketMessage: message => { + const expected = nextRecord(); + if (!expected) { + throw new Error( + `WebSocket fixture messages exhausted before client send: ${stringifyWebSocketMessage(message)}` + ); + } + if (expected.kind !== 'websocket' || expected.websocket.direction !== 'send') { + throw new Error( + `Transport fixture out of order: expected ${describeReplayRecord(expected)}, got WebSocket send` + ); + } + const expectedPayload = deserializeWebSocketPayload(expected.websocket); + if (!webSocketPayloadMatches(expectedPayload, message)) { + throw new Error( + `WebSocket fixture send mismatch: expected ${expectedPayload}, got ${stringifyWebSocketMessage(message)}` + ); + } + replayCursor += 1; + }, + assertComplete: () => { + if (replayCursor !== replayRecords.length) { + throw new Error(`Transport fixture has ${replayRecords.length - replayCursor} unconsumed record(s)`); + } + }, + hasWebSocketRecords: () => replayRecords.some(record => record.kind === 'websocket') + }; +} + +export function createPlaywrightFixtureAdapter(fixture, options = {}) { + const errors = validateFixture(fixture); + if (errors.length > 0) { + throw new Error(errors.join('\n')); + } + + const records = sanitizeFixture(fixture).records; + const pendingRestRequests = []; + let cursor = 0; + let webSocket; + + const nextRecord = () => records[cursor]; + const drain = () => { + while (true) { + const record = nextRecord(); + if (!record) { + return; + } + if (record.kind === 'websocket') { + if (record.websocket.direction === 'receive' && webSocket) { + webSocket.send(deserializeWebSocketPayload(record.websocket)); + cursor += 1; + continue; + } + return; + } + + if (record.rest.direction === 'request') { + const pending = pendingRestRequests.find(entry => !entry.requestMatched); + if (!pending) { + return; + } + const requestKey = `${pending.request.method()} ${urlPath(pending.request.url())}`; + assertRestRequestMatches(record.rest.request, summarizeRequest(pending.request), requestKey); + pending.requestMatched = true; + cursor += 1; + continue; + } + + const pending = pendingRestRequests.find(entry => entry.requestMatched); + if (!pending) { + return; + } + const responseKey = `${record.rest.request.method} ${record.rest.request.url}`; + const requestKey = `${pending.request.method()} ${urlPath(pending.request.url())}`; + if (responseKey !== requestKey) { + throw new Error(`REST fixture request out of order: expected ${responseKey}, got ${requestKey}`); + } + pendingRestRequests.splice(pendingRestRequests.indexOf(pending), 1); + cursor += 1; + void pending.route + .fulfill({ + body: serializeRestBody(record.rest), + contentType: record.rest.headers['content-type'] ?? 'application/json', + headers: record.rest.headers, + status: record.rest.status + }) + .then(pending.resolve, pending.reject); + } + }; + + return { + install: async page => { + await page.route('**/api/**', async (route, request) => { + if (!isNotebookRestUrl(request.url())) { + await route.continue?.(); + return; + } + let resolve; + let reject; + const completed = new Promise((resolvePromise, rejectPromise) => { + resolve = resolvePromise; + reject = rejectPromise; + }); + const pending = { reject, request, requestMatched: false, resolve, route }; + pendingRestRequests.push(pending); + try { + drain(); + const next = nextRecord(); + if (!pending.requestMatched && next?.kind === 'websocket' && next.websocket.direction === 'send') { + pendingRestRequests.splice(pendingRestRequests.indexOf(pending), 1); + throw new Error( + `Transport fixture out of order: expected WebSocket send, got REST ${pending.request.method()} ${urlPath( + pending.request.url() + )}` + ); + } + await completed; + } catch (error) { + const pendingIndex = pendingRestRequests.indexOf(pending); + if (pendingIndex !== -1) { + pendingRestRequests.splice(pendingIndex, 1); + } + throw error; + } + }); + await page.routeWebSocket(/\/ws(?:$|\?)/, ws => { + if (webSocket) { + throw new Error('Transport fixture supports one WebSocket connection per fixture'); + } + webSocket = ws; + const server = options.passthrough && typeof ws.connectToServer === 'function' ? ws.connectToServer() : null; + drain(); + ws.onMessage(message => { + const record = nextRecord(); + if (!record) { + throw new Error( + `WebSocket fixture messages exhausted before client send: ${stringifyWebSocketMessage(message)}` + ); + } + if (record.kind !== 'websocket' || record.websocket.direction !== 'send') { + throw new Error( + `Transport fixture out of order: expected ${describeReplayRecord(record)}, got WebSocket send` + ); + } + const expectedPayload = deserializeWebSocketPayload(record.websocket); + if (!webSocketPayloadMatches(expectedPayload, message)) { + throw new Error( + `WebSocket fixture send mismatch: expected ${expectedPayload}, got ${stringifyWebSocketMessage(message)}` + ); + } + cursor += 1; + if (server) { + server.send(message); + } + drain(); + }); + }); + }, + assertComplete: () => { + if (records.some(record => record.kind === 'websocket') && !webSocket) { + throw new Error('WebSocket fixture was never connected'); + } + if (pendingRestRequests.length > 0 || cursor !== records.length) { + throw new Error(`Transport fixture has ${records.length - cursor} unconsumed record(s)`); + } + } + }; +} + +export function createNotebookTransportRecorder() { + const records = []; + const pending = new Set(); + let sequence = 0; + + const record = value => { + const entry = { + ...value, + sequence: ++sequence + }; + records.push(entry); + return entry; + }; + + return { + install: page => { + page.on('request', request => { + if (!isNotebookRestUrl(request.url())) { + return; + } + record({ + kind: 'rest', + rest: { + direction: 'request', + request: summarizeRequest(request) + } + }); + }); + page.on('response', async response => { + const request = response.request(); + if (!isNotebookRestUrl(request.url())) { + return; + } + const entry = record({ + kind: 'rest', + rest: { + direction: 'response', + headers: filterHeaders(response.headers()), + request: summarizeRequest(request), + status: response.status(), + bodyRaw: '' + } + }); + const bodyRead = safeResponseText(response) + .then(body => { + Object.assign(entry.rest, parseRestBody(body, response.headers())); + }) + .finally(() => pending.delete(bodyRead)); + pending.add(bodyRead); + }); + page.on('websocket', socket => { + if (!isNotebookWebSocketUrl(socket.url())) { + return; + } + socket.on('framesent', frame => record(webSocketRecord('send', framePayload(frame)))); + socket.on('framereceived', frame => record(webSocketRecord('receive', framePayload(frame)))); + }); + }, + stop: async () => { + await Promise.all([...pending]); + }, + snapshot: () => sanitizeFixture({ records: [...records], version: fixtureVersion }), + write: async fixturePath => { + await Promise.all([...pending]); + const sanitized = sanitizeFixture({ records: [...records], version: fixtureVersion }); + mkdirSync(path.dirname(fixturePath), { recursive: true }); + writeFileSync(fixturePath, `${JSON.stringify(sanitized, null, 2)}\n`); + return sanitized; + } + }; +} + +export function serializeRestBody(rest) { + if ('bodyJson' in rest) { + return JSON.stringify(rest.bodyJson); + } + return rest.bodyRaw ?? rest.body ?? ''; +} + +export function parseRestBody(body, headers = {}) { + const contentType = headers['content-type'] ?? headers['Content-Type'] ?? ''; + const trimmed = body.trim(); + if (!trimmed) { + return { bodyRaw: '' }; + } + if (contentType.includes('application/json') || /^[{[]/.test(trimmed)) { + try { + return { bodyJson: JSON.parse(body) }; + } catch { + return { bodyRaw: redactRawSensitiveValues(body) }; + } + } + return { bodyRaw: redactRawSensitiveValues(body) }; +} + +export function stringifyWebSocketMessage(message) { + return Buffer.isBuffer(message) ? message.toString('utf8') : String(message); +} + +export function webSocketPayloadMatches(expectedPayload, actualMessage) { + const expectedBinary = toBinaryBuffer(expectedPayload); + const actualBinary = toBinaryBuffer(actualMessage); + if (expectedBinary || actualBinary) { + return Boolean(expectedBinary && actualBinary && expectedBinary.equals(actualBinary)); + } + + const actualPayload = stringifyWebSocketMessage(actualMessage); + if (looksLikeJson(expectedPayload) && looksLikeJson(actualPayload)) { + try { + return ( + stableJson(normalizeFixtureRecord(JSON.parse(actualPayload))) === + stableJson(normalizeFixtureRecord(JSON.parse(expectedPayload))) + ); + } catch { + return false; + } + } + return expectedPayload === actualPayload; +} + +function summarizeRequest(request) { + const body = request.postData() ?? ''; + return { + headers: filterHeaders(request.headers()), + method: request.method(), + url: urlPath(request.url()), + ...parseRestBody(body, request.headers()) + }; +} + +function sanitizeRecord(record) { + if (record.kind === 'websocket') { + return { + ...record, + websocket: sanitizeWebSocket(record.websocket) + }; + } + if (record.kind !== 'rest') { + return record; + } + return { + ...record, + rest: { + ...record.rest, + ...(record.rest.headers ? { headers: filterHeaders(record.rest.headers) } : {}), + ...(record.rest.request + ? { + request: { + ...record.rest.request, + headers: filterHeaders(record.rest.request.headers) + } + } + : {}) + } + }; +} + +function filterHeaders(headers = {}) { + return Object.fromEntries( + Object.entries(headers) + .map(([key, value]) => [key.toLowerCase(), Array.isArray(value) ? value.join(', ') : String(value ?? '')]) + .filter(([key]) => safeHeaderNames.has(key)) + ); +} + +function webSocketRecord(direction, payload) { + return { + kind: 'websocket', + websocket: { + direction, + ...(Buffer.isBuffer(payload) ? { payloadBase64: payload.toString('base64') } : { payloadText: String(payload) }) + } + }; +} + +export function isNotebookRestUrl(value) { + const url = new URL(value); + return url.pathname === '/api/notebook' || url.pathname.startsWith('/api/notebook/'); +} + +function isNotebookWebSocketUrl(value) { + const url = new URL(value); + return url.pathname === '/ws'; +} + +async function safeResponseText(response) { + return response.text(); +} + +function urlPath(value) { + const url = new URL(value); + for (const [key] of url.searchParams) { + if (shouldRedactField(key)) { + url.searchParams.set(key, `<${key}>`); + } else if (volatileFieldNames.has(key)) { + url.searchParams.set(key, `<${key}>`); + } + } + return `${url.pathname}${url.search}`; +} + +function stableJson(value) { + if (Array.isArray(value)) { + return `[${value.map(entry => stableJson(entry)).join(',')}]`; + } + if (value && typeof value === 'object') { + return `{${Object.keys(value) + .sort() + .map(key => `${JSON.stringify(key)}:${stableJson(value[key])}`) + .join(',')}}`; + } + return JSON.stringify(value); +} + +function sanitizeWebSocket(websocket) { + if (!websocket?.payloadText || !looksLikeJson(websocket.payloadText)) { + return websocket; + } + try { + return { + ...websocket, + payloadText: JSON.stringify(normalizeFixtureRecord(JSON.parse(websocket.payloadText))) + }; + } catch { + return { ...websocket, payloadText: redactRawSensitiveValues(websocket.payloadText) }; + } +} + +function deserializeWebSocketPayload(websocket) { + if ('payloadBase64' in websocket) { + return Buffer.from(websocket.payloadBase64, 'base64'); + } + return websocket.payloadText ?? ''; +} + +function framePayload(frame) { + if (frame && typeof frame === 'object' && 'payload' in frame) { + return frame.payload; + } + return frame; +} + +function assertRestRequestMatches(expectedRequest, actualRequest, requestKey) { + const expected = normalizeFixtureRecord(sanitizeRestRequest(expectedRequest)); + const actual = normalizeFixtureRecord(sanitizeRestRequest(actualRequest)); + if (stableJson(expected) !== stableJson(actual)) { + throw new Error( + `REST fixture request mismatch for ${requestKey}: expected ${stableJson(expected)}, got ${stableJson(actual)}` + ); + } +} + +function sanitizeRestRequest(request) { + return { + ...request, + headers: filterHeaders(request.headers) + }; +} + +function shouldRedactField(key) { + const normalized = key.toLowerCase(); + return ( + sensitiveFieldNames.has(normalized) || + normalized.includes('apikey') || + normalized.includes('credential') || + normalized.includes('password') || + normalized.includes('secret') || + normalized.includes('token') + ); +} + +function redactRawSensitiveValues(value) { + const sensitiveFieldPattern = '(?:api[-_]?key|client[-_]?secret|credential(?:s)?|password|secret|ticket|token)'; + return String(value) + .replace(new RegExp(`([?&]${sensitiveFieldPattern}=)([^&#\\s]+)`, 'gi'), (_match, prefix) => `${prefix}<redacted>`) + .replace( + new RegExp(`((?:[\\"']?${sensitiveFieldPattern}[\\"']?\\s*[:=]\\s*))(?:[\\"']?)([^,}\\]\\s\\"']*)`, 'gi'), + (_match, prefix) => `${prefix}<redacted>` + ); +} + +function describeReplayRecord(record) { + if (!record) { + return 'end of fixture'; + } + if (record.kind === 'rest') { + return `REST ${record.rest.request.method} ${record.rest.request.url}`; + } + return `WebSocket ${record.websocket.direction}`; +} + +function toBinaryBuffer(value) { + if (Buffer.isBuffer(value)) { + return value; + } + if (value instanceof ArrayBuffer) { + return Buffer.from(value); + } + if (ArrayBuffer.isView(value)) { + return Buffer.from(value.buffer, value.byteOffset, value.byteLength); + } + return null; +} + +const validateRestRecord = (errors, prefix, record) => { + if (!record.rest || typeof record.rest !== 'object') { + errors.push(`${prefix}.rest is required`); + return; + } + if (!restDirections.has(record.rest.direction)) { + errors.push(`${prefix}.rest.direction must be request or response`); + } + if (!isHttpMethod(record.rest.request?.method)) { + errors.push(`${prefix}.rest.request.method is required`); + } + if (typeof record.rest.request?.url !== 'string') { + errors.push(`${prefix}.rest.request.url is required`); + } + if (!isHeaderRecord(record.rest.request?.headers)) { + errors.push(`${prefix}.rest.request.headers must be an object`); + } + if (record.rest.direction === 'response') { + if (!Number.isInteger(record.rest.status)) { + errors.push(`${prefix}.rest.status is required for responses`); + } + if (!isHeaderRecord(record.rest.headers)) { + errors.push(`${prefix}.rest.headers must be an object`); + } + if (!hasRestBody(record.rest)) { + errors.push(`${prefix}.rest.bodyJson or bodyRaw is required to preserve response shape`); + } + } else if (!hasRestBody(record.rest.request)) { + errors.push(`${prefix}.rest.request.bodyJson or bodyRaw is required to preserve request shape`); + } +}; + +const validateWebSocketRecord = (errors, prefix, record) => { + if (!record.websocket || typeof record.websocket !== 'object') { + errors.push(`${prefix}.websocket is required`); + return; + } + if (!websocketDirections.has(record.websocket.direction)) { + errors.push(`${prefix}.websocket.direction must be send or receive`); + } + if (!('payloadText' in record.websocket) && !('payloadBase64' in record.websocket)) { + errors.push(`${prefix}.websocket payloadText or payloadBase64 is required to preserve message shape`); + } +}; + +const isHttpMethod = value => typeof value === 'string' && /^[A-Z]+$/.test(value); + +const isHeaderRecord = value => + Boolean(value) && + typeof value === 'object' && + !Array.isArray(value) && + Object.values(value).every(entry => typeof entry === 'string'); + +const hasRestBody = value => Boolean(value) && ('bodyJson' in value || 'bodyRaw' in value || 'body' in value); + +const looksLikeJson = value => /^[{[]/.test(String(value).trim()); diff --git a/zeppelin-web-angular/e2e/tests/notebook/core-contract/capture-fixtures.spec.ts b/zeppelin-web-angular/e2e/tests/notebook/core-contract/capture-fixtures.spec.ts new file mode 100644 index 0000000000..352680b804 --- /dev/null +++ b/zeppelin-web-angular/e2e/tests/notebook/core-contract/capture-fixtures.spec.ts @@ -0,0 +1,503 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You 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 { expect, test } from '@playwright/test'; + +import { + createFixtureReplayAdapter, + createNotebookTransportRecorder, + createPlaywrightFixtureAdapter, + fixtureVersion, + isNotebookRestUrl, + normalizeFixtureRecord, + visitNormalizedFixtureRecords, + validateFixture, + webSocketPayloadMatches +} from '../../../core-contract/notebook-transport-fixture.mjs'; +import { + addPageAnnotationBeforeEach, + createTestNotebook, + navigateToNotebookWithFallback, + PAGES, + performLoginIfRequired, + waitForZeppelinReady +} from '../../../utils'; + +type TestHandler = (value: unknown) => unknown; +type AdapterCall = [kind: string, pattern: unknown, handler: TestHandler]; + +test.describe('Notebook core transport fixture replay', () => { + addPageAnnotationBeforeEach(PAGES.WORKSPACE.NOTEBOOK); + + test('accepts discriminated versioned REST and WebSocket records without changing order or shape', () => { + const fixture = sampleFixture(); + const replayed: unknown[] = []; + + visitNormalizedFixtureRecords(fixture, record => replayed.push(record)); + + expect(replayed).toEqual([ + { + kind: 'websocket', + sequence: 1, + websocket: { direction: 'send', payloadText: '{"op":"GET_NOTE","msgId":"<msgId>"}' } + }, + { + kind: 'rest', + rest: { + bodyJson: { id: 'note-a', paragraphs: [{ id: 'paragraph-a' }] }, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url: '/api/notebook/note-a' }, + status: 200 + }, + sequence: 2 + } + ]); + }); + + test('rejects unsupported versions, ordering loss, and missing discriminated fields', () => { + expect(validateFixture({ version: 999, records: [validRestRecord()] })).toContain( + 'Unsupported fixture version 999' + ); + expect( + validateFixture({ + version: fixtureVersion, + records: [validRestRecord({ sequence: 2 }), validRestRecord({ sequence: 1 })] + }) + ).toContain('records[1].sequence must increase without reordering'); + + const errors = validateFixture({ + version: fixtureVersion, + records: [{ sequence: 1, kind: 'websocket', websocket: { direction: 'receive' } }] + }).join('\n'); + + expect(errors).toContain('websocket payloadText or payloadBase64 is required'); + }); + + test('redacts sensitive fields without replacing stable ids or operation fields', () => { + expect( + normalizeFixtureRecord({ + id: 'stable-id', + operation: 'LIST_REVISION_HISTORY', + ticket: 'secret', + principal: 'anonymous', + data: { noteId: '2A94M5J1Z', revisionId: 'rev-1', stable: 'kept' } + }) + ).toEqual({ + id: 'stable-id', + operation: 'LIST_REVISION_HISTORY', + ticket: '<ticket>', + principal: '<principal>', + data: { noteId: '2A94M5J1Z', revisionId: 'rev-1', stable: 'kept' } + }); + }); + + test('provides a Playwright route adapter for REST responses', async () => { + const adapter = createFixtureReplayAdapter({ + records: [ + validRestRecord({ + sequence: 1, + rest: restResponse('/api/notebook/note-a', { id: 'note-a', paragraphs: [{ id: 'paragraph-a' }] }) + }) + ], + version: fixtureVersion + }); + const fulfilled: unknown[] = []; + await adapter.route( + { + fulfill: async value => fulfilled.push(value) + }, + { + headers: () => ({ accept: 'application/json' }), + method: () => 'GET', + postData: () => null, + url: () => 'http://localhost:8080/api/notebook/note-a' + } + ); + + expect(fulfilled).toEqual([ + { + body: '{"id":"note-a","paragraphs":[{"id":"paragraph-a"}]}', + contentType: 'application/json', + headers: { 'content-type': 'application/json' }, + status: 200 + } + ]); + expect(() => adapter.assertComplete()).not.toThrow(); + }); + + test('captures browser-level REST and WebSocket events with pre-write redaction', async () => { + const handlers = new Map<string, TestHandler[]>(); + const page = { + on: (eventName: string, handler: TestHandler) => + handlers.set(eventName, [...(handlers.get(eventName) ?? []), handler]) + }; + const socketHandlers = new Map<string, TestHandler[]>(); + const socket = { + on: (eventName: string, handler: TestHandler) => + socketHandlers.set(eventName, [...(socketHandlers.get(eventName) ?? []), handler]), + url: () => 'http://localhost:8080/ws' + }; + const recorder = createNotebookTransportRecorder(); + + recorder.install(page); + handlers.get('request')?.[0]?.( + request('POST', 'http://localhost:8080/api/notebook', '{"ticket":"secret","id":"stable-id"}') + ); + handlers.get('response')?.[0]?.( + response(request('GET', 'http://localhost:8080/api/notebook/stable-id'), 200, '{"id":"stable-id"}') + ); + handlers.get('websocket')?.[0]?.(socket); + socketHandlers.get('framesent')?.[0]?.({ payload: '{"op":"GET_NOTE","msgId":"runtime"}' }); + socketHandlers.get('framereceived')?.[0]?.({ payload: '{"op":"NOTE","noteId":"stable-id"}' }); + + await recorder.stop(); + const fixture = recorder.snapshot(); + + expect(validateFixture(fixture)).toEqual([]); + expect(fixture.records.map(record => record.rest?.direction ?? record.websocket?.direction)).toEqual([ + 'request', + 'response', + 'send', + 'receive' + ]); + expect(fixture.records[0].rest.request.bodyJson).toEqual({ id: 'stable-id', ticket: '<ticket>' }); + expect(fixture.records[2].websocket.payloadText).toBe('{"op":"GET_NOTE","msgId":"<msgId>"}'); + }); + + test('redacts WebSocket JSON payloads and normalizes REST query values before writing', async ({}, testInfo) => { + const handlers = new Map<string, TestHandler[]>(); + const page = { + on: (eventName: string, handler: TestHandler) => + handlers.set(eventName, [...(handlers.get(eventName) ?? []), handler]) + }; + const socketHandlers = new Map<string, TestHandler[]>(); + const socket = { + on: (eventName: string, handler: TestHandler) => + socketHandlers.set(eventName, [...(socketHandlers.get(eventName) ?? []), handler]), + url: () => 'http://localhost:8080/ws' + }; + const recorder = createNotebookTransportRecorder(); + const fixturePath = testInfo.outputPath('notebook-transport.json'); + + recorder.install(page); + handlers.get('request')?.[0]?.( + request( + 'GET', + 'http://localhost:8080/api/notebook/note-a?ticket=secret-ticket&token=secret-token&msgId=runtime&view=stable' + ) + ); + handlers.get('websocket')?.[0]?.(socket); + socketHandlers.get('framesent')?.[0]?.({ + payload: + '{"op":"GET_NOTE","id":"stable-id","noteId":"note-a","ticket":"secret-ticket","principal":"alice","msgId":"runtime"}' + }); + + const fixture = await recorder.write(fixturePath); + + expect(JSON.stringify(fixture)).not.toContain('secret-ticket'); + expect(JSON.stringify(fixture)).not.toContain('secret-token'); + expect(JSON.stringify(fixture)).not.toContain('alice'); + expect(fixture.records[0].rest.request.url).toBe( + '/api/notebook/note-a?ticket=%3Cticket%3E&token=%3Ctoken%3E&msgId=%3CmsgId%3E&view=stable' + ); + expect(fixture.records[1].websocket.payloadText).toContain('"op":"GET_NOTE"'); + expect(fixture.records[1].websocket.payloadText).toContain('"id":"stable-id"'); + expect(fixture.records[1].websocket.payloadText).toContain('"noteId":"note-a"'); + }); + + test('matches replayed query placeholders and JSON payloads deterministically', async () => { + const fulfilled: unknown[] = []; + const adapter = createFixtureReplayAdapter({ + records: [ + validRestRecord({ + sequence: 1, + rest: restResponse('/api/notebook/note-a?ticket=%3Cticket%3E&msgId=%3CmsgId%3E', { id: 'note-a' }) + }) + ], + version: fixtureVersion + }); + + await adapter.route( + { fulfill: async value => fulfilled.push(value) }, + request('GET', 'http://localhost:8080/api/notebook/note-a?ticket=runtime-secret&msgId=runtime-id') + ); + + expect(fulfilled).toHaveLength(1); + expect(webSocketPayloadMatches('{"op":"GET_NOTE","msgId":"<msgId>"}', '{"msgId":"runtime","op":"GET_NOTE"}')).toBe( + true + ); + }); + + test('waits for pending response body reads before writing', async ({}, testInfo) => { + const handlers = new Map<string, TestHandler[]>(); + const page = { + on: (eventName: string, handler: TestHandler) => + handlers.set(eventName, [...(handlers.get(eventName) ?? []), handler]) + }; + const recorder = createNotebookTransportRecorder(); + let resolveBody: (value: string) => void = () => undefined; + + recorder.install(page); + handlers.get('response')?.[0]?.( + response(request('GET', 'http://localhost:8080/api/notebook/note-a'), 200, () => { + return new Promise<string>(resolve => { + resolveBody = resolve; + }); + }) + ); + const writePromise = recorder.write(testInfo.outputPath('notebook-transport.json')); + resolveBody('{"id":"note-a"}'); + const fixture = await writePromise; + + expect(fixture.records[0].rest.bodyJson).toEqual({ id: 'note-a' }); + }); + + test('limits REST capture and replay to notebook APIs', async () => { + expect(isNotebookRestUrl('http://localhost:8080/api/notebook')).toBe(true); + expect(isNotebookRestUrl('http://localhost:8080/api/notebook/note-a')).toBe(true); + expect(isNotebookRestUrl('http://localhost:8080/api/security/ticket')).toBe(false); + + const adapter = createFixtureReplayAdapter({ + records: [validRestRecord({ sequence: 1, rest: restResponse('/api/notebook/note-a', { id: 'note-a' }) })], + version: fixtureVersion + }); + const continued: unknown[] = []; + + await adapter.route( + { continue: async () => continued.push('continue') }, + request('GET', 'http://localhost/api/security/ticket') + ); + + expect(continued).toEqual(['continue']); + }); + + test('consumes repeated REST responses in captured sequence and fails on exhausted or out-of-order requests', async () => { + const adapter = createFixtureReplayAdapter({ + records: [ + validRestRecord({ sequence: 1, rest: restResponse('/api/notebook/a', { id: 'a' }) }), + validRestRecord({ sequence: 2, rest: restResponse('/api/notebook/a', { id: 'a-2' }) }), + validRestRecord({ sequence: 3, rest: restResponse('/api/notebook/b', { id: 'b' }) }) + ], + version: fixtureVersion + }); + const bodies: unknown[] = []; + const route = { fulfill: async value => bodies.push(value.body) }; + + await adapter.route(route, request('GET', 'http://localhost:8080/api/notebook/a')); + await adapter.route(route, request('GET', 'http://localhost:8080/api/notebook/a')); + await expect(adapter.route(route, request('GET', 'http://localhost:8080/api/notebook/a'))).rejects.toThrow( + 'expected GET /api/notebook/b' + ); + + expect(bodies).toEqual(['{"id":"a"}', '{"id":"a-2"}']); + }); + + test('installs fixture-only REST and WebSocket routes through the Playwright page adapter API', async () => { + const calls: AdapterCall[] = []; + const page = { + route: async (pattern: unknown, handler: TestHandler) => calls.push(['route', pattern, handler]), + routeWebSocket: async (pattern: unknown, handler: TestHandler) => calls.push(['routeWebSocket', pattern, handler]) + }; + + const adapter = createPlaywrightFixtureAdapter(sampleFixtureWithTwoWebSocketRoundTrips()); + await adapter.install(page); + + expect(calls[0][0]).toBe('route'); + expect(calls[1][0]).toBe('routeWebSocket'); + const wsSends: unknown[] = []; + const serverSends: unknown[] = []; + const clientMessageHandlers: ((message: string) => void)[] = []; + calls[1][2]({ + connectToServer: () => ({ send: message => serverSends.push(message) }), + onMessage: handler => clientMessageHandlers.push(handler), + send: message => wsSends.push(message) + }); + clientMessageHandlers[0]('{"op":"GET_NOTE","msgId":"runtime-1"}'); + clientMessageHandlers[0]('{"op":"RUN_PARAGRAPH","paragraphId":"paragraph-b"}'); + + expect(serverSends).toEqual([]); + expect(wsSends).toEqual(['{"op":"NOTE","noteId":"note-a"}', '{"op":"PARAGRAPH","paragraphId":"paragraph-b"}']); + expect(() => adapter.assertComplete()).not.toThrow(); + expect(() => clientMessageHandlers[0]('{"op":"EXTRA"}')).toThrow('WebSocket fixture messages exhausted'); + }); + + test('only forwards WebSocket client messages when passthrough is explicit', async () => { + const calls: AdapterCall[] = []; + const serverSends: unknown[] = []; + const page = { + route: async () => undefined, + routeWebSocket: async (_pattern: unknown, handler: TestHandler) => + calls.push(['routeWebSocket', _pattern, handler]) + }; + + const adapter = createPlaywrightFixtureAdapter(sampleFixtureWithReceive(), { passthrough: true }); + await adapter.install(page); + + const clientMessageHandlers: ((message: string) => void)[] = []; + calls[0][2]({ + connectToServer: () => ({ send: message => serverSends.push(message) }), + onMessage: handler => clientMessageHandlers.push(handler), + send: () => undefined + }); + clientMessageHandlers[0]('{"op":"GET_NOTE","msgId":"runtime-1"}'); + + expect(serverSends).toEqual(['{"op":"GET_NOTE","msgId":"runtime-1"}']); + expect(() => adapter.assertComplete()).not.toThrow(); + }); + + test('fails on WebSocket send payload mismatch instead of replaying a stale receive', async () => { + const calls: AdapterCall[] = []; + const page = { + route: async () => undefined, + routeWebSocket: async (_pattern: unknown, handler: TestHandler) => + calls.push(['routeWebSocket', _pattern, handler]) + }; + + await createPlaywrightFixtureAdapter(sampleFixtureWithReceive()).install(page); + + const clientMessageHandlers: ((message: string) => void)[] = []; + calls[0][2]({ + connectToServer: () => { + throw new Error('should not connect in fixture-only mode'); + }, + onMessage: handler => clientMessageHandlers.push(handler), + send: () => undefined + }); + + expect(() => clientMessageHandlers[0]('{"op":"RUN_PARAGRAPH"}')).toThrow('WebSocket fixture send mismatch'); + }); + + test('records notebook REST and WebSocket traffic from a real Zeppelin page', async ({ page }) => { + await page.goto('/#/'); + await waitForZeppelinReady(page); + await performLoginIfRequired(page); + const { noteId } = await createTestNotebook(page); + const recorderPage = await page.context().newPage(); + const recorder = createNotebookTransportRecorder(); + + recorder.install(recorderPage); + await navigateToNotebookWithFallback(recorderPage, noteId); + await recorderPage.evaluate(async id => { + await fetch(`/api/notebook/${id}`, { headers: { accept: 'application/json' } }); + }, noteId); + await expect + .poll(async () => recorder.snapshot().records.some(record => record.kind === 'websocket'), { timeout: 15000 }) + .toBe(true); + await recorder.stop(); + const fixture = recorder.snapshot(); + await recorderPage.close(); + + expect(validateFixture(fixture)).toEqual([]); + expect(fixture.records.some(record => record.kind === 'rest')).toBe(true); + expect(fixture.records.some(record => record.kind === 'websocket')).toBe(true); + }); +}); + +const sampleFixture = () => ({ + records: [ + { + kind: 'websocket', + sequence: 1, + websocket: { direction: 'send', payloadText: '{"op":"GET_NOTE","msgId":"client-1"}' } + }, + validRestRecord({ sequence: 2 }) + ], + version: fixtureVersion +}); + +const sampleFixtureWithReceive = () => ({ + records: [ + { + kind: 'websocket', + sequence: 1, + websocket: { direction: 'send', payloadText: '{"op":"GET_NOTE","msgId":"<msgId>"}' } + }, + { + kind: 'websocket', + sequence: 2, + websocket: { direction: 'receive', payloadText: '{"op":"NOTE"}' } + } + ], + version: fixtureVersion +}); + +const sampleFixtureWithTwoWebSocketRoundTrips = () => ({ + records: [ + { + kind: 'websocket', + sequence: 1, + websocket: { direction: 'send', payloadText: '{"op":"GET_NOTE","msgId":"<msgId>"}' } + }, + { + kind: 'websocket', + sequence: 2, + websocket: { direction: 'receive', payloadText: '{"op":"NOTE","noteId":"note-a"}' } + }, + { + kind: 'websocket', + sequence: 3, + websocket: { + direction: 'send', + payloadText: '{"op":"RUN_PARAGRAPH","paragraphId":"paragraph-b"}' + } + }, + { + kind: 'websocket', + sequence: 4, + websocket: { + direction: 'receive', + payloadText: '{"op":"PARAGRAPH","paragraphId":"paragraph-b"}' + } + } + ], + version: fixtureVersion +}); + +const validRestRecord = (overrides = {}) => ({ + kind: 'rest', + rest: restResponse('/api/notebook/note-a', { id: 'note-a', paragraphs: [{ id: 'paragraph-a' }] }), + sequence: 1, + ...overrides +}); + +const restResponse = (url: string, bodyJson: unknown) => ({ + bodyJson, + direction: 'response', + headers: { 'content-type': 'application/json' }, + request: { bodyRaw: '', headers: { accept: 'application/json' }, method: 'GET', url }, + status: 200 +}); + +const request = (method: string, url: string, body = '', headers = { accept: 'application/json' }) => ({ + headers: () => headers, + method: () => method, + postData: () => body, + url: () => url +}); + +const response = ( + sourceRequest: ReturnType<typeof request>, + status: number, + body: string | (() => Promise<string>), + headers = { 'content-type': 'application/json' } +) => ({ + headers: () => headers, + request: () => sourceRequest, + status: () => status, + text: async () => (typeof body === 'function' ? body() : body) +}); diff --git a/zeppelin-web-angular/package.json b/zeppelin-web-angular/package.json index 8acdda44a6..798f07ac55 100644 --- a/zeppelin-web-angular/package.json +++ b/zeppelin-web-angular/package.json @@ -14,6 +14,7 @@ "build:projects": "npm run build-project:sdk && npm run build-project:vis", "build-project:sdk": "ng build --project zeppelin-sdk", "check:websocket-contract": "node --test scripts/check-websocket-contract.test.js && node scripts/check-websocket-contract.js", + "check:core-contract-fixtures": "node --test e2e/core-contract/*.test.mjs", "build-project:vis": "ng build --project zeppelin-visualization", "lint": "cross-env NODE_OPTIONS='--max-old-space-size=8192' ng lint && npm run lint:react && prettier --check \"**/*.{ts,tsx,mts,js,json,css,html}\"", "lint:fix": "cross-env NODE_OPTIONS='--max-old-space-size=8192' ng lint --fix && npm run lint:fix:react && prettier --write \"**/*.{ts,tsx,mts,js,json,css,html}\"", diff --git a/zeppelin-web-angular/pom.xml b/zeppelin-web-angular/pom.xml index 3f2fee17ff..3653b15fb3 100644 --- a/zeppelin-web-angular/pom.xml +++ b/zeppelin-web-angular/pom.xml @@ -129,6 +129,18 @@ </configuration> </execution> + <execution> + <id>npm check core contract fixtures</id> + <goals> + <goal>npm</goal> + </goals> + <phase>test</phase> + <configuration> + <skip>${skipTests}</skip> + <arguments>run check:core-contract-fixtures</arguments> + </configuration> + </execution> + <execution> <id>npm e2e</id> <goals>
