feat: implement journal stream via Shared Worker

Move MQTT connection and message log management to a Shared Worker
using Comlink. This allows multiple browser tabs to share a single
connection and synchronized state.
This commit is contained in:
hypercross 2026-07-12 10:08:14 +08:00
parent fc6e37a13d
commit 2187b7ed82
4 changed files with 742 additions and 393 deletions

7
package-lock.json generated
View File

@ -14,6 +14,7 @@
"@thisbeyond/solid-dnd": "^0.7.5", "@thisbeyond/solid-dnd": "^0.7.5",
"aedes": "^1.1.1", "aedes": "^1.1.1",
"chokidar": "^5.0.0", "chokidar": "^5.0.0",
"comlink": "^4.4.2",
"commander": "^14.0.3", "commander": "^14.0.3",
"csv-parse": "^6.1.0", "csv-parse": "^6.1.0",
"csv-stringify": "^6.7.0", "csv-stringify": "^6.7.0",
@ -4332,6 +4333,12 @@
"node": ">= 0.8" "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": { "node_modules/commander": {
"version": "14.0.3", "version": "14.0.3",
"resolved": "https://registry.npmjs.org/commander/-/commander-14.0.3.tgz", "resolved": "https://registry.npmjs.org/commander/-/commander-14.0.3.tgz",

View File

@ -47,6 +47,7 @@
"@thisbeyond/solid-dnd": "^0.7.5", "@thisbeyond/solid-dnd": "^0.7.5",
"aedes": "^1.1.1", "aedes": "^1.1.1",
"chokidar": "^5.0.0", "chokidar": "^5.0.0",
"comlink": "^4.4.2",
"commander": "^14.0.3", "commander": "^14.0.3",
"csv-parse": "^6.1.0", "csv-parse": "^6.1.0",
"csv-stringify": "^6.7.0", "csv-stringify": "^6.7.0",

View File

@ -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 * This is the main-thread side of the journal stream. It spawns a Shared
* connection, local message log, per-sender sequence tracking, and the * Worker that owns the MQTT connection. All state flows from the worker
* revealed-paths set (populated by the link reducer). * to this store via Comlink callbacks.
* *
* Session lifecycle (create/list/delete) is handled via MQTT retained * The public API is unchanged components still call `useJournalStream()`,
* topics ttrpg/$SESSIONS for the manifest and ttrpg/{id}/meta per session. * `sendMessage()`, etc. exactly as before.
*
* No persistence here that's the CLI server's job via JSONL append.
*/ */
import { createStore, produce } from "solid-js/store"; import { createStore, produce } from "solid-js/store";
import { createSignal } from "solid-js"; import { createSignal } from "solid-js";
import * as Comlink from "comlink";
import type { StreamMessage } from "../journal/registry"; import type { StreamMessage } from "../journal/registry";
import { getMessageType, validatePayload } 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 { import {
loadPersisted, loadPersisted,
saveName, saveName,
@ -32,38 +33,17 @@ import {
export interface JournalStreamState { export interface JournalStreamState {
sessionId: string | null; sessionId: string | null;
/** Human-readable name for the current session (from manifest) */
sessionName: string | null; sessionName: string | null;
/** Full message log, oldest-first */
messages: StreamMessage[]; messages: StreamMessage[];
/** Last sequence number per sender */
senderSeq: Record<string, number>; senderSeq: Record<string, number>;
/**
* 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<string, Set<string>>; revealedPaths: Record<string, Set<string>>;
/** MQTT connection status */
connected: boolean; connected: boolean;
/** Granular connection state for UI indicators */
connectionStatus: "disconnected" | "connecting" | "connected" | "error"; connectionStatus: "disconnected" | "connecting" | "connected" | "error";
/** Last connection error message, if any */
connectionError: string | null; connectionError: string | null;
/** This client's identity */
myName: string; myName: string;
/** Role: gm | player | observer. Immutable while connected. */
myRole: "gm" | "player" | "observer"; myRole: "gm" | "player" | "observer";
/** Broker URL, set after connect */
brokerUrl: string | null; brokerUrl: string | null;
/** Active player list (keyed by player name) */
players: Record<string, { role: string }>; players: Record<string, { role: string }>;
/**
* 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<string, string>; stats: Record<string, string>;
} }
@ -73,8 +53,24 @@ export interface SessionMeta {
players: string[]; players: string[];
} }
export interface SessionManifest { export type { SessionManifest } from "../../workers/journal.worker";
sessions: Record<string, SessionMeta>;
// ---------------------------------------------------------------------------
// Worker connection
// ---------------------------------------------------------------------------
let _workerAPI: Comlink.Remote<JournalWorkerAPI> | null = null;
let _unsubscribe: (() => void) | null = null;
function getWorkerAPI(): Comlink.Remote<JournalWorkerAPI> {
if (!_workerAPI) {
const worker = new SharedWorker(
new URL("../../workers/journal.worker.ts", import.meta.url),
);
_workerAPI = Comlink.wrap<JournalWorkerAPI>(worker.port);
worker.port.start();
}
return _workerAPI;
} }
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@ -84,7 +80,6 @@ export interface SessionManifest {
const persisted = loadPersisted(); const persisted = loadPersisted();
const urlParams = readUrlParams(); const urlParams = readUrlParams();
// URL params override localStorage if present
const initialName = urlParams.playerName ?? persisted.myName; const initialName = urlParams.playerName ?? persisted.myName;
const initialSession = urlParams.sessionId ?? persisted.lastSessionId; const initialSession = urlParams.sessionId ?? persisted.lastSessionId;
@ -104,7 +99,6 @@ const [state, setState] = createStore<JournalStreamState>({
stats: {}, stats: {},
}); });
// Sync initial URL params if they came from localStorage (not URL)
if (initialName && !urlParams.playerName) syncUrlParam("player", initialName); if (initialName && !urlParams.playerName) syncUrlParam("player", initialName);
if (initialSession && !urlParams.sessionId) if (initialSession && !urlParams.sessionId)
syncUrlParam("session", initialSession); syncUrlParam("session", initialSession);
@ -117,56 +111,10 @@ const [sessionList, setSessionList] = createSignal<SessionManifest>({
export { sessionList as sessions }; 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 { function runReducer(msg: StreamMessage): void {
const def = getMessageType(msg.type); const def = getMessageType(msg.type);
if (def?.reducer) { if (def?.reducer) {
@ -174,260 +122,186 @@ function runReducer(msg: StreamMessage): void {
} }
} }
const $SESSIONS = "ttrpg/$SESSIONS"; /** Replay all messages through reducers to rebuild derived state. */
function replayReducers(): void {
// --------------------------------------------------------------------------- // Reset derived state
// Hydration (initial load from server) setState("revealedPaths", {});
// --------------------------------------------------------------------------- setState("stats", {});
for (const msg of state.messages) {
/**
* 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<void> {
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<string, number> = {};
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); runReducer(msg);
} catch {
/* skip corrupt */
}
} }
}
// ---------------------------------------------------------------------------
// Worker patch handler
// ---------------------------------------------------------------------------
function handlePatch(patch: WorkerPatch): void {
switch (patch.type) {
case "fullState": {
setState( setState(
produce((s) => { produce((s) => {
s.messages = messages; s.sessionId = patch.state.sessionId;
s.senderSeq = senderSeq; s.sessionName = patch.state.sessionName;
s.sessionId = sessionId; 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;
}
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);
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;
}
}
}
// ---------------------------------------------------------------------------
// Subscribe to worker (called once per tab)
// ---------------------------------------------------------------------------
let _subscribed = false;
async function ensureSubscribed(): Promise<void> {
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
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
/** export function setMyName(name: string): void {
* Connect to the MQTT broker, subscribe to the session stream, session setState("myName", name);
* list, and session meta. Must be called after `hydrateFromServer`. 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( export async function connectStream(
sessionId: string, sessionId: string,
brokerUrl: string, brokerUrl: string,
): Promise<void> { ): Promise<void> {
const { default: mqtt } = await import("mqtt"); const api = getWorkerAPI();
setState("connectionStatus", "connecting"); setState("connectionStatus", "connecting");
setState("connectionError", null); setState("connectionError", null);
return new Promise<void>((resolve, reject) => { try {
const client = mqtt.connect(brokerUrl, { await api.connect(sessionId, brokerUrl, state.myName, state.myRole);
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); saveBrokerUrl(brokerUrl);
saveSessionId(sessionId); saveSessionId(sessionId);
} catch (err) {
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("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( setState(
produce((s) => { "connectionError",
delete s.players[playerName]; err instanceof Error ? err.message : "Connection failed",
}),
); );
throw err;
} }
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);
}
}
});
});
} }
// ---------------------------------------------------------------------------
// 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<T>( export function sendMessage<T>(
type: string, type: string,
payload: T, payload: T,
): ):
{ success: true; msg: StreamMessage<T> } | { success: false; error: string } { | { success: true; msg: StreamMessage<T> }
if (!_mqttClient || !_mqttConnected) { | { success: false; error: string } {
return { success: false, error: "Not connected to stream" }; const api = getWorkerAPI();
}
const sessionId = state.sessionId;
if (!sessionId) {
return { success: false, error: "No active session" };
}
// Validate locally first
const validation = validatePayload(type, payload); const validation = validatePayload(type, payload);
if (!validation.success) { if (!validation.success) {
return { success: false, error: validation.error }; 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 sender = state.myName;
const seq = (state.senderSeq[sender] ?? 0) + 1; const seq = (state.senderSeq[sender] ?? 0) + 1;
const id = makeMessageId(sender, seq); const id = `${sender}-${seq}`;
const msg: StreamMessage<T> = { const msg: StreamMessage<T> = {
id, id,
@ -439,116 +313,52 @@ export function sendMessage<T>(
reverted: false, 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 }; 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(): export function revertLatest():
{ success: true } | { success: false; error: string } { | { success: true }
if (!_mqttClient || !_mqttConnected) { | { success: false; error: string } {
return { success: false, error: "Not connected to stream" }; const api = getWorkerAPI();
} return api.revertLatest() as unknown as
| { success: true }
const sessionId = state.sessionId; | { success: false; error: string };
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 { export function disconnectStream(): void {
if (_mqttClient) { const api = getWorkerAPI();
// Clear our presence before disconnecting (tombstone) api.disconnect();
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
removeUrlParam("autojoin"); removeUrlParam("autojoin");
} }
// --------------------------------------------------------------------------- export function createSession(
// Derived / helpers 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 { export function canRevert(): boolean {
const sender = state.myName; const sender = state.myName;
const seq = state.senderSeq[sender]; const seq = state.senderSeq[sender];
if (!seq) return false; 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; return msg !== undefined && !msg.reverted;
} }
@ -561,3 +371,11 @@ export function useJournalStream() {
} }
export { state as journalStreamState }; 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<void> {
// Hydration is now handled by the worker during connect()
}

View File

@ -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<string, SessionMeta>;
}
/** Lightweight state snapshot sent to tabs. */
export interface WorkerState {
sessionId: string | null;
sessionName: string | null;
messages: StreamMessage[];
senderSeq: Record<string, number>;
connected: boolean;
connectionStatus: "disconnected" | "connecting" | "connected" | "error";
connectionError: string | null;
myName: string;
myRole: "gm" | "player" | "observer";
brokerUrl: string | null;
players: Record<string, { role: string }>;
}
/** 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<string, { role: string }> }
| { 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<PatchCallback>();
// ---------------------------------------------------------------------------
// 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<void> {
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<string, number> = {};
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<void> {
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<void>((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<void> {
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<PatchCallback>): () => 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);