From 98c4a195e896e56efbc8496b2d38a77acbc0b811 Mon Sep 17 00:00:00 2001 From: hypercross Date: Mon, 6 Jul 2026 15:53:52 +0800 Subject: [PATCH] refactor: replace websocket-stream with native ws utility Remove the `websocket-stream` dependency in favor of `createWebSocketStream` from the `ws` package. This simplifies the MQTT broker setup in the journal server and removes unnecessary transitive dependencies. Also improve the JournalPanel UI with smoother transitions and better visibility handling. --- package-lock.json | 67 +----------- package.json | 1 - src/cli/journal.ts | 36 ++++--- src/components/journal/JournalPanel.tsx | 26 +++-- src/components/stores/journalStream.ts | 131 ++++++++++++------------ 5 files changed, 108 insertions(+), 153 deletions(-) diff --git a/package-lock.json b/package-lock.json index 9318c03..cc5b115 100644 --- a/package-lock.json +++ b/package-lock.json @@ -29,7 +29,6 @@ "solid-js": "^1.9.3", "three": "^0.183.2", "three-3mf-exporter": "^45.1.0", - "websocket-stream": "^5.5.2", "ws": "^8.21.0", "yarn-spinner-loader": "^0.1.0", "zod": "^4.4.3" @@ -3724,12 +3723,6 @@ "integrity": "sha512-8+9WqebbFzpX9OR+Wa6O29asIogeRMzcGtAINdpMHHyAg10f05aSFVBbcEqGf/PXw1EjAZ+q2/bEBg3DvurK3Q==", "license": "Python-2.0" }, - "node_modules/async-limiter": { - "version": "1.0.1", - "resolved": "https://registry.npmjs.org/async-limiter/-/async-limiter-1.0.1.tgz", - "integrity": "sha512-csOlWGAcRFJaI6m+F2WKdnMKr4HhdhFVBk0H/QbJFMCr+uO2kwohwXQPxw/9OCxp05r5ghVBFSyioixx3gfkNQ==", - "license": "MIT" - }, "node_modules/asynckit": { "version": "0.4.0", "resolved": "https://registry.npmjs.org/asynckit/-/asynckit-0.4.0.tgz", @@ -5267,18 +5260,6 @@ "node": ">= 0.4" } }, - "node_modules/duplexify": { - "version": "3.7.1", - "resolved": "https://registry.npmjs.org/duplexify/-/duplexify-3.7.1.tgz", - "integrity": "sha512-07z8uv2wMyS51kKhD1KsdXJg5WQ6t93RneqRxUHnskXVtlYYkLqM0gqStQZ3pj073g687jPCHrqNfCzawLYh5g==", - "license": "MIT", - "dependencies": { - "end-of-stream": "^1.0.0", - "inherits": "^2.0.1", - "readable-stream": "^2.0.0", - "stream-shift": "^1.0.0" - } - }, "node_modules/electron-to-chromium": { "version": "1.5.331", "resolved": "https://registry.npmjs.org/electron-to-chromium/-/electron-to-chromium-1.5.331.tgz", @@ -5306,15 +5287,6 @@ "dev": true, "license": "MIT" }, - "node_modules/end-of-stream": { - "version": "1.4.5", - "resolved": "https://registry.npmjs.org/end-of-stream/-/end-of-stream-1.4.5.tgz", - "integrity": "sha512-ooEGc6HP26xXq/N+GCGOT0JKCLDGrq2bQUZrQ7gyrJiZANJ/8YDTxTpQBXGMn+WbIQXNVpyWymm7KYVICQnyOg==", - "license": "MIT", - "dependencies": { - "once": "^1.4.0" - } - }, "node_modules/enhanced-resolve": { "version": "5.20.1", "resolved": "https://registry.npmjs.org/enhanced-resolve/-/enhanced-resolve-5.20.1.tgz", @@ -8838,6 +8810,7 @@ "version": "1.4.0", "resolved": "https://registry.npmjs.org/once/-/once-1.4.0.tgz", "integrity": "sha512-lNaJgI+2Q5URQBkccEKHTQOPaXdUxnZZElQTZY0MFUAuaEqe1E+Nyvgdz/aIyNi6Z9MzO5dv1H8n58/GELp3+w==", + "dev": true, "license": "ISC", "dependencies": { "wrappy": "1" @@ -9783,12 +9756,6 @@ "node": ">= 0.8" } }, - "node_modules/stream-shift": { - "version": "1.0.3", - "resolved": "https://registry.npmjs.org/stream-shift/-/stream-shift-1.0.3.tgz", - "integrity": "sha512-76ORR0DO1o1hlKwTbi/DM3EXWGf3ZJYO8cXX5RJwnul2DEg2oyoZyjLNoQM8WsvZiFKCRfC1O0J7iCvie3RZmQ==", - "license": "MIT" - }, "node_modules/string_decoder": { "version": "1.1.1", "resolved": "https://registry.npmjs.org/string_decoder/-/string_decoder-1.1.1.tgz", @@ -10297,12 +10264,6 @@ "node": ">=0.8.0" } }, - "node_modules/ultron": { - "version": "1.1.1", - "resolved": "https://registry.npmjs.org/ultron/-/ultron-1.1.1.tgz", - "integrity": "sha512-UIEXBNeYmKptWH6z8ZnqTeS8fV74zG0/eRU9VGkpzz+LIJNs8W/zM/L+7ctCkRrgbNnnR0xxw4bKOr0cW0N0Og==", - "license": "MIT" - }, "node_modules/undici-types": { "version": "6.21.0", "resolved": "https://registry.npmjs.org/undici-types/-/undici-types-6.21.0.tgz", @@ -10586,31 +10547,6 @@ "node": ">=12" } }, - "node_modules/websocket-stream": { - "version": "5.5.2", - "resolved": "https://registry.npmjs.org/websocket-stream/-/websocket-stream-5.5.2.tgz", - "integrity": "sha512-8z49MKIHbGk3C4HtuHWDtYX8mYej1wWabjthC/RupM9ngeukU4IWoM46dgth1UOS/T4/IqgEdCDJuMe2039OQQ==", - "license": "BSD-2-Clause", - "dependencies": { - "duplexify": "^3.5.1", - "inherits": "^2.0.1", - "readable-stream": "^2.3.3", - "safe-buffer": "^5.1.2", - "ws": "^3.2.0", - "xtend": "^4.0.0" - } - }, - "node_modules/websocket-stream/node_modules/ws": { - "version": "3.3.3", - "resolved": "https://registry.npmjs.org/ws/-/ws-3.3.3.tgz", - "integrity": "sha512-nnWLa/NwZSt4KQJu51MYlCcSQ5g7INpOrOMt4XV8j4dqTXdmlUmSHQ8/oLC069ckre0fRsgfvsKwbTdtKLCDkA==", - "license": "MIT", - "dependencies": { - "async-limiter": "~1.0.0", - "safe-buffer": "~5.1.0", - "ultron": "~1.1.0" - } - }, "node_modules/whatwg-encoding": { "version": "2.0.0", "resolved": "https://registry.npmjs.org/whatwg-encoding/-/whatwg-encoding-2.0.0.tgz", @@ -10741,6 +10677,7 @@ "version": "1.0.2", "resolved": "https://registry.npmjs.org/wrappy/-/wrappy-1.0.2.tgz", "integrity": "sha512-l4Sp/DRseor9wL6EvV2+TuQn63dMkPjZ/sp9XkghTEbV9KlPS1xUsZ3u7/IQO4wxtcFB4bgpQPRcR3QCvezPcQ==", + "dev": true, "license": "ISC" }, "node_modules/write-file-atomic": { diff --git a/package.json b/package.json index 4181d4b..3896d16 100644 --- a/package.json +++ b/package.json @@ -52,7 +52,6 @@ "solid-js": "^1.9.3", "three": "^0.183.2", "three-3mf-exporter": "^45.1.0", - "websocket-stream": "^5.5.2", "ws": "^8.21.0", "yarn-spinner-loader": "^0.1.0", "zod": "^4.4.3" diff --git a/src/cli/journal.ts b/src/cli/journal.ts index 1b59c0a..4809196 100644 --- a/src/cli/journal.ts +++ b/src/cli/journal.ts @@ -2,7 +2,10 @@ * Journal stream persistence — MQTT broker + JSONL append + session manifest * * Runs an aedes MQTT broker over WebSocket, attached to the existing - * HTTP server via the upgrade event. No separate TCP port. + * HTTP server. Uses ws v8's built-in createWebSocketStream. + * + * Pattern from aedes docs: + * https://github.com/moscajs/aedes#mqtt-server-over-websocket * * Persistence uses aedes's internal subscribe/publish API, so the * server process doesn't need to connect to itself as an MQTT client. @@ -20,7 +23,7 @@ import { } from "fs"; import type { Server as HttpServer } from "http"; import type { AedesPublishPacket } from "aedes"; -import websocketStream from "websocket-stream"; +import { WebSocketServer, createWebSocketStream } from "ws"; // --------------------------------------------------------------------------- // Types @@ -101,18 +104,21 @@ export async function createJournalServer( const { Aedes: AedesFactory } = await import("aedes"); const broker = new AedesFactory(); - // ---- WebSocket server (attached to HTTP server) ---- - const { WebSocketServer } = await import("ws"); - const wsServer = new WebSocketServer({ noServer: true }); + // ---- WebSocket server (attached to the existing HTTP server) ---- + // Pattern from aedes docs: pass { server } to WSS, use 'connection' event + const wss = new WebSocketServer({ server: httpServer }); - httpServer.on("upgrade", (req, socket, head) => { - wsServer.handleUpgrade(req, socket, head, (ws) => { - // websocket-stream's types reference an older ws version; cast works. - const stream = websocketStream( - ws as unknown as Parameters[0], - ); - broker.handle(stream); - }); + wss.on("connection", (ws, req) => { + console.log(`[journal] ws connected from ${req.socket.remoteAddress}`); + const stream = createWebSocketStream(ws); + broker.handle(stream, req); + }); + + broker.on("client", (client: import("aedes").Client) => { + console.log(`[journal] mqtt client ready: ${client.id}`); + }); + broker.on("clientDisconnect", (client: import("aedes").Client) => { + console.log(`[journal] mqtt client disconnected: ${client.id}`); }); // ---- Persistence (internal subscribe, no separate MQTT client) ---- @@ -149,10 +155,10 @@ export async function createJournalServer( return { broker, - wsServer, + wsServer: wss, async close() { console.log("[journal] shutting down..."); - await new Promise((r) => wsServer.close(() => r())); + await new Promise((r) => wss.close(() => r())); await new Promise((r) => broker.close(() => r())); }, }; diff --git a/src/components/journal/JournalPanel.tsx b/src/components/journal/JournalPanel.tsx index e19d9e7..d2dd146 100644 --- a/src/components/journal/JournalPanel.tsx +++ b/src/components/journal/JournalPanel.tsx @@ -25,17 +25,24 @@ export const JournalPanel: Component = (props) => { const stream = useJournalStream(); return ( - + <> + {/* Backdrop — hidden on desktop since panel doesn't overlay content */}
+ {/* Panel — always mounted for exit animation, visibility toggled */} - + ); }; diff --git a/src/components/stores/journalStream.ts b/src/components/stores/journalStream.ts index ec74d96..29d05c2 100644 --- a/src/components/stores/journalStream.ts +++ b/src/components/stores/journalStream.ts @@ -198,77 +198,82 @@ export async function connectStream( setState("connectionStatus", "connecting"); setState("connectionError", null); - const client = mqtt.connect(brokerUrl, { - clientId: `${state.myName}-${Date.now()}`, - protocol: brokerUrl.startsWith("wss") ? "wss" : "ws", - reconnectPeriod: 2000, - }); - - _mqttClient = client; - - client.on("connect", () => { - _mqttConnected = true; - setState("connected", true); - setState("connectionStatus", "connected"); - setState("connectionError", null); - setState("brokerUrl", brokerUrl); - - // Persist connection info for next time - if (typeof localStorage !== "undefined") { - localStorage.setItem(LS_BROKER_URL, brokerUrl); - localStorage.setItem(LS_LAST_SESSION, sessionId); - } - - client.subscribe(`ttrpg/${sessionId}/stream`, { qos: 1 }, (err) => { - if (err) console.error("[stream] stream sub err:", err); + return new Promise((resolve, reject) => { + const client = mqtt.connect(brokerUrl, { + clientId: `${state.myName}-${Date.now()}`, + protocol: brokerUrl.startsWith("wss") ? "wss" : "ws", + reconnectPeriod: 2000, }); - client.subscribe($SESSIONS, { qos: 1 }, (err) => { - if (err) console.error("[stream] sessions sub err:", err); + + _mqttClient = client; + + client.on("connect", () => { + _mqttConnected = true; + setState("connected", true); + setState("connectionStatus", "connected"); + setState("connectionError", null); + setState("brokerUrl", brokerUrl); + + // Persist connection info for next time + if (typeof localStorage !== "undefined") { + localStorage.setItem(LS_BROKER_URL, brokerUrl); + localStorage.setItem(LS_LAST_SESSION, sessionId); + } + + client.subscribe(`ttrpg/${sessionId}/stream`, { qos: 1 }, (err) => { + if (err) console.error("[stream] stream sub err:", err); + }); + client.subscribe($SESSIONS, { qos: 1 }, (err) => { + if (err) console.error("[stream] sessions sub err:", err); + }); + client.subscribe(`ttrpg/${sessionId}/meta`, { qos: 1 }); + + resolve(); }); - client.subscribe(`ttrpg/${sessionId}/meta`, { qos: 1 }); - }); - client.on("message", (topic, payload) => { - const raw = payload.toString(); - const parts = topic.split("/"); + client.on("error", (err) => { + console.error("[stream] mqtt error:", err); + setState("connectionStatus", "error"); + setState("connectionError", err.message); + reject(err); + }); - if (topic === $SESSIONS) { - // Session manifest update - try { - const manifest: SessionManifest = JSON.parse(raw); - setSessionList(manifest); - } catch (e) { - console.error("[stream] manifest parse err:", e); + client.on("close", () => { + _mqttConnected = false; + setState("connected", false); + setState("connectionStatus", "disconnected"); + }); + + client.on("message", (topic, payload) => { + const raw = payload.toString(); + const parts = topic.split("/"); + + if (topic === $SESSIONS) { + // Session manifest update + try { + const manifest: SessionManifest = JSON.parse(raw); + setSessionList(manifest); + } catch (e) { + console.error("[stream] manifest parse err:", e); + } + return; } - return; - } - if (parts.length >= 3 && parts[0] === "ttrpg" && parts[2] === "meta") { - // Session meta update — manifest will be republished by server, - // handled via $SESSIONS subscription above. - return; - } - - if (parts.length >= 3 && parts[0] === "ttrpg" && parts[2] === "stream") { - try { - const msg: StreamMessage = JSON.parse(raw); - receiveMessage(msg); - } catch (e) { - console.error("[stream] malformed message:", e); + if (parts.length >= 3 && parts[0] === "ttrpg" && parts[2] === "meta") { + // Session meta update — manifest will be republished by server, + // handled via $SESSIONS subscription above. + return; } - } - }); - client.on("close", () => { - _mqttConnected = false; - setState("connected", false); - setState("connectionStatus", "disconnected"); - }); - - client.on("error", (err) => { - console.error("[stream] mqtt error:", err); - setState("connectionStatus", "error"); - setState("connectionError", err.message); + if (parts.length >= 3 && parts[0] === "ttrpg" && parts[2] === "stream") { + try { + const msg: StreamMessage = JSON.parse(raw); + receiveMessage(msg); + } catch (e) { + console.error("[stream] malformed message:", e); + } + } + }); }); }