diff --git a/README.md b/README.md index 31e6efac..a16ba775 100644 --- a/README.md +++ b/README.md @@ -55,6 +55,7 @@ js_eval("console.log")('hello, world') - [done] CommonJS module system .py loader, loads Python modules for use by JS - [done] Python host environment supplies event loop, including EventEmitter, setTimeout, etc. - [done] Python host environment supplies XMLHttpRequest +- [done] Python host environment supplies WebSocket - [done] Python TypedArrays coerce to JS TypeArrays - [done] JS TypedArrays coerce to Python TypeArrays - [done] Python lists coerce to JS Arrays diff --git a/python/pythonmonkey/__init__.py b/python/pythonmonkey/__init__.py index 957713ff..b101ea8f 100644 --- a/python/pythonmonkey/__init__.py +++ b/python/pythonmonkey/__init__.py @@ -15,3 +15,4 @@ require("timers") require("url") require("XMLHttpRequest") +require("WebSocket") diff --git a/python/pythonmonkey/builtin_modules/WebSocket-internal.d.ts b/python/pythonmonkey/builtin_modules/WebSocket-internal.d.ts new file mode 100644 index 00000000..9c3cafe1 --- /dev/null +++ b/python/pythonmonkey/builtin_modules/WebSocket-internal.d.ts @@ -0,0 +1,30 @@ +/** + * @file WebSocket-internal.d.ts + * @brief TypeScript type declarations for the internal WebSocket helpers + * @author Dan Desjardins + * @date September 2026 + * + * @copyright Copyright (c) 2026 Distributive Corp. + */ + +/** + * Open a WebSocket connection and pump its messages into the callbacks. + * Resolves once the connection has closed; `onClose` is called exactly once. + */ +export declare function wsConnect( + url: string, + protocols: string[], + // called before the 'open' event, with the negotiated subprotocol and the send/close functions + onOpen: ( + protocol: string, + sendText: (data: string) => Promise, + sendBinary: (data: Uint8Array) => Promise, + close: (code: number, reason: string) => Promise, + ) => void, + onMessage: (data: string | Uint8Array, isBinary: boolean) => void, + onError: (message: string) => void, + onClose: (code: number, reason: string) => void, + // the debug logging function + /** See `pm.bootstrap.require("debug")` */ + debug: (selector: string) => ((...args: string[]) => void), +): Promise; diff --git a/python/pythonmonkey/builtin_modules/WebSocket-internal.py b/python/pythonmonkey/builtin_modules/WebSocket-internal.py new file mode 100644 index 00000000..e7164152 --- /dev/null +++ b/python/pythonmonkey/builtin_modules/WebSocket-internal.py @@ -0,0 +1,87 @@ +# @file WebSocket-internal.py +# @brief internal helper functions for WebSocket, backed by aiohttp +# @author Dan Desjardins +# @date September 2026 +# @copyright Copyright (c) 2026 Distributive Corp. + +import aiohttp +from typing import Callable, List, Union + + +async def wsConnect( + url: str, + protocols: List[str], + onOpen: Callable[[str, Callable, Callable, Callable], None], + onMessage: Callable[[Union[str, bytearray], bool], None], + onError: Callable[[str], None], + onClose: Callable[[int, str], None], + debug: Callable[[str], Callable[..., None]], + / +): + """ + Open a WebSocket and pump its messages into the JS-side callbacks. + + onOpen receives the negotiated subprotocol and the send/close functions in + the same synchronous call that precedes the JS 'open' event, so a message + sent from an 'open' handler can never race the functions being wired up. + Every path ends with exactly one onClose call. + """ + log = debug('ws:io') + session = aiohttp.ClientSession() + + try: + ws = await session.ws_connect(url, protocols=tuple(protocols)) + except Exception as e: + onError(str(e)) + await session.close() + onClose(1006, '') + return + + def guard(fn): + # JS calls these without awaiting the returned promise, so failures are + # reported through onError instead of becoming unhandled rejections. + async def wrapper(*args): + if ws.closed: + return + try: + await fn(*args) + except Exception as e: + onError(str(e)) + return wrapper + + sendText = guard(ws.send_str) + sendBinary = guard(lambda data: ws.send_bytes(bytes(data))) + closeConn = guard(lambda code, reason: ws.close(code=code, message=reason.encode('utf-8'))) + + onOpen(ws.protocol or '', sendText, sendBinary, closeConn) + + close_reason = '' + try: + while True: + # receive() rather than `async for`, which swallows the CLOSE frame's reason + msg = await ws.receive() + log('received', msg.type.name) + if msg.type == aiohttp.WSMsgType.TEXT: + onMessage(msg.data, False) + elif msg.type == aiohttp.WSMsgType.BINARY: + onMessage(bytearray(msg.data), True) + elif msg.type == aiohttp.WSMsgType.ERROR: + onError(str(ws.exception())) + break + elif msg.type == aiohttp.WSMsgType.CLOSE: + close_reason = msg.extra or '' + break + elif msg.type == aiohttp.WSMsgType.CLOSED: + break + # CLOSING means our own close() is in flight; the next receive() finishes + # the handshake and reports CLOSED with the real close code + except Exception as e: + onError(str(e)) + finally: + await session.close() + # close_code is unset when the connection dropped without a close frame + onClose(ws.close_code or 1006, close_reason) + + +# Module exports +exports['wsConnect'] = wsConnect # type: ignore diff --git a/python/pythonmonkey/builtin_modules/WebSocket.js b/python/pythonmonkey/builtin_modules/WebSocket.js new file mode 100644 index 00000000..db374a5c --- /dev/null +++ b/python/pythonmonkey/builtin_modules/WebSocket.js @@ -0,0 +1,187 @@ +/** + * @file WebSocket.js + * Implement the WebSocket API, backed by Python's aiohttp + * WebSocket client (WebSocket-internal.py). + * @author Dan Desjardins + * @date September 2026 + * + * @copyright Copyright (c) 2026 Distributive Corp. + */ +'use strict'; + +const { EventTarget, Event } = require('event-target'); +const { DOMException } = require('dom-exception'); +const { URL } = require('url'); +const { wsConnect } = require('WebSocket-internal'); +const debug = globalThis.python.eval('__import__("pythonmonkey").bootstrap.require')('debug'); + +// exposed +class MessageEvent extends Event +{ + constructor(type, eventInitDict = {}) + { + super(type); + this.data = eventInitDict.data; + } +} + +// exposed +class CloseEvent extends Event +{ + constructor(type, eventInitDict = {}) + { + super(type); + this.code = eventInitDict.code ?? 1000; + this.reason = eventInitDict.reason ?? ''; + this.wasClean = eventInitDict.wasClean ?? true; + } +} + +/** + * Implement the `WebSocket` API according to the spec, backed by aiohttp. + * @see https://websockets.spec.whatwg.org/ + */ +class WebSocket extends EventTarget +{ + /** @readonly */ static CONNECTING = 0; + /** @readonly */ static OPEN = 1; + /** @readonly */ static CLOSING = 2; + /** @readonly */ static CLOSED = 3; + + /** @readonly */ CONNECTING = 0; + /** @readonly */ OPEN = 1; + /** @readonly */ CLOSING = 2; + /** @readonly */ CLOSED = 3; + + // event handlers -- EventTarget#dispatchEvent auto-invokes these + onopen = null; + onmessage = null; + onerror = null; + onclose = null; + + #readyState = WebSocket.CONNECTING; + #conn = null; + #url; + #protocol = ''; + + // engine.io-client's WebSocket transport calls `this.ws._socket.unref()` on + // open when its autoUnref option is set, a Node `ws`-library detail that a + // browser-style WebSocket doesn't have. Nothing here needs unref'ing. + _socket = { unref() {}, ref() {} }; + + /** + * @param {string | URL} url + * @param {string | string[]} [protocols] + */ + constructor(url, protocols) + { + super(); + const parsedURL = new URL(url); + if (!['ws:', 'wss:'].includes(parsedURL.protocol)) + throw new DOMException(`Invalid WebSocket URL scheme "${parsedURL.protocol}"`, 'SyntaxError'); + this.#url = parsedURL.href; + + const protoArray = protocols ? (Array.isArray(protocols) ? protocols : [protocols]) : []; + + // aiohttp's ws_connect() expects a plain http(s):// URL, not ws(s):// + const httpURL = this.#url.replace(/^ws/, 'http'); + + debug('ws:connect')(`connecting to ${httpURL}`); + + wsConnect( + httpURL, + protoArray, + (protocol, sendText, sendBinary, closeFn) => // onOpen + { + this.#conn = { sendText, sendBinary, close: closeFn }; + if (this.#readyState === WebSocket.CLOSING) // close() was called while connecting + { + closeFn(1000, ''); + return; + } + this.#protocol = protocol; + this.#readyState = WebSocket.OPEN; + debug('ws:open')(`connected to ${this.#url}`); + this.dispatchEvent(new Event('open')); + }, + (data, isBinary) => // onMessage + { + // copy binary payloads out of the Python bytearray into a JS-owned ArrayBuffer + const payload = isBinary ? new Uint8Array(data).buffer : data; + this.dispatchEvent(new MessageEvent('message', { data: payload })); + }, + (message) => // onError + { + debug('ws:error')(message); + this.dispatchEvent(new Event('error')); + }, + (code, reason) => // onClose + { + this.#readyState = WebSocket.CLOSED; + debug('ws:close')(`closed, code=${code} reason=${reason}`); + // 1006 is reserved for connections that dropped without a close handshake + this.dispatchEvent(new CloseEvent('close', { code, reason, wasClean: code !== 1006 })); + }, + debug, + ).catch((e) => // only reachable if a callback above threw past the Python side + { + debug('ws:error')(String(e)); + if (this.#readyState === WebSocket.CLOSED) + return; + this.#readyState = WebSocket.CLOSED; + this.dispatchEvent(new Event('error')); + this.dispatchEvent(new CloseEvent('close', { code: 1006, reason: String(e), wasClean: false })); + }); + } + + get readyState() { return this.#readyState; } + get url() { return this.#url; } + get protocol() { return this.#protocol; } + get bufferedAmount() { return 0; } // not tracked + + /** + * @param {string | ArrayBuffer | ArrayBufferView} data + */ + send(data) + { + if (this.#readyState === WebSocket.CONNECTING) + throw new DOMException('WebSocket is still connecting (readyState CONNECTING)', 'InvalidStateError'); + if (this.#readyState !== WebSocket.OPEN) + return; // per spec: silently discard if not OPEN + if (typeof data === 'string') + this.#conn.sendText(data); + else if (data instanceof ArrayBuffer) + this.#conn.sendBinary(new Uint8Array(data)); + else if (ArrayBuffer.isView(data)) + this.#conn.sendBinary(new Uint8Array(data.buffer, data.byteOffset, data.byteLength)); + else + throw new TypeError('WebSocket.send() data must be a string, ArrayBuffer or ArrayBufferView'); + } + + /** + * @param {number} [code] + * @param {string} [reason] + */ + close(code = 1000, reason = '') + { + if (this.#readyState === WebSocket.CLOSING || this.#readyState === WebSocket.CLOSED) + return; + this.#readyState = WebSocket.CLOSING; + if (this.#conn) // otherwise still connecting; onOpen closes it + this.#conn.close(code, reason); + } +} + +/* A side-effect of loading this module is to add WebSocket and related symbols to the global + * object, matching XMLHttpRequest.js, so code written for browsers works without a require(). + */ +if (!globalThis.WebSocket) + globalThis.WebSocket = WebSocket; +if (!globalThis.MessageEvent) + globalThis.MessageEvent = MessageEvent; +if (!globalThis.CloseEvent) + globalThis.CloseEvent = CloseEvent; + +exports.WebSocket = WebSocket; +exports.MessageEvent = MessageEvent; +exports.CloseEvent = CloseEvent; diff --git a/tests/python/test_websocket.py b/tests/python/test_websocket.py new file mode 100644 index 00000000..0c9d6d5d --- /dev/null +++ b/tests/python/test_websocket.py @@ -0,0 +1,106 @@ +import asyncio +import socket +import aiohttp +import aiohttp.web +import pythonmonkey as pm + + +def free_port(): + with socket.socket() as s: + s.bind(('127.0.0.1', 0)) + return s.getsockname()[1] + + +async def ws_handler(request): + ws = aiohttp.web.WebSocketResponse(protocols=('chat',)) + await ws.prepare(request) + async for msg in ws: + if msg.type == aiohttp.WSMsgType.TEXT: + if msg.data == 'close-me': + await ws.close(code=4000, message=b'bye') + else: + await ws.send_str('echo:' + msg.data) + elif msg.type == aiohttp.WSMsgType.BINARY: + await ws.send_bytes(bytes(reversed(msg.data))) + return ws + + +def test_websocket(): + async def async_fn(): + port = free_port() + app = aiohttp.web.Application() + app.router.add_get('/ws', ws_handler) + runner = aiohttp.web.AppRunner(app) + await runner.setup() + await aiohttp.web.TCPSite(runner, '127.0.0.1', port).start() + + # text and binary round trips, subprotocol negotiation, server-initiated close + log = await pm.eval(""" + (port) => new Promise((resolve, reject) => { + const log = []; + const ws = new WebSocket(`ws://127.0.0.1:${port}/ws`, ['chat']); + log.push('state=' + ws.readyState); + ws.onopen = () => { + log.push(`open protocol=${ws.protocol} state=${ws.readyState}`); + ws.send('hello'); + }; + ws.onmessage = (ev) => { + if (typeof ev.data === 'string') { + log.push('text:' + ev.data); + ws.send(new Uint8Array([1, 2, 3])); + } else { + log.push('binary:' + Array.from(new Uint8Array(ev.data)).join(',')); + ws.send('close-me'); + } + }; + ws.onerror = () => log.push('error'); + ws.onclose = (ev) => { + log.push(`close code=${ev.code} reason=${ev.reason} clean=${ev.wasClean} state=${ws.readyState}`); + resolve(log); + }; + setTimeout(() => reject(new Error('timeout: ' + log.join(' | '))), 5000); + }) + """)(port) + assert list(log) == [ + 'state=0', + 'open protocol=chat state=1', + 'text:echo:hello', + 'binary:3,2,1', + 'close code=4000 reason=bye clean=true state=3', + ] + + # refused connection fires error then close, and ends CLOSED + log = await pm.eval(""" + (port) => new Promise((resolve, reject) => { + const log = []; + const ws = new WebSocket(`ws://127.0.0.1:${port}/`); + ws.onopen = () => log.push('open'); + ws.onerror = () => log.push('error state=' + ws.readyState); + ws.onclose = (ev) => { + log.push(`close code=${ev.code} clean=${ev.wasClean} state=${ws.readyState}`); + resolve(log); + }; + setTimeout(() => reject(new Error('timeout: ' + log.join(' | '))), 5000); + }) + """)(free_port()) + assert list(log) == ['error state=0', 'close code=1006 clean=false state=3'] + + # close() while still connecting never fires 'open' + log = await pm.eval(""" + (port) => new Promise((resolve, reject) => { + const log = []; + const ws = new WebSocket(`ws://127.0.0.1:${port}/ws`); + ws.close(); + log.push('state=' + ws.readyState); + ws.onopen = () => log.push('open'); + ws.onclose = (ev) => { + log.push(`close code=${ev.code} state=${ws.readyState}`); + resolve(log); + }; + setTimeout(() => reject(new Error('timeout: ' + log.join(' | '))), 5000); + }) + """)(port) + assert list(log) == ['state=2', 'close code=1000 state=3'] + + await runner.cleanup() + asyncio.run(async_fn())