Moves plex's API client and tools out of the shared hal sdk into the module, so a Plex API change rebuilds only plex. Tools: status, search, sessions, recently-added, refresh. And a real event design: it emits playback started/stopped and item.added by watching the server, and consumes module.*.download.completed to rescan so a downloader's fetch becomes a visible item. Typechecks against the sdk; manifest parses with its emits/ consumes and broker own-secret. Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF
69 lines
2.8 KiB
TypeScript
69 lines
2.8 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 0046/0047):
|
|
// 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("module.plex.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("module.plex.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("module.plex.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("module.*.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");
|