442 lines
13 KiB
TypeScript
442 lines
13 KiB
TypeScript
|
|
/**
|
||
|
|
* 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;
|
||
|
|
/** 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,
|
||
|
|
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");
|
||
|
|
|
||
|
|
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("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);
|
||
|
|
});
|
||
|
|
|
||
|
|
client.on("error", (err) => {
|
||
|
|
console.error("[stream] mqtt error:", 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[] = []): 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);
|
||
|
|
}
|
||
|
|
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
// 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 };
|