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.
This commit is contained in:
+21
-15
@@ -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<typeof websocketStream>[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<void>((r) => wsServer.close(() => r()));
|
||||
await new Promise<void>((r) => wss.close(() => r()));
|
||||
await new Promise<void>((r) => broker.close(() => r()));
|
||||
},
|
||||
};
|
||||
|
||||
Reference in New Issue
Block a user