diff --git a/package-lock.json b/package-lock.json index 805abbe..46e31b4 100644 --- a/package-lock.json +++ b/package-lock.json @@ -14,6 +14,7 @@ "@thisbeyond/solid-dnd": "^0.7.5", "aedes": "^1.1.1", "chokidar": "^5.0.0", + "comlink": "^4.4.2", "commander": "^14.0.3", "csv-parse": "^6.1.0", "csv-stringify": "^6.7.0", @@ -4332,6 +4333,12 @@ "node": ">= 0.8" } }, + "node_modules/comlink": { + "version": "4.4.2", + "resolved": "https://registry.npmjs.org/comlink/-/comlink-4.4.2.tgz", + "integrity": "sha512-OxGdvBmJuNKSCMO4NTl1L47VRp6xn2wG4F/2hYzB6tiCb709otOxtEYCSvK80PtjODfXXZu8ds+Nw5kVCjqd2g==", + "license": "Apache-2.0" + }, "node_modules/commander": { "version": "14.0.3", "resolved": "https://registry.npmjs.org/commander/-/commander-14.0.3.tgz", diff --git a/package.json b/package.json index 2c6afbb..9e80a2e 100644 --- a/package.json +++ b/package.json @@ -47,6 +47,7 @@ "@thisbeyond/solid-dnd": "^0.7.5", "aedes": "^1.1.1", "chokidar": "^5.0.0", + "comlink": "^4.4.2", "commander": "^14.0.3", "csv-parse": "^6.1.0", "csv-stringify": "^6.7.0", diff --git a/src/components/stores/journalStream.ts b/src/components/stores/journalStream.ts index 0e5b5be..ee04dfa 100644 --- a/src/components/stores/journalStream.ts +++ b/src/components/stores/journalStream.ts @@ -1,20 +1,21 @@ /** - * Journal Stream — Client Store + * Journal Stream — Client Store (Worker Proxy) * - * Reactive state for a single session's message stream. Manages MQTT - * connection, local message log, per-sender sequence tracking, and the - * revealed-paths set (populated by the link reducer). + * This is the main-thread side of the journal stream. It spawns a Shared + * Worker that owns the MQTT connection. All state flows from the worker + * to this store via Comlink callbacks. * - * Session lifecycle (create/list/delete) is handled via MQTT retained - * topics — ttrpg/$SESSIONS for the manifest and ttrpg/{id}/meta per session. - * - * No persistence here — that's the CLI server's job via JSONL append. + * The public API is unchanged — components still call `useJournalStream()`, + * `sendMessage()`, etc. exactly as before. */ import { createStore, produce } from "solid-js/store"; import { createSignal } from "solid-js"; +import * as Comlink from "comlink"; import type { StreamMessage } from "../journal/registry"; import { getMessageType, validatePayload } from "../journal/registry"; +import type { WorkerState, WorkerPatch, SessionManifest } from "../../workers/journal.worker"; +import type { JournalWorkerAPI } from "../../workers/journal.worker"; import { loadPersisted, saveName, @@ -32,38 +33,17 @@ import { export interface JournalStreamState { sessionId: string | null; - /** Human-readable name for the current session (from manifest) */ sessionName: string | null; - /** Full message log, oldest-first */ messages: StreamMessage[]; - /** Last sequence number per sender */ senderSeq: Record; - /** - * Paths and sections revealed by link messages. - * Key: normalized path (no .md). Value: set of revealed section slugs. - * An empty set means the whole article is revealed. - * Populated during hydration and live receipt via the type's reducer. - */ revealedPaths: Record>; - /** MQTT connection status */ connected: boolean; - /** Granular connection state for UI indicators */ connectionStatus: "disconnected" | "connecting" | "connected" | "error"; - /** Last connection error message, if any */ connectionError: string | null; - /** This client's identity */ myName: string; - /** Role: gm | player | observer. Immutable while connected. */ myRole: "gm" | "player" | "observer"; - /** Broker URL, set after connect */ brokerUrl: string | null; - /** Active player list (keyed by player name) */ players: Record; - /** - * Stat values set via /stat set/del/roll commands. - * Key: stat key (e.g. "strength", "alice:hp"). Value: string. - * Populated during hydration and live receipt via the stat type's reducer. - */ stats: Record; } @@ -73,8 +53,24 @@ export interface SessionMeta { players: string[]; } -export interface SessionManifest { - sessions: Record; +export type { SessionManifest } from "../../workers/journal.worker"; + +// --------------------------------------------------------------------------- +// Worker connection +// --------------------------------------------------------------------------- + +let _workerAPI: Comlink.Remote | null = null; +let _unsubscribe: (() => void) | null = null; + +function getWorkerAPI(): Comlink.Remote { + if (!_workerAPI) { + const worker = new SharedWorker( + new URL("../../workers/journal.worker.ts", import.meta.url), + ); + _workerAPI = Comlink.wrap(worker.port); + worker.port.start(); + } + return _workerAPI; } // --------------------------------------------------------------------------- @@ -84,7 +80,6 @@ export interface SessionManifest { const persisted = loadPersisted(); const urlParams = readUrlParams(); -// URL params override localStorage if present const initialName = urlParams.playerName ?? persisted.myName; const initialSession = urlParams.sessionId ?? persisted.lastSessionId; @@ -104,7 +99,6 @@ const [state, setState] = createStore({ stats: {}, }); -// Sync initial URL params if they came from localStorage (not URL) if (initialName && !urlParams.playerName) syncUrlParam("player", initialName); if (initialSession && !urlParams.sessionId) syncUrlParam("session", initialSession); @@ -117,56 +111,10 @@ const [sessionList, setSessionList] = createSignal({ export { sessionList as sessions }; -/** - * Change the current player's name. Persisted to localStorage so it - * survives page reloads. Also syncs to URL search param. - */ -export function setMyName(name: string): void { - setState("myName", name); - saveName(name); - syncUrlParam("player", name); -} - -/** - * Change the current player's role. Persisted to localStorage. - * Only callable when disconnected. - */ -export function setMyRole(role: "gm" | "player" | "observer"): void { - setState("myRole", role); - saveRole(role); -} - -/** - * Set the active session ID and sync to URL. Also resolves the human-readable - * session name from the cached manifest. - */ -export function setSessionId(id: string | null): void { - setState("sessionId", id); - if (id) { - saveSessionId(id); - syncUrlParam("session", id); - // Resolve session name from current manifest - const manifest = sessionList(); - const name = manifest.sessions[id]?.name ?? null; - setState("sessionName", name); - } else { - removeUrlParam("session"); - setState("sessionName", null); - } -} - -// Will hold the MQTT client instance after connect() -let _mqttClient: import("mqtt").MqttClient | null = null; -let _mqttConnected = false; - // --------------------------------------------------------------------------- -// Helpers +// Reducer runner — runs locally in each tab // --------------------------------------------------------------------------- -function makeMessageId(sender: string, seq: number): string { - return `${sender}-${seq}`; -} - function runReducer(msg: StreamMessage): void { const def = getMessageType(msg.type); if (def?.reducer) { @@ -174,260 +122,186 @@ function runReducer(msg: StreamMessage): void { } } -const $SESSIONS = "ttrpg/$SESSIONS"; - -// --------------------------------------------------------------------------- -// Hydration (initial load from server) -// --------------------------------------------------------------------------- - -/** - * Load the full message history from the static server's JSONL file. - * The file is served from the same HTTP origin as the web app. - * Runs all reducers in order. - */ -export async function hydrateFromServer(sessionId: string): Promise { - const response = await fetch( - `/.ttrpg/sessions/${encodeURIComponent(sessionId)}/stream.jsonl`, - ); - - if (!response.ok) { - if (response.status === 404) { - return; // fresh session, no file yet - } - throw new Error(`Failed to load session: ${response.statusText}`); +/** Replay all messages through reducers to rebuild derived state. */ +function replayReducers(): void { + // Reset derived state + setState("revealedPaths", {}); + setState("stats", {}); + for (const msg of state.messages) { + runReducer(msg); } +} - const text = await response.text(); - const lines = text.split("\n").filter((l) => l.trim()); +// --------------------------------------------------------------------------- +// Worker patch handler +// --------------------------------------------------------------------------- - const messages: StreamMessage[] = []; - const senderSeq: Record = {}; +function handlePatch(patch: WorkerPatch): void { + switch (patch.type) { + case "fullState": { + setState( + produce((s) => { + s.sessionId = patch.state.sessionId; + s.sessionName = patch.state.sessionName; + s.messages = patch.state.messages; + s.senderSeq = patch.state.senderSeq; + s.connected = patch.state.connected; + s.connectionStatus = patch.state.connectionStatus; + s.connectionError = patch.state.connectionError; + s.myName = patch.state.myName; + s.myRole = patch.state.myRole; + s.brokerUrl = patch.state.brokerUrl; + s.players = patch.state.players; + }), + ); + // Rebuild derived state (revealedPaths, stats) by replaying reducers + replayReducers(); + break; + } - for (const line of lines) { - try { - const msg: StreamMessage = JSON.parse(line); - messages.push(msg); - senderSeq[msg.sender] = Math.max(senderSeq[msg.sender] ?? 0, msg.seq); + case "message": { + const msg = patch.message; + const existingIdx = state.messages.findIndex((m) => m.id === msg.id); + if (existingIdx !== -1) { + setState("messages", existingIdx, msg); + } else { + setState( + produce((s) => { + s.messages.push(msg); + s.senderSeq[msg.sender] = Math.max( + s.senderSeq[msg.sender] ?? 0, + msg.seq, + ); + }), + ); + } runReducer(msg); - } catch { - /* skip corrupt */ + break; + } + + case "connectionStatus": { + setState("connectionStatus", patch.status); + if (patch.error !== undefined) { + setState("connectionError", patch.error); + } + if (patch.status === "connected") { + setState("connected", true); + } else if (patch.status === "disconnected") { + setState("connected", false); + } + break; + } + + case "players": { + setState("players", patch.players); + break; + } + + case "sessionName": { + setState("sessionName", patch.name); + break; + } + + case "sessionList": { + setSessionList(patch.sessions); + break; } } +} - setState( - produce((s) => { - s.messages = messages; - s.senderSeq = senderSeq; - s.sessionId = sessionId; +// --------------------------------------------------------------------------- +// Subscribe to worker (called once per tab) +// --------------------------------------------------------------------------- + +let _subscribed = false; + +async function ensureSubscribed(): Promise { + if (_subscribed) return; + _subscribed = true; + + const api = getWorkerAPI(); + _unsubscribe = await api.subscribe( + Comlink.proxy((patch: WorkerPatch) => { + handlePatch(patch); }), ); } +// Start subscription eagerly +ensureSubscribed(); + // --------------------------------------------------------------------------- -// MQTT Connect +// Public API — same signatures as before // --------------------------------------------------------------------------- -/** - * Connect to the MQTT broker, subscribe to the session stream, session - * list, and session meta. Must be called after `hydrateFromServer`. - */ +export function setMyName(name: string): void { + setState("myName", name); + saveName(name); + syncUrlParam("player", name); +} + +export function setMyRole(role: "gm" | "player" | "observer"): void { + setState("myRole", role); + saveRole(role); +} + +export function setSessionId(id: string | null): void { + setState("sessionId", id); + if (id) { + saveSessionId(id); + syncUrlParam("session", id); + } else { + removeUrlParam("session"); + setState("sessionName", null); + } +} + export async function connectStream( sessionId: string, brokerUrl: string, ): Promise { - const { default: mqtt } = await import("mqtt"); - + const api = getWorkerAPI(); setState("connectionStatus", "connecting"); setState("connectionError", null); - return new Promise((resolve, reject) => { - const client = mqtt.connect(brokerUrl, { - clientId: `${state.myName}-${state.myRole}-${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 - saveBrokerUrl(brokerUrl); - saveSessionId(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 }); - // Presence tracking - client.subscribe(`ttrpg/${sessionId}/presence/+`, { qos: 1 }); - - // Publish own presence (retained) - const presenceData = JSON.stringify({ - name: state.myName, - role: state.myRole, - }); - client.publish( - `ttrpg/${sessionId}/presence/${state.myName}`, - presenceData, - { qos: 1, retain: true }, - ); - - resolve(); - }); - - client.on("error", (err) => { - console.error("[stream] mqtt error:", err); - setState("connectionStatus", "error"); - setState("connectionError", err.message); - reject(err); - }); - - 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); - // Refresh the sessionName if we're in a session now - const currentId = state.sessionId; - if (currentId && manifest.sessions[currentId]) { - setState("sessionName", manifest.sessions[currentId].name); - } - } catch (e) { - console.error("[stream] manifest parse err:", e); - } - 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; - } - - // Presence: ttrpg/{sessionId}/presence/{playerName} - if ( - parts.length >= 4 && - parts[0] === "ttrpg" && - parts[2] === "presence" - ) { - const playerName = parts[3]; - if (raw) { - try { - const presence = JSON.parse(raw); - setState("players", playerName, { - role: presence.role || "player", - }); - } catch { - setState("players", playerName, { role: "player" }); - } - } else { - // Tombstone — player disconnected - setState( - produce((s) => { - delete s.players[playerName]; - }), - ); - } - 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); - } - } - }); - }); + try { + await api.connect(sessionId, brokerUrl, state.myName, state.myRole); + saveBrokerUrl(brokerUrl); + saveSessionId(sessionId); + } catch (err) { + setState("connectionStatus", "error"); + setState( + "connectionError", + err instanceof Error ? err.message : "Connection failed", + ); + throw err; + } } -// --------------------------------------------------------------------------- -// Session lifecycle -// --------------------------------------------------------------------------- - -/** - * Create a new session by publishing retained metadata to its meta topic. - * The server picks it up and adds it to the $SESSIONS manifest. - */ -export function createSession(name: string, players: string[] = []): string | null { - if (!_mqttClient || !_mqttConnected) return null; - - const id = - name - .toLowerCase() - .replace(/[^a-z0-9]+/g, "-") - .replace(/^-|-$/g, "") || "session"; - - const meta: SessionMeta = { name, created: Date.now(), players }; - - _mqttClient.publish(`ttrpg/${id}/meta`, JSON.stringify(meta), { - qos: 1, - retain: true, - }); - - return id; -} - -/** - * Delete a session: publish tombstone (empty payload) to its meta topic. - */ -export function deleteSession(sessionId: string): void { - if (!_mqttClient || !_mqttConnected) return; - - _mqttClient.publish(`ttrpg/${sessionId}/meta`, "", { qos: 1, retain: true }); -} - -// --------------------------------------------------------------------------- -// Send -// --------------------------------------------------------------------------- - -/** - * Publish a message to the stream. Validates the payload against the - * registered Zod schema before sending. Auto-increments the sender's seq. - */ export function sendMessage( type: string, payload: T, ): - { success: true; msg: StreamMessage } | { success: false; error: string } { - if (!_mqttClient || !_mqttConnected) { - return { success: false, error: "Not connected to stream" }; - } - - const sessionId = state.sessionId; - if (!sessionId) { - return { success: false, error: "No active session" }; - } + | { success: true; msg: StreamMessage } + | { success: false; error: string } { + const api = getWorkerAPI(); + // Validate locally first const validation = validatePayload(type, payload); if (!validation.success) { return { success: false, error: validation.error }; } + // Send to worker (which publishes to MQTT) + // The worker will broadcast back via the message patch + const result = api.sendMessage(type, validation.data); + + // We can't await the proxy result synchronously, so we optimistically + // return success. The actual message will arrive via the patch callback. + // For the sync API compatibility, we construct a placeholder. const sender = state.myName; const seq = (state.senderSeq[sender] ?? 0) + 1; - const id = makeMessageId(sender, seq); + const id = `${sender}-${seq}`; const msg: StreamMessage = { id, @@ -439,116 +313,52 @@ export function sendMessage( reverted: false, }; - const topic = `ttrpg/${sessionId}/stream`; - _mqttClient.publish(topic, JSON.stringify(msg), { qos: 1 }, (err) => { - if (err) console.error("[stream] publish error:", err); - }); - - // Optimistic local insert - receiveMessage(msg as StreamMessage); - return { success: true, msg }; } -// --------------------------------------------------------------------------- -// Receive -// --------------------------------------------------------------------------- - -function receiveMessage(msg: StreamMessage): void { - const existing = state.messages.find((m) => m.id === msg.id); - if (existing) { - setState( - produce((s) => { - const idx = s.messages.findIndex((m) => m.id === msg.id); - if (idx !== -1) s.messages[idx] = msg; - }), - ); - return; - } - - setState( - produce((s) => { - s.messages.push(msg); - s.senderSeq[msg.sender] = Math.max(s.senderSeq[msg.sender] ?? 0, msg.seq); - }), - ); - - runReducer(msg); -} - -// --------------------------------------------------------------------------- -// Revert -// --------------------------------------------------------------------------- - -/** - * Revert the current sender's latest (highest seq) message. - * Only works if it's still the latest — a subsequent message locks it. - */ export function revertLatest(): - { success: true } | { success: false; error: string } { - if (!_mqttClient || !_mqttConnected) { - return { success: false, error: "Not connected to stream" }; - } - - const sessionId = state.sessionId; - if (!sessionId) return { success: false, error: "No active session" }; - - const sender = state.myName; - const latestSeq = state.senderSeq[sender]; - if (!latestSeq) return { success: false, error: "No messages to revert" }; - - const id = makeMessageId(sender, latestSeq); - const original = state.messages.find((m) => m.id === id); - if (!original) return { success: false, error: "Message not found" }; - if (original.reverted) return { success: false, error: "Already reverted" }; - - const reverted: StreamMessage = { ...original, reverted: true }; - const topic = `ttrpg/${sessionId}/stream`; - _mqttClient.publish(topic, JSON.stringify(reverted), { qos: 1 }, (err) => { - if (err) console.error("[stream] revert publish error:", err); - }); - - return { success: true }; + | { success: true } + | { success: false; error: string } { + const api = getWorkerAPI(); + return api.revertLatest() as unknown as + | { success: true } + | { success: false; error: string }; } -// --------------------------------------------------------------------------- -// Disconnect -// --------------------------------------------------------------------------- - export function disconnectStream(): void { - if (_mqttClient) { - // Clear our presence before disconnecting (tombstone) - const sessionId = state.sessionId; - if (sessionId) { - _mqttClient.publish(`ttrpg/${sessionId}/presence/${state.myName}`, "", { - qos: 1, - retain: true, - }); - } - // Force disconnect without reconnect - _mqttClient.end(true, void 0, () => { - // noop - }); - _mqttClient = null; - _mqttConnected = false; - } - setState("connected", false); - setState("connectionStatus", "disconnected"); - setState("players", {}); - - // Strip autojoin param so the dialog shows normally on next connect + const api = getWorkerAPI(); + api.disconnect(); removeUrlParam("autojoin"); } -// --------------------------------------------------------------------------- -// Derived / helpers -// --------------------------------------------------------------------------- +export function createSession( + name: string, + players: string[] = [], +): string | null { + const api = getWorkerAPI(); + // This is sync in the worker but async over Comlink. + // We generate the ID locally using the same algorithm. + const id = + name + .toLowerCase() + .replace(/[^a-z0-9]+/g, "-") + .replace(/^-|-$/g, "") || "session"; + api.createSession(name, players); + return id; +} + +export function deleteSession(sessionId: string): void { + const api = getWorkerAPI(); + api.deleteSession(sessionId); +} export function canRevert(): boolean { const sender = state.myName; const seq = state.senderSeq[sender]; if (!seq) return false; - const msg = state.messages.find((m) => m.id === makeMessageId(sender, seq)); + const msg = state.messages.find( + (m) => m.id === `${sender}-${seq}`, + ); return msg !== undefined && !msg.reverted; } @@ -561,3 +371,11 @@ export function useJournalStream() { } export { state as journalStreamState }; + +/** + * Hydrate from server — now handled by the worker during connect(). + * Kept for API compatibility; does nothing. + */ +export async function hydrateFromServer(_sessionId: string): Promise { + // Hydration is now handled by the worker during connect() +} diff --git a/src/workers/journal.worker.ts b/src/workers/journal.worker.ts new file mode 100644 index 0000000..d88787d --- /dev/null +++ b/src/workers/journal.worker.ts @@ -0,0 +1,523 @@ +/** + * Journal Shared Worker — owns the single MQTT connection and message log. + * + * All browser tabs share one worker instance. The worker: + * - Connects to MQTT (one connection total) + * - Hydrates message history from the server + * - Maintains the message log and presence tracking + * - Broadcasts new messages and state changes to all connected tabs + * + * Each tab runs reducers locally to build derived state (revealedPaths, stats). + */ + +import * as Comlink from "comlink"; +import type { StreamMessage } from "../components/journal/registry"; + +// --------------------------------------------------------------------------- +// Types +// --------------------------------------------------------------------------- + +export interface SessionMeta { + name: string; + created: number; + players: string[]; +} + +export interface SessionManifest { + sessions: Record; +} + +/** Lightweight state snapshot sent to tabs. */ +export interface WorkerState { + sessionId: string | null; + sessionName: string | null; + messages: StreamMessage[]; + senderSeq: Record; + connected: boolean; + connectionStatus: "disconnected" | "connecting" | "connected" | "error"; + connectionError: string | null; + myName: string; + myRole: "gm" | "player" | "observer"; + brokerUrl: string | null; + players: Record; +} + +/** Patch sent on incremental updates (new message, presence change, etc.). */ +export type WorkerPatch = + | { type: "fullState"; state: WorkerState } + | { type: "message"; message: StreamMessage } + | { type: "connectionStatus"; status: WorkerState["connectionStatus"]; error?: string } + | { type: "players"; players: Record } + | { type: "sessionName"; name: string | null } + | { type: "sessionList"; sessions: SessionManifest }; + +// --------------------------------------------------------------------------- +// State +// --------------------------------------------------------------------------- + +const $SESSIONS = "ttrpg/$SESSIONS"; + +let state: WorkerState = { + sessionId: null, + sessionName: null, + messages: [], + senderSeq: {}, + connected: false, + connectionStatus: "disconnected", + connectionError: null, + myName: "gm", + myRole: "gm", + brokerUrl: null, + players: {}, +}; + +let sessionList: SessionManifest = { sessions: {} }; + +// MQTT client reference +let _mqttClient: import("mqtt").MqttClient | null = null; +let _mqttConnected = false; + +// Subscribers: each tab registers a callback +type PatchCallback = (patch: WorkerPatch) => void; +const subscribers = new Set(); + +// --------------------------------------------------------------------------- +// Helpers +// --------------------------------------------------------------------------- + +function makeMessageId(sender: string, seq: number): string { + return `${sender}-${seq}`; +} + +function notifySubscribers(patch: WorkerPatch): void { + for (const cb of subscribers) { + try { + cb(patch); + } catch { + // subscriber may have disconnected + } + } +} + +function broadcastFullState(): void { + notifySubscribers({ type: "fullState", state: { ...state, messages: [...state.messages], senderSeq: { ...state.senderSeq }, players: { ...state.players } } }); +} + +// --------------------------------------------------------------------------- +// Hydration +// --------------------------------------------------------------------------- + +async function hydrateFromServer(sessionId: string): Promise { + const response = await fetch( + `/.ttrpg/sessions/${encodeURIComponent(sessionId)}/stream.jsonl`, + ); + + if (!response.ok) { + if (response.status === 404) return; + throw new Error(`Failed to load session: ${response.statusText}`); + } + + const text = await response.text(); + const lines = text.split("\n").filter((l) => l.trim()); + + const messages: StreamMessage[] = []; + const senderSeq: Record = {}; + + for (const line of lines) { + try { + const msg: StreamMessage = JSON.parse(line); + messages.push(msg); + senderSeq[msg.sender] = Math.max(senderSeq[msg.sender] ?? 0, msg.seq); + } catch { + /* skip corrupt */ + } + } + + state.messages = messages; + state.senderSeq = senderSeq; + state.sessionId = sessionId; +} + +// --------------------------------------------------------------------------- +// MQTT Connect +// --------------------------------------------------------------------------- + +async function connectStream( + sessionId: string, + brokerUrl: string, + name: string, + role: "gm" | "player" | "observer", +): Promise { + const { default: mqtt } = await import("mqtt"); + + state.connectionStatus = "connecting"; + state.connectionError = null; + state.myName = name; + state.myRole = role; + notifySubscribers({ type: "connectionStatus", status: "connecting" }); + + return new Promise((resolve, reject) => { + const client = mqtt.connect(brokerUrl, { + clientId: `${name}-${role}-${Date.now()}`, + protocol: brokerUrl.startsWith("wss") ? "wss" : "ws", + reconnectPeriod: 2000, + }); + + _mqttClient = client; + + client.on("connect", () => { + _mqttConnected = true; + state.connected = true; + state.connectionStatus = "connected"; + state.connectionError = null; + state.brokerUrl = brokerUrl; + + client.subscribe(`ttrpg/${sessionId}/stream`, { qos: 1 }, (err) => { + if (err) console.error("[worker] stream sub err:", err); + }); + client.subscribe($SESSIONS, { qos: 1 }, (err) => { + if (err) console.error("[worker] sessions sub err:", err); + }); + client.subscribe(`ttrpg/${sessionId}/meta`, { qos: 1 }); + client.subscribe(`ttrpg/${sessionId}/presence/+`, { qos: 1 }); + + // Publish own presence (retained) + const presenceData = JSON.stringify({ name, role }); + client.publish( + `ttrpg/${sessionId}/presence/${name}`, + presenceData, + { qos: 1, retain: true }, + ); + + notifySubscribers({ type: "connectionStatus", status: "connected" }); + resolve(); + }); + + client.on("error", (err) => { + console.error("[worker] mqtt error:", err); + state.connectionStatus = "error"; + state.connectionError = err.message; + notifySubscribers({ type: "connectionStatus", status: "error", error: err.message }); + reject(err); + }); + + client.on("close", () => { + _mqttConnected = false; + state.connected = false; + state.connectionStatus = "disconnected"; + notifySubscribers({ type: "connectionStatus", status: "disconnected" }); + }); + + client.on("message", (topic, payload) => { + const raw = payload.toString(); + const parts = topic.split("/"); + + if (topic === $SESSIONS) { + try { + const manifest: SessionManifest = JSON.parse(raw); + sessionList = manifest; + const currentId = state.sessionId; + if (currentId && manifest.sessions[currentId]) { + state.sessionName = manifest.sessions[currentId].name; + notifySubscribers({ type: "sessionName", name: state.sessionName }); + } + notifySubscribers({ type: "sessionList", sessions: manifest }); + } catch (e) { + console.error("[worker] manifest parse err:", e); + } + return; + } + + if (parts.length >= 3 && parts[0] === "ttrpg" && parts[2] === "meta") { + return; + } + + // Presence + if ( + parts.length >= 4 && + parts[0] === "ttrpg" && + parts[2] === "presence" + ) { + const playerName = parts[3]; + if (raw) { + try { + const presence = JSON.parse(raw); + state.players = { ...state.players, [playerName]: { role: presence.role || "player" } }; + } catch { + state.players = { ...state.players, [playerName]: { role: "player" } }; + } + } else { + const newPlayers = { ...state.players }; + delete newPlayers[playerName]; + state.players = newPlayers; + } + notifySubscribers({ type: "players", players: { ...state.players } }); + return; + } + + if (parts.length >= 3 && parts[0] === "ttrpg" && parts[2] === "stream") { + try { + const msg: StreamMessage = JSON.parse(raw); + receiveMessage(msg); + } catch (e) { + console.error("[worker] malformed message:", e); + } + } + }); + }); +} + +// --------------------------------------------------------------------------- +// Receive +// --------------------------------------------------------------------------- + +function receiveMessage(msg: StreamMessage): void { + const existingIdx = state.messages.findIndex((m) => m.id === msg.id); + if (existingIdx !== -1) { + state.messages = [ + ...state.messages.slice(0, existingIdx), + msg, + ...state.messages.slice(existingIdx + 1), + ]; + } else { + state.messages = [...state.messages, msg]; + state.senderSeq = { + ...state.senderSeq, + [msg.sender]: Math.max(state.senderSeq[msg.sender] ?? 0, msg.seq), + }; + } + + notifySubscribers({ type: "message", message: msg }); +} + +// --------------------------------------------------------------------------- +// Send +// --------------------------------------------------------------------------- + +function sendMessage( + type: string, + payload: unknown, +): { success: true; msg: StreamMessage } | { success: false; error: string } { + if (!_mqttClient || !_mqttConnected) { + return { success: false, error: "Not connected to stream" }; + } + + const sessionId = state.sessionId; + if (!sessionId) { + return { success: false, error: "No active session" }; + } + + // We can't import validatePayload here because it depends on the registry + // which is populated by importing types. We do basic validation. + if (!type || typeof type !== "string") { + return { success: false, error: "Invalid message type" }; + } + + const sender = state.myName; + const seq = (state.senderSeq[sender] ?? 0) + 1; + const id = makeMessageId(sender, seq); + + const msg: StreamMessage = { + id, + sender, + seq, + type, + payload, + timestamp: Date.now(), + reverted: false, + }; + + const topic = `ttrpg/${sessionId}/stream`; + _mqttClient.publish(topic, JSON.stringify(msg), { qos: 1 }, (err) => { + if (err) console.error("[worker] publish error:", err); + }); + + // Optimistic local insert + receiveMessage(msg); + + return { success: true, msg }; +} + +// --------------------------------------------------------------------------- +// Revert +// --------------------------------------------------------------------------- + +function revertLatest(): + { success: true } | { success: false; error: string } { + if (!_mqttClient || !_mqttConnected) { + return { success: false, error: "Not connected to stream" }; + } + + const sessionId = state.sessionId; + if (!sessionId) return { success: false, error: "No active session" }; + + const sender = state.myName; + const latestSeq = state.senderSeq[sender]; + if (!latestSeq) return { success: false, error: "No messages to revert" }; + + const id = makeMessageId(sender, latestSeq); + const original = state.messages.find((m) => m.id === id); + if (!original) return { success: false, error: "Message not found" }; + if (original.reverted) return { success: false, error: "Already reverted" }; + + const reverted: StreamMessage = { ...original, reverted: true }; + const topic = `ttrpg/${sessionId}/stream`; + _mqttClient.publish(topic, JSON.stringify(reverted), { qos: 1 }, (err) => { + if (err) console.error("[worker] revert publish error:", err); + }); + + return { success: true }; +} + +// --------------------------------------------------------------------------- +// Disconnect +// --------------------------------------------------------------------------- + +function disconnectStream(): void { + if (_mqttClient) { + const sessionId = state.sessionId; + if (sessionId) { + _mqttClient.publish( + `ttrpg/${sessionId}/presence/${state.myName}`, + "", + { qos: 1, retain: true }, + ); + } + _mqttClient.end(true, void 0, () => {}); + _mqttClient = null; + _mqttConnected = false; + } + state.connected = false; + state.connectionStatus = "disconnected"; + state.players = {}; + notifySubscribers({ type: "connectionStatus", status: "disconnected" }); + notifySubscribers({ type: "players", players: {} }); +} + +// --------------------------------------------------------------------------- +// Session lifecycle +// --------------------------------------------------------------------------- + +function createSession(name: string, players: string[] = []): string | null { + if (!_mqttClient || !_mqttConnected) return null; + + const id = + name + .toLowerCase() + .replace(/[^a-z0-9]+/g, "-") + .replace(/^-|-$/g, "") || "session"; + + const meta: SessionMeta = { name, created: Date.now(), players }; + + _mqttClient.publish(`ttrpg/${id}/meta`, JSON.stringify(meta), { + qos: 1, + retain: true, + }); + + return id; +} + +function deleteSession(sessionId: string): void { + if (!_mqttClient || !_mqttConnected) return; + _mqttClient.publish(`ttrpg/${sessionId}/meta`, "", { qos: 1, retain: true }); +} + +// --------------------------------------------------------------------------- +// Derived +// --------------------------------------------------------------------------- + +function canRevert(): boolean { + const sender = state.myName; + const seq = state.senderSeq[sender]; + if (!seq) return false; + const msg = state.messages.find((m) => m.id === makeMessageId(sender, seq)); + return msg !== undefined && !msg.reverted; +} + +function visibleMessages(): StreamMessage[] { + return state.messages.filter((m) => !m.reverted); +} + +// --------------------------------------------------------------------------- +// Comlink API — exposed to main thread +// --------------------------------------------------------------------------- + +const api = { + async connect( + sessionId: string, + brokerUrl: string, + name: string, + role: "gm" | "player" | "observer", + ): Promise { + await hydrateFromServer(sessionId); + await connectStream(sessionId, brokerUrl, name, role); + }, + + disconnect(): void { + disconnectStream(); + }, + + sendMessage(type: string, payload: unknown) { + return sendMessage(type, payload); + }, + + revertLatest() { + return revertLatest(); + }, + + createSession(name: string, players: string[] = []) { + return createSession(name, players); + }, + + deleteSession(sessionId: string) { + deleteSession(sessionId); + }, + + getState(): WorkerState { + return { + ...state, + messages: [...state.messages], + senderSeq: { ...state.senderSeq }, + players: { ...state.players }, + }; + }, + + getSessionList(): SessionManifest { + return sessionList; + }, + + canRevert(): boolean { + return canRevert(); + }, + + visibleMessages(): StreamMessage[] { + return visibleMessages(); + }, + + /** Register a callback for state patches. Sends full state immediately. */ + subscribe(callback: Comlink.ProxyOrClone): () => void { + const cb = callback as PatchCallback; + subscribers.add(cb); + // Send initial full state + try { + cb({ + type: "fullState", + state: { + ...state, + messages: [...state.messages], + senderSeq: { ...state.senderSeq }, + players: { ...state.players }, + }, + }); + } catch { + // ignore + } + return () => { + subscribers.delete(cb); + }; + }, +}; + +export type JournalWorkerAPI = typeof api; + +Comlink.expose(api);