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

462 lines
14 KiB
TypeScript
Raw Normal View History

/**
* 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<string, number>;
/**
* Paths revealed by article.reveal messages.
* Populated during hydration and live receipt via the type's reducer.
*/
revealedPaths: Set<string>;
/** 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<string, SessionMeta>;
}
// ---------------------------------------------------------------------------
// 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<JournalStreamState>({
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<SessionManifest>({
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<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);
} 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<void> {
const { default: mqtt } = await import("mqtt");
setState("connectionStatus", "connecting");
setState("connectionError", null);
return new Promise<void>((resolve, reject) => {
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 });
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);
} 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);
}
}
});
});
}
// ---------------------------------------------------------------------------
// 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<T>(
type: string,
payload: T,
):
{ success: true; msg: StreamMessage<T> } | { 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<T> = {
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 };