/** * Journal Stream — Client Store * * 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 article.reveal reducer). * * 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. */ import { createStore, produce } from "solid-js/store"; import { createSignal } from "solid-js"; import type { StreamMessage } from "../journal/registry"; import { getMessageType, validatePayload } from "../journal/registry"; // --------------------------------------------------------------------------- // Types // --------------------------------------------------------------------------- export interface JournalStreamState { sessionId: string | null; /** Full message log, oldest-first */ messages: StreamMessage[]; /** Last sequence number per sender */ senderSeq: Record; /** * Paths revealed by article.reveal messages. * Populated during hydration and live receipt via the type's reducer. */ revealedPaths: Set; /** 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; /** Broker URL, set after connect */ brokerUrl: string | null; } export interface SessionMeta { name: string; created: number; players: string[]; } export interface SessionManifest { sessions: Record; } // --------------------------------------------------------------------------- // localStorage keys // --------------------------------------------------------------------------- const LS_PLAYER_NAME = "ttrpg.playerName"; const LS_BROKER_URL = "ttrpg.brokerUrl"; const LS_LAST_SESSION = "ttrpg.lastSessionId"; function loadPersisted(): { myName: string; brokerUrl: string | null; lastSessionId: string | null; } { if (typeof localStorage === "undefined") { return { myName: "gm", brokerUrl: null, lastSessionId: null }; } return { myName: localStorage.getItem(LS_PLAYER_NAME) || "gm", brokerUrl: localStorage.getItem(LS_BROKER_URL), lastSessionId: localStorage.getItem(LS_LAST_SESSION), }; } // --------------------------------------------------------------------------- // Store // --------------------------------------------------------------------------- const persisted = loadPersisted(); const [state, setState] = createStore({ sessionId: persisted.lastSessionId, messages: [], senderSeq: {}, revealedPaths: new Set(), connected: false, connectionStatus: "disconnected", connectionError: null, myName: persisted.myName, brokerUrl: persisted.brokerUrl, }); export { setState as journalSetState }; const [sessionList, setSessionList] = createSignal({ sessions: {}, }); export { sessionList as sessions }; /** * Change the current player's name. Persisted to localStorage so it * survives page reloads. */ export function setMyName(name: string): void { setState("myName", name); if (typeof localStorage !== "undefined") { localStorage.setItem(LS_PLAYER_NAME, name); } } // Will hold the MQTT client instance after connect() let _mqttClient: import("mqtt").MqttClient | null = null; let _mqttConnected = false; // --------------------------------------------------------------------------- // Helpers // --------------------------------------------------------------------------- function makeMessageId(sender: string, seq: number): string { return `${sender}-${seq}`; } function runReducer(msg: StreamMessage): void { const def = getMessageType(msg.type); if (def?.reducer) { def.reducer(msg.payload, msg); } } 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}`); } 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); runReducer(msg); } catch { /* skip corrupt */ } } setState( produce((s) => { s.messages = messages; s.senderSeq = senderSeq; s.sessionId = sessionId; }), ); } // --------------------------------------------------------------------------- // MQTT Connect // --------------------------------------------------------------------------- /** * Connect to the MQTT broker, subscribe to the session stream, session * list, and session meta. Must be called after `hydrateFromServer`. */ export async function connectStream( sessionId: string, brokerUrl: string, ): Promise { const { default: mqtt } = await import("mqtt"); 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); }); client.subscribe($SESSIONS, { qos: 1 }, (err) => { if (err) console.error("[stream] sessions sub err:", err); }); client.subscribe(`ttrpg/${sessionId}/meta`, { qos: 1 }); }); 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; } 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); } } }); 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); }); } // --------------------------------------------------------------------------- // 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[] = []): void { if (!_mqttClient || !_mqttConnected) return; 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, }); } /** * 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" }; } const validation = validatePayload(type, payload); if (!validation.success) { return { success: false, error: validation.error }; } const sender = state.myName; const seq = (state.senderSeq[sender] ?? 0) + 1; const id = makeMessageId(sender, seq); const msg: StreamMessage = { id, sender, seq, type, payload: validation.data, timestamp: Date.now(), 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 }; } // --------------------------------------------------------------------------- // Disconnect // --------------------------------------------------------------------------- export function disconnectStream(): void { if (_mqttClient) { _mqttClient.end(true); _mqttClient = null; _mqttConnected = false; } setState("connected", false); setState("connectionStatus", "disconnected"); } // --------------------------------------------------------------------------- // Derived / helpers // --------------------------------------------------------------------------- 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)); return msg !== undefined && !msg.reverted; } export function visibleMessages(): StreamMessage[] { return state.messages.filter((m) => !m.reverted); } export function useJournalStream() { return state; } export { state as journalStreamState };