plex: full nox module — client, tools and events (ADR 0044/0046)
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
This commit is contained in:
@@ -0,0 +1,131 @@
|
||||
// The Plex API client — plex's own code, living in the module (novox/hq ADR 0044). Moved out of the
|
||||
// shared hal sdk, where a change to Plex's API rebuilt everything; here it rebuilds only plex. Both
|
||||
// this module's tools and its events entrypoint import it, and nothing outside plex does.
|
||||
|
||||
import { existsSync, readFileSync } from "node:fs";
|
||||
import { join } from "node:path";
|
||||
|
||||
export interface PlexLibrary {
|
||||
key: string;
|
||||
title: string;
|
||||
type: string;
|
||||
count?: number;
|
||||
}
|
||||
|
||||
export interface PlexSession {
|
||||
key: string;
|
||||
title: string;
|
||||
user: string;
|
||||
player: string;
|
||||
state: string;
|
||||
type: string;
|
||||
}
|
||||
|
||||
export interface PlexItem {
|
||||
title: string;
|
||||
type: string;
|
||||
year?: number;
|
||||
summary?: string;
|
||||
addedAt?: string;
|
||||
}
|
||||
|
||||
export class PlexClient {
|
||||
readonly baseUrl: string;
|
||||
|
||||
constructor(
|
||||
url: string,
|
||||
private readonly token: string,
|
||||
) {
|
||||
this.baseUrl = url.replace(/\/$/, "");
|
||||
}
|
||||
|
||||
/**
|
||||
* Build from the module's resolved environment. The token is read from MESH_PLEX_TOKEN, or
|
||||
* discovered from the server's own Preferences.xml under the data directory — the same file Plex
|
||||
* writes it to, so a running server needs nothing configured by hand.
|
||||
*/
|
||||
static fromEnv(env: NodeJS.ProcessEnv = process.env): PlexClient {
|
||||
const url = env.MESH_PLEX_URL ?? `http://127.0.0.1:${env.PLEX_PORT ?? "32400"}`;
|
||||
const dataDir = env.MESH_PLEX_DATA_DIR ?? "/var/lib/plex";
|
||||
const token = env.MESH_PLEX_TOKEN ?? PlexClient.detectToken(dataDir);
|
||||
if (!token) throw new Error("no Plex token — set MESH_PLEX_TOKEN or make the data dir readable");
|
||||
return new PlexClient(url, token);
|
||||
}
|
||||
|
||||
/** Discover the token from the server's Preferences.xml, falling back to null. */
|
||||
static detectToken(dataDir: string): string | null {
|
||||
const prefs = join(dataDir, "config", "Library", "Application Support", "Plex Media Server", "Preferences.xml");
|
||||
if (existsSync(prefs)) {
|
||||
const match = readFileSync(prefs, "utf8").match(/PlexOnlineToken="([^"]+)"/);
|
||||
if (match) return match[1];
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private async get(path: string): Promise<any> {
|
||||
const url = `${this.baseUrl}${path}`;
|
||||
const sep = url.includes("?") ? "&" : "?";
|
||||
const res = await fetch(`${url}${sep}X-Plex-Token=${this.token}`, { headers: { Accept: "application/json" } });
|
||||
if (!res.ok) throw new Error(`Plex API ${path}: ${res.status} ${await res.text()}`);
|
||||
return res.json();
|
||||
}
|
||||
|
||||
async getServerInfo(): Promise<{ name: string; version: string; platform: string }> {
|
||||
const mc = (await this.get("/")).MediaContainer;
|
||||
return { name: mc.friendlyName || mc.machineIdentifier, version: mc.version, platform: mc.platform };
|
||||
}
|
||||
|
||||
async getLibraries(): Promise<PlexLibrary[]> {
|
||||
const dirs = (await this.get("/library/sections")).MediaContainer?.Directory ?? [];
|
||||
return dirs.map((d: any) => ({ key: d.key, title: d.title, type: d.type, count: d.count }));
|
||||
}
|
||||
|
||||
async getSessions(): Promise<PlexSession[]> {
|
||||
const sessions = (await this.get("/status/sessions")).MediaContainer?.Metadata ?? [];
|
||||
return sessions.map((s: any) => ({
|
||||
key: s.sessionKey ?? s.ratingKey,
|
||||
title: s.title + (s.grandparentTitle ? ` (${s.grandparentTitle})` : ""),
|
||||
user: s.User?.title ?? "unknown",
|
||||
player: s.Player?.title ?? s.Player?.product ?? "unknown",
|
||||
state: s.Player?.state ?? "unknown",
|
||||
type: s.type,
|
||||
}));
|
||||
}
|
||||
|
||||
async search(query: string): Promise<PlexItem[]> {
|
||||
const hubs = (await this.get(`/hubs/search?query=${encodeURIComponent(query)}&limit=20`)).MediaContainer?.Hub ?? [];
|
||||
const results: PlexItem[] = [];
|
||||
for (const hub of hubs) {
|
||||
for (const m of hub.Metadata ?? []) {
|
||||
results.push({
|
||||
title: m.title + (m.grandparentTitle ? ` (${m.grandparentTitle})` : ""),
|
||||
type: m.type,
|
||||
year: m.year,
|
||||
summary: m.summary?.slice(0, 200),
|
||||
});
|
||||
}
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
||||
async getRecentlyAdded(limit = 20): Promise<PlexItem[]> {
|
||||
const items = (await this.get(`/library/recentlyAdded?X-Plex-Container-Size=${limit}`)).MediaContainer?.Metadata ?? [];
|
||||
return items.map((m: any) => ({
|
||||
title: m.title + (m.grandparentTitle ? ` (${m.grandparentTitle})` : ""),
|
||||
type: m.type,
|
||||
year: m.year,
|
||||
summary: m.summary?.slice(0, 200),
|
||||
addedAt: m.addedAt ? new Date(m.addedAt * 1000).toISOString() : undefined,
|
||||
}));
|
||||
}
|
||||
|
||||
/** Ask Plex to rescan a library section — how a "new media arrived" event becomes a visible item. */
|
||||
async refreshLibrary(key: string): Promise<void> {
|
||||
await this.get(`/library/sections/${key}/refresh`);
|
||||
}
|
||||
|
||||
/** Rescan every library, for when what arrived is not known to belong to one. */
|
||||
async refreshAll(): Promise<void> {
|
||||
for (const library of await this.getLibraries()) await this.refreshLibrary(library.key);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
// 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");
|
||||
@@ -4,6 +4,17 @@
|
||||
"capabilities": [
|
||||
"container-runtime"
|
||||
],
|
||||
"emits": [
|
||||
"module.plex.playback.started",
|
||||
"module.plex.playback.stopped",
|
||||
"module.plex.item.added"
|
||||
],
|
||||
"consumes": [
|
||||
"module.*.download.completed"
|
||||
],
|
||||
"own-secrets": {
|
||||
"broker": "/var/lib/plex/broker"
|
||||
},
|
||||
"listens": [
|
||||
{
|
||||
"port": 32400,
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"name": "@novox/module-plex",
|
||||
"version": "0.1.0",
|
||||
"description": "plex — media server. Its API client, tools and events live here (novox/hq ADR 0044).",
|
||||
"type": "module",
|
||||
"private": true,
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.0.0",
|
||||
"typescript": "^5.6.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
// plex's tools — moved here from the shared sdk (novox/hq ADR 0044), importing plex's own client.
|
||||
// They return structured data; the mesh serves them through the sdk's tool harness.
|
||||
|
||||
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { PlexClient } from "../client.js";
|
||||
|
||||
export function getPlexTools(plex: PlexClient): ToolDefinition[] {
|
||||
return [
|
||||
{
|
||||
name: "plex_status",
|
||||
description: "Plex server status: server info, libraries, active sessions, recently added.",
|
||||
input: {},
|
||||
run: async () => {
|
||||
const [server, libraries, sessions, recent] = await Promise.all([
|
||||
plex.getServerInfo(),
|
||||
plex.getLibraries(),
|
||||
plex.getSessions(),
|
||||
plex.getRecentlyAdded(10),
|
||||
]);
|
||||
return { server, libraries, sessions, recentlyAdded: recent };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "plex_search",
|
||||
description: "Search across all Plex libraries — movies, shows, episodes, music.",
|
||||
input: { query: { type: "string", description: "the search query" } },
|
||||
run: async (args) => ({ query: String(args.query), results: await plex.search(String(args.query)) }),
|
||||
},
|
||||
{
|
||||
name: "plex_sessions",
|
||||
description: "Active Plex playback sessions — who is watching what, and where.",
|
||||
input: {},
|
||||
run: async () => {
|
||||
const sessions = await plex.getSessions();
|
||||
return { count: sessions.length, sessions };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "plex_recently_added",
|
||||
description: "Recently added media in Plex.",
|
||||
input: { limit: { type: "number", description: "how many items (default 20)" } },
|
||||
run: async (args) => ({ items: await plex.getRecentlyAdded(args.limit ? Number(args.limit) : 20) }),
|
||||
},
|
||||
{
|
||||
name: "plex_refresh",
|
||||
description: "Ask Plex to rescan its libraries so new files on disk become visible items.",
|
||||
input: { library: { type: "string", description: "a library section key; omitted rescans all" } },
|
||||
run: async (args) => {
|
||||
if (args.library) {
|
||||
await plex.refreshLibrary(String(args.library));
|
||||
return { refreshed: String(args.library) };
|
||||
}
|
||||
await plex.refreshAll();
|
||||
return { refreshed: "all" };
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
// The tools exist only when a token can be found; without one, plex contributes none rather than
|
||||
// failing the whole runtime.
|
||||
registerModuleTools("plex", (env) => {
|
||||
try {
|
||||
return getPlexTools(PlexClient.fromEnv(env));
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
});
|
||||
@@ -0,0 +1,12 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ES2022",
|
||||
"module": "NodeNext",
|
||||
"moduleResolution": "NodeNext",
|
||||
"strict": true,
|
||||
"esModuleInterop": true,
|
||||
"skipLibCheck": true,
|
||||
"noEmit": true
|
||||
},
|
||||
"include": ["client.ts", "index.ts", "tools/index.ts"]
|
||||
}
|
||||
Reference in New Issue
Block a user