Files
mesh-catalog/modules/plex/index.ts
jschoubben 7b06a7a408 Event names are local now, in the manifests and in the code
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 `#`.
2026-09-27 14:42:28 +02:00

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");