Every module named its events the way the old bus spelled a routing key — `module.<module>.<verb>`. Design 29 says a module names an event locally and the mesh works out where it lands, so all 37 were stale against a rule already decided. On the new bus that derives into a namespace belonging to a module called "module", so no cross-module subscription in the mesh matched anything: nothing failed, nothing reacted (novox/hq 04-ISSUES/127). 36 manifests converted, and 43 files of module code with them. The code mattered as much as the manifests: the runtime builds the subject from what `emit()` is handed, so a converted manifest with unconverted code would have had the permission and the subject disagree. Three things the new check found on the way: - `photos` emitted an event its manifest never declared, which the new bus refuses outright. Declared. - `showcase` waited for an event nothing emits, so its demo could never be triggered — only `showcase` may publish under its own name. It emits both halves now. - `distribution` declared an event named after a different module. It emits `image.pushed` under its own name. An event about a *role* belongs on the seat, where the name outlives whoever holds it, but the sdk has no way to publish on a seat yet, so that stays recorded rather than declared. The audit logger's "everything" pattern is `**` rather than the old bus's `#`.
69 lines
2.7 KiB
TypeScript
69 lines
2.7 KiB
TypeScript
// plex's events. The tool runtime imports this once the broker is bound, and it does two things:
|
|
// it watches the server and emits what happened, and it reacts to the mesh's media events.
|
|
//
|
|
// Emits (novox/hq ADR 0041/0042):
|
|
// module.plex.playback.started / .stopped — someone began or ended watching
|
|
// module.plex.item.added — a new item appeared in a library
|
|
// Consumes:
|
|
// module.*.download.completed — a downloader finished; rescan so the file shows up
|
|
//
|
|
// The polling is deliberately unhurried: Plex is a neighbour on the same node, and an event a few
|
|
// seconds late is an event, whereas hammering the server for immediacy nobody asked for is not.
|
|
|
|
import { emit, on } from "@novox/mesh-sdk/events";
|
|
import { PlexClient, type PlexSession } from "./client.js";
|
|
|
|
const plex = PlexClient.fromEnv();
|
|
|
|
// Playback, by diffing the set of active sessions. Primed silently on the first look so a server
|
|
// that was already streaming when this started does not announce it as freshly begun.
|
|
const active = new Map<string, PlexSession>();
|
|
let playbackPrimed = false;
|
|
async function pollSessions(): Promise<void> {
|
|
const sessions = await plex.getSessions();
|
|
const now = new Map(sessions.map((s) => [s.key, s]));
|
|
if (playbackPrimed) {
|
|
for (const [key, s] of now) {
|
|
if (!active.has(key)) await emit("playback.started", { title: s.title, user: s.user, player: s.player, kind: s.type });
|
|
}
|
|
for (const [key, s] of active) {
|
|
if (!now.has(key)) await emit("playback.stopped", { title: s.title, user: s.user, player: s.player });
|
|
}
|
|
}
|
|
active.clear();
|
|
for (const [key, s] of now) active.set(key, s);
|
|
playbackPrimed = true;
|
|
}
|
|
|
|
// New items, by diffing recently-added. Primed silently too, or a restart would re-announce the
|
|
// whole recent list as new.
|
|
const seen = new Set<string>();
|
|
let itemsPrimed = false;
|
|
async function pollRecent(): Promise<void> {
|
|
const items = await plex.getRecentlyAdded(20);
|
|
for (const item of items) {
|
|
const id = `${item.title}@${item.addedAt ?? ""}`;
|
|
if (!seen.has(id)) {
|
|
if (itemsPrimed) await emit("item.added", item);
|
|
seen.add(id);
|
|
}
|
|
}
|
|
itemsPrimed = true;
|
|
}
|
|
|
|
// A downloader finished somewhere on the mesh: rescan, so what it fetched becomes a visible item
|
|
// rather than a file Plex has not noticed. Idempotent — a rescan too many costs a little disk I/O.
|
|
await on("*.download.completed", async () => {
|
|
await plex.refreshAll();
|
|
});
|
|
|
|
const tick = (fn: () => Promise<void>, everyMs: number): void => {
|
|
const run = (): void => void fn().catch((err) => console.error(`[plex] ${err}`));
|
|
setInterval(run, everyMs);
|
|
run();
|
|
};
|
|
tick(pollSessions, 15_000);
|
|
tick(pollRecent, 60_000);
|
|
|
|
console.log("[plex] watching sessions and recently-added, reacting to downloads");
|