ttrpg-tools/src/components/stores/journalStream.ts

388 lines
11 KiB
TypeScript

/**
* Journal Stream — Client Store (Worker Proxy)
*
* 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.
*
* 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,
saveRole,
saveSessionId,
saveBrokerUrl,
syncUrlParam,
removeUrlParam,
readUrlParams,
} from "./persistence";
// ---------------------------------------------------------------------------
// Types
// ---------------------------------------------------------------------------
export interface JournalStreamState {
sessionId: string | null;
sessionName: string | null;
messages: StreamMessage[];
senderSeq: Record<string, number>;
revealedPaths: Record<string, Set<string>>;
connected: boolean;
connectionStatus: "disconnected" | "connecting" | "connected" | "error";
connectionError: string | null;
myName: string;
myRole: "gm" | "player" | "observer";
brokerUrl: string | null;
players: Record<string, { role: string }>;
variables: Record<string, string>;
}
export interface SessionMeta {
name: string;
created: number;
players: string[];
}
export type { SessionManifest } from "../../workers/journal.worker";
// ---------------------------------------------------------------------------
// Worker connection
// ---------------------------------------------------------------------------
let _workerAPI: Comlink.Remote<JournalWorkerAPI> | 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;
}
// ---------------------------------------------------------------------------
// Store
// ---------------------------------------------------------------------------
const persisted = loadPersisted();
const urlParams = readUrlParams();
const initialName = urlParams.playerName ?? persisted.myName;
const initialSession = urlParams.sessionId ?? persisted.lastSessionId;
const [state, setState] = createStore<JournalStreamState>({
sessionId: initialSession,
sessionName: null,
messages: [],
senderSeq: {},
revealedPaths: {},
connected: false,
connectionStatus: "disconnected",
connectionError: null,
myName: initialName,
myRole: (persisted.myRole as "gm" | "player" | "observer") || "gm",
brokerUrl: persisted.brokerUrl,
players: {},
variables: {},
});
if (initialName && !urlParams.playerName) syncUrlParam("player", initialName);
if (initialSession && !urlParams.sessionId)
syncUrlParam("session", initialSession);
export { setState as journalSetState };
const [sessionList, setSessionList] = createSignal<SessionManifest>({
sessions: {},
});
export { sessionList as sessions };
// ---------------------------------------------------------------------------
// Reducer runner — runs locally in each tab
// ---------------------------------------------------------------------------
function runReducer(msg: StreamMessage): void {
const def = getMessageType(msg.type);
if (def?.reducer) {
def.reducer(msg.payload, msg);
}
}
/** Replay all messages through reducers to rebuild derived state. */
function replayReducers(): void {
// Reset derived state
setState("revealedPaths", {});
setState("variables", {});
for (const msg of state.messages) {
runReducer(msg);
}
}
// ---------------------------------------------------------------------------
// Worker patch handler
// ---------------------------------------------------------------------------
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;
}),
);
// Persist name/role so URL and localStorage stay in sync
if (patch.state.myName) {
saveName(patch.state.myName);
syncUrlParam("player", patch.state.myName);
}
saveRole(patch.state.myRole);
// Rebuild derived state (revealedPaths, variables) 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();
await api.subscribe(
Comlink.proxy((patch: WorkerPatch) => {
handlePatch(patch);
}),
);
}
// Start subscription eagerly
ensureSubscribed();
// ---------------------------------------------------------------------------
// Public API — same signatures as before
// ---------------------------------------------------------------------------
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<void> {
const api = getWorkerAPI();
setState("connectionStatus", "connecting");
setState("connectionError", null);
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;
}
}
export function sendMessage<T>(
type: string,
payload: T,
):
| { success: true; msg: StreamMessage<T> }
| { 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 = `${sender}-${seq}`;
const msg: StreamMessage<T> = {
id,
sender,
seq,
type,
payload: validation.data,
timestamp: Date.now(),
reverted: false,
};
return { success: true, msg };
}
export function revertLatest():
| { success: true }
| { success: false; error: string } {
const api = getWorkerAPI();
return api.revertLatest() as unknown as
| { success: true }
| { success: false; error: string };
}
export function disconnectStream(): void {
const api = getWorkerAPI();
api.disconnect();
removeUrlParam("autojoin");
}
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 === `${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 };
/**
* 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()
}