242 lines
6.8 KiB
TypeScript
242 lines
6.8 KiB
TypeScript
/**
|
|
* Journal stream persistence — MQTT broker + JSONL append + session manifest
|
|
*
|
|
* Runs an aedes MQTT broker over WebSocket, attached to the existing
|
|
* HTTP server. Uses ws v8's built-in createWebSocketStream.
|
|
*
|
|
* Pattern from aedes docs:
|
|
* https://github.com/moscajs/aedes#mqtt-server-over-websocket
|
|
*
|
|
* Persistence uses aedes's internal subscribe/publish API, so the
|
|
* server process doesn't need to connect to itself as an MQTT client.
|
|
*/
|
|
|
|
import { join, resolve } from "path";
|
|
import {
|
|
appendFileSync,
|
|
existsSync,
|
|
mkdirSync,
|
|
readFileSync,
|
|
rmdirSync,
|
|
unlinkSync,
|
|
writeFileSync,
|
|
} from "fs";
|
|
import type { Server as HttpServer } from "http";
|
|
import type { AedesPublishPacket } from "aedes";
|
|
import { WebSocketServer, createWebSocketStream } from "ws";
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Types
|
|
// ---------------------------------------------------------------------------
|
|
|
|
interface SessionMeta {
|
|
name: string;
|
|
created: number;
|
|
players: string[];
|
|
}
|
|
|
|
interface SessionManifest {
|
|
sessions: Record<string, SessionMeta>;
|
|
}
|
|
|
|
export interface JournalServer {
|
|
broker: import("aedes").Aedes;
|
|
wsServer: import("ws").WebSocketServer;
|
|
close(): Promise<void>;
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Path helpers
|
|
// ---------------------------------------------------------------------------
|
|
|
|
function dataDir(contentDir: string): string {
|
|
const d = resolve(contentDir, ".ttrpg");
|
|
mkdirSync(d, { recursive: true });
|
|
return d;
|
|
}
|
|
|
|
function manifestPath(dataRoot: string): string {
|
|
return join(dataRoot, "manifest.json");
|
|
}
|
|
|
|
function sessionDir(dataRoot: string, id: string): string {
|
|
return join(dataRoot, "sessions", id);
|
|
}
|
|
|
|
function streamPath(dataRoot: string, id: string): string {
|
|
return join(sessionDir(dataRoot, id), "stream.jsonl");
|
|
}
|
|
|
|
function loadManifest(dataRoot: string): SessionManifest {
|
|
const p = manifestPath(dataRoot);
|
|
if (!existsSync(p)) return { sessions: {} };
|
|
try {
|
|
return JSON.parse(readFileSync(p, "utf-8"));
|
|
} catch {
|
|
return { sessions: {} };
|
|
}
|
|
}
|
|
|
|
function saveManifest(dataRoot: string, m: SessionManifest): void {
|
|
writeFileSync(manifestPath(dataRoot), JSON.stringify(m, null, 2), "utf-8");
|
|
}
|
|
|
|
function appendStream(dataRoot: string, id: string, line: string): void {
|
|
const sd = sessionDir(dataRoot, id);
|
|
mkdirSync(sd, { recursive: true });
|
|
appendFileSync(streamPath(dataRoot, id), line + "\n", "utf-8");
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Server factory
|
|
// ---------------------------------------------------------------------------
|
|
|
|
const $SESSIONS = "ttrpg/$SESSIONS";
|
|
|
|
export async function createJournalServer(
|
|
contentDir: string,
|
|
httpServer: HttpServer,
|
|
): Promise<JournalServer> {
|
|
const root = dataDir(contentDir);
|
|
console.log(`[journal] data dir: ${root}`);
|
|
|
|
// ---- MQTT Broker ----
|
|
const { Aedes: AedesFactory } = await import("aedes");
|
|
// aedes requires listen() to initialize persistence before handling connections.
|
|
const broker = await AedesFactory.createBroker();
|
|
|
|
// ---- WebSocket server (attached to the existing HTTP server) ----
|
|
// Pattern from aedes docs: pass { server } to WSS, use 'connection' event
|
|
const wss = new WebSocketServer({ server: httpServer });
|
|
|
|
wss.on("connection", (ws, req) => {
|
|
const stream = createWebSocketStream(ws);
|
|
broker.handle(stream, req);
|
|
});
|
|
|
|
broker.on("client", (client: import("aedes").Client) => {
|
|
console.log(`[journal] client ready: ${client.id}`);
|
|
});
|
|
broker.on("clientDisconnect", (client: import("aedes").Client) => {
|
|
console.log(`[journal] client disconnected: ${client.id}`);
|
|
});
|
|
|
|
// ---- Persistence (internal subscribe, no separate MQTT client) ----
|
|
broker.subscribe(
|
|
"ttrpg/+/stream",
|
|
(packet, cb) => {
|
|
const sessionId = extractSessionId(packet.topic);
|
|
if (sessionId) {
|
|
try {
|
|
appendStream(root, sessionId, packet.payload.toString());
|
|
} catch (e) {
|
|
console.error(`[journal] stream append err for ${sessionId}:`, e);
|
|
}
|
|
}
|
|
cb();
|
|
},
|
|
() => {},
|
|
);
|
|
|
|
broker.subscribe(
|
|
"ttrpg/+/meta",
|
|
(packet, cb) => {
|
|
const sessionId = extractSessionId(packet.topic);
|
|
if (sessionId) {
|
|
handleMetaChange(root, sessionId, broker, packet.payload.toString());
|
|
}
|
|
cb();
|
|
},
|
|
() => {},
|
|
);
|
|
|
|
// Publish the initial manifest on startup so clients see existing sessions
|
|
const initialManifest = loadManifest(root);
|
|
broker.publish(
|
|
{
|
|
topic: $SESSIONS,
|
|
payload: Buffer.from(JSON.stringify(initialManifest)),
|
|
qos: 1,
|
|
retain: true,
|
|
cmd: "publish" as const,
|
|
dup: false,
|
|
},
|
|
() => {},
|
|
);
|
|
|
|
console.log("[journal] persistence listener active");
|
|
console.log("[journal] broker available at ws://<host>:<port>");
|
|
|
|
return {
|
|
broker,
|
|
wsServer: wss,
|
|
async close() {
|
|
console.log("[journal] shutting down...");
|
|
await new Promise<void>((r) => wss.close(() => r()));
|
|
await new Promise<void>((r) => broker.close(() => r()));
|
|
},
|
|
};
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Helpers
|
|
// ---------------------------------------------------------------------------
|
|
|
|
function extractSessionId(topic: string): string | null {
|
|
const parts = topic.split("/");
|
|
return parts.length >= 3 ? parts[1] : null;
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Session lifecycle
|
|
// ---------------------------------------------------------------------------
|
|
|
|
function handleMetaChange(
|
|
root: string,
|
|
sessionId: string,
|
|
broker: import("aedes").Aedes,
|
|
rawPayload: string,
|
|
): void {
|
|
const manifest = loadManifest(root);
|
|
try {
|
|
const meta = JSON.parse(rawPayload) as SessionMeta | null;
|
|
|
|
if (!meta || !meta.name) {
|
|
if (manifest.sessions[sessionId]) {
|
|
delete manifest.sessions[sessionId];
|
|
const sp = streamPath(root, sessionId);
|
|
if (existsSync(sp)) unlinkSync(sp);
|
|
const sd = sessionDir(root, sessionId);
|
|
try {
|
|
rmdirSync(sd);
|
|
} catch {
|
|
/* */
|
|
}
|
|
console.log(`[journal] session deleted: ${sessionId}`);
|
|
}
|
|
} else {
|
|
manifest.sessions[sessionId] = {
|
|
name: meta.name,
|
|
created: meta.created || Date.now(),
|
|
players: meta.players || [],
|
|
};
|
|
console.log(`[journal] session updated: ${sessionId} (${meta.name})`);
|
|
}
|
|
|
|
saveManifest(root, manifest);
|
|
broker.publish(
|
|
{
|
|
topic: $SESSIONS,
|
|
payload: Buffer.from(JSON.stringify(manifest)),
|
|
qos: 1,
|
|
retain: true,
|
|
cmd: "publish" as const,
|
|
dup: false,
|
|
},
|
|
() => {},
|
|
);
|
|
} catch (e) {
|
|
console.error(`[journal] meta parse err for ${sessionId}:`, e);
|
|
}
|
|
}
|