228 lines
6.5 KiB
TypeScript
228 lines
6.5 KiB
TypeScript
|
|
/**
|
||
|
|
* Journal stream persistence — MQTT broker + JSONL append + session manifest
|
||
|
|
*
|
||
|
|
* Starts an aedes MQTT broker and connects to itself to:
|
||
|
|
* 1. Persist all stream messages to .ttrpg/sessions/{id}/stream.jsonl
|
||
|
|
* 2. Manage session lifecycle via retained meta topics
|
||
|
|
*
|
||
|
|
* No HTTP routes — the JSONL files are served by the existing static server.
|
||
|
|
*/
|
||
|
|
|
||
|
|
import { join, resolve } from "path";
|
||
|
|
import {
|
||
|
|
appendFileSync,
|
||
|
|
existsSync,
|
||
|
|
mkdirSync,
|
||
|
|
readFileSync,
|
||
|
|
rmdirSync,
|
||
|
|
unlinkSync,
|
||
|
|
writeFileSync,
|
||
|
|
} from "fs";
|
||
|
|
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
// Types
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
|
||
|
|
interface SessionMeta {
|
||
|
|
name: string;
|
||
|
|
created: number;
|
||
|
|
players: string[];
|
||
|
|
}
|
||
|
|
|
||
|
|
interface SessionManifest {
|
||
|
|
sessions: Record<string, SessionMeta>;
|
||
|
|
}
|
||
|
|
|
||
|
|
export interface JournalServer {
|
||
|
|
/** MQTT broker instance */
|
||
|
|
broker: import("aedes").Aedes;
|
||
|
|
/** TCP server the broker is attached to */
|
||
|
|
tcpServer: import("net").Server;
|
||
|
|
/** Persistence client (server subscribes to itself) */
|
||
|
|
persistenceClient: import("mqtt").MqttClient | null;
|
||
|
|
/** Clean shutdown */
|
||
|
|
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");
|
||
|
|
}
|
||
|
|
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
// Manifest
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
|
||
|
|
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");
|
||
|
|
}
|
||
|
|
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
// JSONL persistence
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
|
||
|
|
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
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
|
||
|
|
export async function createJournalServer(
|
||
|
|
contentDir: string,
|
||
|
|
mqttPort: number,
|
||
|
|
): Promise<JournalServer> {
|
||
|
|
const root = dataDir(contentDir);
|
||
|
|
console.log(`[journal] data dir: ${root}`);
|
||
|
|
|
||
|
|
// ---- MQTT Broker ----
|
||
|
|
const { Aedes: AedesFactory } = await import("aedes");
|
||
|
|
const { createServer } = await import("net");
|
||
|
|
|
||
|
|
const broker = new AedesFactory();
|
||
|
|
const tcpServer = createServer(broker.handle.bind(broker));
|
||
|
|
|
||
|
|
tcpServer.listen(mqttPort, "0.0.0.0", () => {
|
||
|
|
console.log(`[journal] MQTT broker on port ${mqttPort}`);
|
||
|
|
});
|
||
|
|
|
||
|
|
// ---- Persistence + Manifest client ----
|
||
|
|
let client: import("mqtt").MqttClient | null = null;
|
||
|
|
|
||
|
|
const { default: mqtt } = await import("mqtt");
|
||
|
|
client = mqtt.connect(`mqtt://127.0.0.1:${mqttPort}`, {
|
||
|
|
clientId: "journal-server",
|
||
|
|
});
|
||
|
|
|
||
|
|
client.on("connect", () => {
|
||
|
|
// 1. Persistence: append every stream message to JSONL
|
||
|
|
client!.subscribe("ttrpg/+/stream", { qos: 1 }, (err) => {
|
||
|
|
if (err) console.error("[journal] stream sub err:", err);
|
||
|
|
});
|
||
|
|
|
||
|
|
// 2. Session meta: manage manifest from retained meta topics
|
||
|
|
client!.subscribe("ttrpg/+/meta", { qos: 1 }, (err) => {
|
||
|
|
if (err) console.error("[journal] meta sub err:", err);
|
||
|
|
});
|
||
|
|
});
|
||
|
|
|
||
|
|
client.on("message", (topic, payload) => {
|
||
|
|
const parts = topic.split("/");
|
||
|
|
if (parts.length < 3) return;
|
||
|
|
const sessionId = parts[1];
|
||
|
|
const subtopic = parts[2];
|
||
|
|
|
||
|
|
if (subtopic === "stream") {
|
||
|
|
// Persist every stream message
|
||
|
|
try {
|
||
|
|
appendStream(root, sessionId, payload.toString());
|
||
|
|
} catch (e) {
|
||
|
|
console.error(`[journal] stream append err for ${sessionId}:`, e);
|
||
|
|
}
|
||
|
|
} else if (subtopic === "meta") {
|
||
|
|
// Session lifecycle
|
||
|
|
handleMetaChange(root, sessionId, client!, payload.toString());
|
||
|
|
}
|
||
|
|
});
|
||
|
|
|
||
|
|
console.log("[journal] persistence listener active");
|
||
|
|
|
||
|
|
return {
|
||
|
|
broker,
|
||
|
|
tcpServer,
|
||
|
|
persistenceClient: client,
|
||
|
|
async close() {
|
||
|
|
console.log("[journal] shutting down...");
|
||
|
|
client?.end(true);
|
||
|
|
await new Promise<void>((r) => broker.close(() => r()));
|
||
|
|
await new Promise<void>((r) => tcpServer.close(() => r()));
|
||
|
|
},
|
||
|
|
};
|
||
|
|
}
|
||
|
|
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
// Session lifecycle
|
||
|
|
// ---------------------------------------------------------------------------
|
||
|
|
|
||
|
|
const $SESSIONS = "ttrpg/$SESSIONS";
|
||
|
|
|
||
|
|
function handleMetaChange(
|
||
|
|
root: string,
|
||
|
|
sessionId: string,
|
||
|
|
client: import("mqtt").MqttClient,
|
||
|
|
rawPayload: string,
|
||
|
|
): void {
|
||
|
|
const manifest = loadManifest(root);
|
||
|
|
|
||
|
|
try {
|
||
|
|
const meta = JSON.parse(rawPayload) as SessionMeta | null;
|
||
|
|
|
||
|
|
if (!meta || !meta.name) {
|
||
|
|
// Tombstone: empty or invalid payload → delete session
|
||
|
|
if (manifest.sessions[sessionId]) {
|
||
|
|
delete manifest.sessions[sessionId];
|
||
|
|
|
||
|
|
// Wipe files
|
||
|
|
const sp = streamPath(root, sessionId);
|
||
|
|
if (existsSync(sp)) unlinkSync(sp);
|
||
|
|
const sd = sessionDir(root, sessionId);
|
||
|
|
try {
|
||
|
|
rmdirSync(sd);
|
||
|
|
} catch {
|
||
|
|
/* not empty, fine */
|
||
|
|
}
|
||
|
|
|
||
|
|
console.log(`[journal] session deleted: ${sessionId}`);
|
||
|
|
}
|
||
|
|
} else {
|
||
|
|
// Create or update
|
||
|
|
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);
|
||
|
|
|
||
|
|
// Republish full manifest as a retained message
|
||
|
|
client.publish($SESSIONS, JSON.stringify(manifest), {
|
||
|
|
qos: 1,
|
||
|
|
retain: true,
|
||
|
|
});
|
||
|
|
} catch (e) {
|
||
|
|
console.error(`[journal] meta parse err for ${sessionId}:`, e);
|
||
|
|
}
|
||
|
|
}
|