diff --git a/modules/grafana/client.ts b/modules/grafana/client.ts new file mode 100644 index 0000000..41395ae --- /dev/null +++ b/modules/grafana/client.ts @@ -0,0 +1,110 @@ +// Grafana's API client — grafana's own code, living in the module (novox/hq ADR 0044). Ported from +// the shared hal sdk, where a change here rebuilt everything; here it rebuilds only grafana. Both +// this module's tools and its events entrypoint import it, and nothing outside grafana does. + +export interface GrafanaHealth { + database: string; + version: string; + commit: string; +} + +export interface GrafanaDatasource { + id: number; + uid: string; + name: string; + type: string; + url: string; + isDefault: boolean; + database?: string; +} + +export interface GrafanaDashboard { + uid: string; + title: string; + url: string; + tags: string[]; + folderTitle?: string; +} + +export interface GrafanaAlert { + /** The rule name (labels.alertname), the stable identity a firing alert is diffed on. */ + name: string; + /** Grafana unified-alerting state: "Normal" | "Pending" | "Alerting". */ + state: string; + labels: Record; + activeAt?: string; +} + +export class GrafanaClient { + readonly baseUrl: string; + private readonly authHeader: string; + + constructor(url: string, authHeader: string) { + this.baseUrl = url.replace(/\/$/, ""); + this.authHeader = authHeader; + } + + /** + * Build from the module's resolved environment. Auth is a service-account/API token + * (MESH_GRAFANA_TOKEN, sent as Bearer) when present, else HTTP basic with the admin password the + * module keeps as its own secret (MESH_GRAFANA_PASSWORD, user MESH_GRAFANA_USER, default admin). + * Throws when neither is configured — the module then contributes nothing rather than failing. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): GrafanaClient { + const url = env.MESH_GRAFANA_URL ?? `http://127.0.0.1:${env.GRAFANA_PORT ?? "3000"}`; + const token = env.MESH_GRAFANA_TOKEN; + if (token) return new GrafanaClient(url, `Bearer ${token}`); + const password = env.MESH_GRAFANA_PASSWORD; + if (password) { + const user = env.MESH_GRAFANA_USER ?? "admin"; + return new GrafanaClient(url, `Basic ${Buffer.from(`${user}:${password}`).toString("base64")}`); + } + throw new Error("no Grafana auth — set MESH_GRAFANA_TOKEN or MESH_GRAFANA_PASSWORD"); + } + + private async get(path: string): Promise { + const res = await fetch(`${this.baseUrl}${path}`, { + headers: { Authorization: this.authHeader, Accept: "application/json" }, + }); + if (!res.ok) throw new Error(`Grafana API ${path}: ${res.status} ${await res.text()}`); + return res.json(); + } + + async health(): Promise { + const h = await this.get("/api/health"); + return { database: h.database ?? "unknown", version: h.version ?? "unknown", commit: h.commit ?? "unknown" }; + } + + async listDatasources(): Promise { + const arr = (await this.get("/api/datasources")) as any[]; + return arr.map((d) => ({ + id: d.id, uid: d.uid, name: d.name, type: d.type, url: d.url, + isDefault: !!d.isDefault, database: d.database || undefined, + })); + } + + async listDashboards(query?: string): Promise { + const params = new URLSearchParams({ type: "dash-db" }); + if (query) params.set("query", query); + const arr = (await this.get(`/api/search?${params.toString()}`)) as any[]; + return arr.map((d) => ({ + uid: d.uid, title: d.title, url: d.url, tags: d.tags ?? [], folderTitle: d.folderTitle || undefined, + })); + } + + /** + * Active alert instances from unified alerting's Prometheus-compatible surface. Grafana without + * alerting configured answers this with an empty set (or a 404, surfaced by get) — callers treat + * "no alerts" and "no alerting" alike. + */ + async listAlerts(): Promise { + const data = (await this.get("/api/prometheus/grafana/api/v1/alerts")).data ?? {}; + const alerts = (data.alerts ?? []) as any[]; + return alerts.map((a) => ({ + name: a.labels?.alertname ?? "unknown", + state: a.state ?? "unknown", + labels: a.labels ?? {}, + activeAt: a.activeAt || undefined, + })); + } +} diff --git a/modules/grafana/index.ts b/modules/grafana/index.ts new file mode 100644 index 0000000..bb497c8 --- /dev/null +++ b/modules/grafana/index.ts @@ -0,0 +1,58 @@ +// grafana's events. The tool runtime imports this once the broker is bound. It watches unified +// alerting and announces when an alert instance starts firing. +// +// Emits (novox/hq ADR 0046/0047): +// module.grafana.alert.firing — an alert instance entered the Alerting state +// +// A Grafana with no alerting configured simply never has a firing alert, so this observes nothing +// and emits nothing — no error, no noise. + +import { emit } from "@novox/mesh-sdk/events"; +import { GrafanaClient, type GrafanaAlert } from "./client.js"; + +// Constructed lazily so an unconfigured node (no auth) loads this entrypoint without crashing the +// events host — it simply watches nothing. +let grafana: GrafanaClient | undefined; +try { + grafana = GrafanaClient.fromEnv(); +} catch (err) { + console.log(`[grafana] not configured, not watching alerts: ${err}`); +} + +// Firing alerts, by diffing the set currently in the Alerting state. Primed silently on the first +// look so alerts already firing when this started are not announced as freshly firing. +const firing = new Set(); +let primed = false; + +const alertKey = (a: GrafanaAlert): string => + `${a.name}:${Object.entries(a.labels).sort().map(([k, v]) => `${k}=${v}`).join(",")}`; + +async function pollAlerts(client: GrafanaClient): Promise { + const now = new Set(); + const byKey = new Map(); + for (const a of await client.listAlerts()) { + if (a.state.toLowerCase() !== "alerting") continue; + const key = alertKey(a); + now.add(key); + byKey.set(key, a); + } + if (primed) { + for (const key of now) { + if (!firing.has(key)) { + const a = byKey.get(key)!; + await emit("module.grafana.alert.firing", { name: a.name, labels: a.labels, activeAt: a.activeAt }); + } + } + } + firing.clear(); + for (const key of now) firing.add(key); + primed = true; +} + +if (grafana) { + const client = grafana; + const run = (): void => void pollAlerts(client).catch((err) => console.error(`[grafana] ${err}`)); + setInterval(run, 30_000); + run(); + console.log("[grafana] watching for firing alerts"); +} diff --git a/modules/grafana/module.json b/modules/grafana/module.json index 5dcbbd5..e9ae8fa 100644 --- a/modules/grafana/module.json +++ b/modules/grafana/module.json @@ -1,8 +1,12 @@ { "module": "grafana", "version": "1", + "emits": [ + "module.grafana.alert.firing" + ], "own-secrets": { - "admin": "/var/lib/grafana-module/admin.secret" + "admin": "/var/lib/grafana-module/admin.secret", + "broker": "/var/lib/grafana-module/broker" }, "capabilities": [ "container-runtime" diff --git a/modules/grafana/package.json b/modules/grafana/package.json new file mode 100644 index 0000000..49003a9 --- /dev/null +++ b/modules/grafana/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-grafana", + "version": "0.1.0", + "description": "grafana — monitoring dashboards. 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" + } +} diff --git a/modules/grafana/tools/index.ts b/modules/grafana/tools/index.ts new file mode 100644 index 0000000..55fc89b --- /dev/null +++ b/modules/grafana/tools/index.ts @@ -0,0 +1,55 @@ +// grafana's tools — moved here from the shared hal sdk (novox/hq ADR 0044), importing grafana'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 { GrafanaClient } from "../client.js"; + +export function getGrafanaTools(grafana: GrafanaClient): ToolDefinition[] { + return [ + { + name: "grafana_status", + description: "Grafana server health — database state, version, build commit.", + input: {}, + run: async () => grafana.health(), + }, + { + name: "grafana_list_datasources", + description: "List Grafana data sources — name, type, backing URL, which is default.", + input: {}, + run: async () => { + const datasources = await grafana.listDatasources(); + return { count: datasources.length, datasources }; + }, + }, + { + name: "grafana_list_dashboards", + description: "List Grafana dashboards, optionally filtered by a name query.", + input: { query: { type: "string", description: "filter dashboards by name (optional)" } }, + run: async (args) => { + const query = args.query ? String(args.query) : undefined; + const dashboards = await grafana.listDashboards(query); + return { count: dashboards.length, dashboards }; + }, + }, + { + name: "grafana_alerts", + description: "Active Grafana alert instances and their state (Alerting, Pending, Normal).", + input: {}, + run: async () => { + const alerts = await grafana.listAlerts(); + const firing = alerts.filter((a) => a.state.toLowerCase() === "alerting"); + return { count: alerts.length, firing: firing.length, alerts }; + }, + }, + ]; +} + +// The tools exist only when Grafana auth can be resolved; without it, grafana contributes none +// rather than failing the whole tool runtime. +registerModuleTools("grafana", (env) => { + try { + return getGrafanaTools(GrafanaClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/grafana/tsconfig.json b/modules/grafana/tsconfig.json new file mode 100644 index 0000000..3677859 --- /dev/null +++ b/modules/grafana/tsconfig.json @@ -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"] +} diff --git a/modules/nextcloud/client.ts b/modules/nextcloud/client.ts new file mode 100644 index 0000000..5418924 --- /dev/null +++ b/modules/nextcloud/client.ts @@ -0,0 +1,93 @@ +// Nextcloud's client — nextcloud's own code, living in the module (novox/hq ADR 0044). Both this +// module's tools and its events entrypoint import it, and nothing outside nextcloud does. +// +// Nextcloud is administered two ways, and this client speaks both: +// - occ, its admin CLI, is a PHP script inside the container runnable only as the web user. We +// reach it with `docker exec`, the same side channel an operator would use by hand — turned +// into something the mesh can call. Users and apps come from here. +// - the OCS Sharing API answers over HTTP with the admin credentials. Shares come from here, +// because occ has no version-stable "list every share" across the releases we run. + +import { execFileSync } from "node:child_process"; + +export interface NextcloudUser { + uid: string; + displayName: string; +} + +export interface NextcloudShare { + /** The OCS share id — the stable identity a new share is diffed on. */ + id: string; + path: string; + shareType: number; + shareWith?: string; + owner: string; +} + +export class NextcloudClient { + constructor( + private readonly container: string, + private readonly ocsUrl: string, + private readonly adminUser: string, + private readonly adminPassword: string, + ) {} + + /** + * Build from the module's resolved environment. occ needs only the container name (default + * "nextcloud"); the OCS surface needs the admin password the module keeps as its own secret + * (MESH_NEXTCLOUD_ADMIN_PASSWORD, user MESH_NEXTCLOUD_ADMIN_USER default admin, URL the local + * container). The admin password is treated as the "configured for mesh administration" signal: + * throws without it, and the module then contributes nothing rather than failing. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): NextcloudClient { + const container = env.MESH_NEXTCLOUD_CONTAINER ?? "nextcloud"; + const ocsUrl = env.MESH_NEXTCLOUD_URL ?? `http://127.0.0.1:${env.NEXTCLOUD_PORT ?? "80"}`; + const adminUser = env.MESH_NEXTCLOUD_ADMIN_USER ?? "admin"; + const adminPassword = env.MESH_NEXTCLOUD_ADMIN_PASSWORD; + if (!adminPassword) throw new Error("no Nextcloud admin password — set MESH_NEXTCLOUD_ADMIN_PASSWORD"); + return new NextcloudClient(container, ocsUrl.replace(/\/$/, ""), adminUser, adminPassword); + } + + /** Run occ inside the container as the web user, returning its stdout, throwing its own message. */ + occ(args: string[]): string { + try { + return execFileSync("docker", ["exec", "-u", "www-data", this.container, "php", "occ", ...args], { + encoding: "utf8", timeout: 60_000, + }).trim(); + } catch (err: any) { + const detail = String(err?.stderr ?? err?.stdout ?? err?.message ?? "").trim(); + throw new Error(detail || `occ produced no output — is the ${this.container} container running?`); + } + } + + listUsers(): NextcloudUser[] { + // user:list --output=json answers an object of uid → display name. + const raw = this.occ(["user:list", "--output=json"]); + const map = JSON.parse(raw || "{}") as Record; + return Object.entries(map).map(([uid, displayName]) => ({ uid, displayName })); + } + + listApps(): { enabled: string[]; disabled: string[] } { + const raw = this.occ(["app:list", "--output=json"]); + const parsed = JSON.parse(raw || "{}") as { enabled?: Record; disabled?: Record }; + return { enabled: Object.keys(parsed.enabled ?? {}), disabled: Object.keys(parsed.disabled ?? {}) }; + } + + /** List every share, over the OCS Sharing API with the admin credentials. */ + async listShares(): Promise { + const auth = Buffer.from(`${this.adminUser}:${this.adminPassword}`).toString("base64"); + const res = await fetch( + `${this.ocsUrl}/ocs/v2.php/apps/files_sharing/api/v1/shares?format=json`, + { headers: { Authorization: `Basic ${auth}`, "OCS-APIRequest": "true", Accept: "application/json" } }, + ); + if (!res.ok) throw new Error(`Nextcloud OCS shares: ${res.status} ${await res.text()}`); + const rows = ((await res.json())?.ocs?.data ?? []) as any[]; + return rows.map((s) => ({ + id: String(s.id), + path: s.path ?? "", + shareType: Number(s.share_type ?? -1), + shareWith: s.share_with || undefined, + owner: s.uid_owner ?? "unknown", + })); + } +} diff --git a/modules/nextcloud/index.ts b/modules/nextcloud/index.ts new file mode 100644 index 0000000..3fc26f1 --- /dev/null +++ b/modules/nextcloud/index.ts @@ -0,0 +1,57 @@ +// nextcloud's events. The tool runtime imports this once the broker is bound. It watches the user +// list and the share list and announces new arrivals. +// +// Emits (novox/hq ADR 0046/0047): +// module.nextcloud.user.created — a user account appeared (occ user:list) +// module.nextcloud.share.created — a share appeared (OCS shares) +// +// Both are diffed and primed silently on the first look, so a restart does not re-announce every +// existing user and share as freshly created. + +import { emit } from "@novox/mesh-sdk/events"; +import { NextcloudClient } from "./client.js"; + +// Constructed lazily so an unconfigured node (no admin password) loads this entrypoint without +// crashing the events host — it simply watches nothing. +let nextcloud: NextcloudClient | undefined; +try { + nextcloud = NextcloudClient.fromEnv(); +} catch (err) { + console.log(`[nextcloud] not configured, not watching: ${err}`); +} + +const knownUsers = new Set(); +let usersPrimed = false; +async function pollUsers(client: NextcloudClient): Promise { + const users = client.listUsers(); + for (const u of users) { + if (knownUsers.has(u.uid)) continue; + if (usersPrimed) await emit("module.nextcloud.user.created", { uid: u.uid, displayName: u.displayName }); + knownUsers.add(u.uid); + } + usersPrimed = true; +} + +const knownShares = new Set(); +let sharesPrimed = false; +async function pollShares(client: NextcloudClient): Promise { + const shares = await client.listShares(); + for (const s of shares) { + if (knownShares.has(s.id)) continue; + if (sharesPrimed) await emit("module.nextcloud.share.created", { id: s.id, path: s.path, shareType: s.shareType, shareWith: s.shareWith, owner: s.owner }); + knownShares.add(s.id); + } + sharesPrimed = true; +} + +if (nextcloud) { + const client = nextcloud; + const tick = (fn: (c: NextcloudClient) => Promise): void => { + const run = (): void => void fn(client).catch((err) => console.error(`[nextcloud] ${err}`)); + setInterval(run, 60_000); + run(); + }; + tick(pollUsers); + tick(pollShares); + console.log("[nextcloud] watching users and shares"); +} diff --git a/modules/nextcloud/module.json b/modules/nextcloud/module.json index 20faa8d..8d56fee 100644 --- a/modules/nextcloud/module.json +++ b/modules/nextcloud/module.json @@ -21,8 +21,13 @@ "postgres-database": "/var/lib/nextcloud-module/database.secret", "s3-bucket": "/var/lib/nextcloud-module/store.secret" }, + "emits": [ + "module.nextcloud.user.created", + "module.nextcloud.share.created" + ], "own-secrets": { - "admin": "/var/lib/nextcloud-module/admin.secret" + "admin": "/var/lib/nextcloud-module/admin.secret", + "broker": "/var/lib/nextcloud-module/broker" }, "capabilities": [ "container-runtime" diff --git a/modules/nextcloud/package.json b/modules/nextcloud/package.json new file mode 100644 index 0000000..841026f --- /dev/null +++ b/modules/nextcloud/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-nextcloud", + "version": "0.1.0", + "description": "nextcloud — file sync and share. Its 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" + } +} diff --git a/modules/nextcloud/tools/index.ts b/modules/nextcloud/tools/index.ts new file mode 100644 index 0000000..224e8b7 --- /dev/null +++ b/modules/nextcloud/tools/index.ts @@ -0,0 +1,57 @@ +// nextcloud's tools — importing nextcloud's own client (novox/hq ADR 0044). occ runs inside the +// container; shares come over OCS. They return structured data; the mesh serves them through the +// sdk's tool harness. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { NextcloudClient } from "../client.js"; + +export function getNextcloudTools(nextcloud: NextcloudClient): ToolDefinition[] { + return [ + { + name: "nextcloud_users", + description: "List Nextcloud user accounts — uid and display name — via occ.", + input: {}, + run: async () => { + const users = nextcloud.listUsers(); + return { count: users.length, users }; + }, + }, + { + name: "nextcloud_shares", + description: "List Nextcloud shares — path, type, who it is shared with — via the OCS API.", + input: {}, + run: async () => { + const shares = await nextcloud.listShares(); + return { count: shares.length, shares }; + }, + }, + { + name: "nextcloud_apps", + description: "List Nextcloud apps, split into enabled and disabled, via occ.", + input: {}, + run: async () => nextcloud.listApps(), + }, + { + name: "nextcloud_occ", + description: + "Run an arbitrary occ admin command, e.g. status, 'config:system:get trusted_domains', " + + "user:list. occ is Nextcloud's CLI inside the container, run as the web user.", + input: { args: { type: "array", description: 'occ arguments, e.g. ["config:system:get","trusted_domains"]' } }, + run: async (args) => { + const occArgs = (args.args ?? []) as unknown[]; + if (!Array.isArray(occArgs) || occArgs.length === 0) throw new Error('args must be a non-empty array, e.g. ["status"]'); + return { output: nextcloud.occ(occArgs.map(String)) || "(no output)" }; + }, + }, + ]; +} + +// The tools exist only when the admin password can be resolved; without it, nextcloud contributes +// none rather than failing the whole tool runtime. +registerModuleTools("nextcloud", (env) => { + try { + return getNextcloudTools(NextcloudClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/nextcloud/tsconfig.json b/modules/nextcloud/tsconfig.json new file mode 100644 index 0000000..3677859 --- /dev/null +++ b/modules/nextcloud/tsconfig.json @@ -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"] +} diff --git a/modules/nodered/client.ts b/modules/nodered/client.ts new file mode 100644 index 0000000..7b92f59 --- /dev/null +++ b/modules/nodered/client.ts @@ -0,0 +1,86 @@ +// Node-RED's admin-API client — nodered's own code, living in the module (novox/hq ADR 0044). Only +// this module's tools import it; nodered has nothing to poll, so there is no events entrypoint. +// +// Node-RED exposes a runtime admin API under its base URL: GET/POST /flows for the whole flow +// configuration, GET /nodes for installed node modules. A default install has no auth; when +// adminAuth is on, a bearer token (minted at /auth/token) is required. + +export interface NodeRedFlow { + /** The tab (flow) node id. */ + id: string; + label: string; + disabled: boolean; +} + +export interface NodeRedNodeModule { + name: string; + version: string; + types: string[]; +} + +export class NodeRedClient { + readonly baseUrl: string; + + constructor( + url: string, + private readonly token?: string, + ) { + this.baseUrl = url.replace(/\/$/, ""); + } + + /** + * Build from the module's resolved environment. MESH_NODERED_URL locates the admin API and is the + * "this node runs Node-RED" signal — throws when unset, and the module then contributes nothing + * rather than failing on every node. MESH_NODERED_TOKEN is the bearer token when adminAuth is on; + * a default install needs none. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): NodeRedClient { + const url = env.MESH_NODERED_URL; + if (!url) throw new Error("no Node-RED URL — set MESH_NODERED_URL"); + return new NodeRedClient(url, env.MESH_NODERED_TOKEN); + } + + private headers(extra: Record = {}): Record { + return { Accept: "application/json", ...(this.token ? { Authorization: `Bearer ${this.token}` } : {}), ...extra }; + } + + private async req(path: string, init: RequestInit = {}): Promise { + const res = await fetch(`${this.baseUrl}${path}`, init); + if (!res.ok) throw new Error(`Node-RED ${path}: ${res.status} ${await res.text()}`); + return res.json(); + } + + /** The full flow configuration — the flat array of every node across every tab. */ + async getConfig(): Promise { + const body = await this.req("/flows", { headers: this.headers() }); + // /flows answers a bare array by default, or { rev, flows } to a v2-aware client. + return Array.isArray(body) ? body : (body.flows ?? []); + } + + /** The tabs (flows), each a node of type "tab" in the configuration. */ + async listFlows(): Promise<{ flows: NodeRedFlow[]; nodeCount: number }> { + const config = await this.getConfig(); + const flows = config + .filter((n) => n.type === "tab") + .map((n) => ({ id: n.id, label: n.label ?? "(unnamed)", disabled: !!n.disabled })); + return { flows, nodeCount: config.length }; + } + + async listNodes(): Promise { + const modules = (await this.req("/nodes", { headers: this.headers() })) as any[]; + return modules.map((m) => ({ name: m.name, version: m.version, types: m.types ?? [] })); + } + + /** + * Replace the whole flow configuration and deploy. Returns the new revision. `type` maps to + * Node-RED's deployment types — "full" (default), "nodes", or "flows". + */ + async deployFlows(config: any[], type = "full"): Promise<{ rev?: string; nodeCount: number }> { + const body = await this.req("/flows", { + method: "POST", + headers: this.headers({ "Content-Type": "application/json", "Node-RED-Deployment-Type": type }), + body: JSON.stringify(config), + }); + return { rev: body?.rev, nodeCount: config.length }; + } +} diff --git a/modules/nodered/module.json b/modules/nodered/module.json index b5d90d9..1bca4dd 100644 --- a/modules/nodered/module.json +++ b/modules/nodered/module.json @@ -1,6 +1,12 @@ { "module": "nodered", "version": "1", + "emits": [ + "module.nodered.flows.deployed" + ], + "own-secrets": { + "broker": "/var/lib/nodered/broker" + }, "capabilities": [ "container-runtime" ], diff --git a/modules/nodered/package.json b/modules/nodered/package.json new file mode 100644 index 0000000..cfa5768 --- /dev/null +++ b/modules/nodered/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-nodered", + "version": "0.1.0", + "description": "nodered — flow-based automation. Its client and tools 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" + } +} diff --git a/modules/nodered/tools/index.ts b/modules/nodered/tools/index.ts new file mode 100644 index 0000000..bd6c636 --- /dev/null +++ b/modules/nodered/tools/index.ts @@ -0,0 +1,64 @@ +// nodered's tools — importing nodered's own admin-API client (novox/hq ADR 0044). They return +// structured data; the mesh serves them through the sdk's tool harness. +// +// The deploy tool is nodered's one event source (novox/hq ADR 0046/0047): a successful deploy +// emits module.nodered.flows.deployed. nodered has nothing to observe on a timer, so there is no +// separate events entrypoint — the emit rides the action that causes it. The emit is best-effort: +// if no broker is bound, the deploy still succeeds and the announcement is simply skipped. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { emit } from "@novox/mesh-sdk/events"; +import { NodeRedClient } from "../client.js"; + +export function getNodeRedTools(nodered: NodeRedClient): ToolDefinition[] { + return [ + { + name: "nodered_list_flows", + description: "List Node-RED flows (tabs) — id, label, disabled state — and the total node count.", + input: {}, + run: async () => nodered.listFlows(), + }, + { + name: "nodered_list_nodes", + description: "List the Node-RED node modules installed in the runtime and their versions.", + input: {}, + run: async () => { + const nodes = await nodered.listNodes(); + return { count: nodes.length, nodes }; + }, + }, + { + name: "nodered_deploy", + description: + "Replace the whole Node-RED flow configuration and deploy it. `flows` is the full node " + + "array (as GET /flows returns). Emits module.nodered.flows.deployed on success.", + input: { + flows: { type: "array", description: "the full flow configuration — every node across every tab" }, + type: { type: "string", description: "deployment type: full (default), nodes, or flows" }, + }, + run: async (args) => { + const flows = args.flows as unknown[]; + if (!Array.isArray(flows)) throw new Error("flows must be an array of Node-RED nodes"); + const type = args.type ? String(args.type) : "full"; + const result = await nodered.deployFlows(flows, type); + // Best-effort announcement — a deploy must not fail because the broker is unbound here. + try { + await emit("module.nodered.flows.deployed", { rev: result.rev, nodeCount: result.nodeCount, type }); + } catch (err) { + console.error(`[nodered] deployed but could not emit: ${err}`); + } + return result; + }, + }, + ]; +} + +// The tools exist only when a Node-RED URL is configured; without one, nodered contributes none +// rather than failing the whole tool runtime. +registerModuleTools("nodered", (env) => { + try { + return getNodeRedTools(NodeRedClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/nodered/tsconfig.json b/modules/nodered/tsconfig.json new file mode 100644 index 0000000..426d382 --- /dev/null +++ b/modules/nodered/tsconfig.json @@ -0,0 +1,12 @@ +{ + "compilerOptions": { + "target": "ES2022", + "module": "NodeNext", + "moduleResolution": "NodeNext", + "strict": true, + "esModuleInterop": true, + "skipLibCheck": true, + "noEmit": true + }, + "include": ["client.ts", "tools/index.ts"] +} diff --git a/modules/tautulli/client.ts b/modules/tautulli/client.ts new file mode 100644 index 0000000..6c308c6 --- /dev/null +++ b/modules/tautulli/client.ts @@ -0,0 +1,98 @@ +// Tautulli's API client — tautulli's own code, living in the module (novox/hq ADR 0044). Both this +// module's tools and its events entrypoint import it, and nothing outside tautulli does. +// +// Tautulli speaks one endpoint: GET /api/v2?apikey=…&cmd=…&, answering +// { response: { result: "success" | "error", message, data } }. This client unwraps that envelope +// and hands back only the data. + +export interface TautulliSession { + user: string; + title: string; + mediaType: string; + state: string; + progressPercent: number; + player: string; +} + +export interface TautulliWatch { + /** Tautulli's history row id — the stable identity a recorded watch is diffed on. */ + id: number; + user: string; + title: string; + mediaType: string; + /** "watched" | "watching" | ... — Tautulli's own watched_status label. */ + watchedStatus: string; + percentComplete: number; + /** Unix seconds the play started, as Tautulli reports it. */ + date?: number; +} + +export interface TautulliHomeStat { + statId: string; + rows: Array>; +} + +export class TautulliClient { + readonly baseUrl: string; + + constructor( + url: string, + private readonly apiKey: string, + ) { + this.baseUrl = url.replace(/\/$/, ""); + } + + /** + * Build from the module's resolved environment. The API key is read from MESH_TAUTULLI_APIKEY + * (Tautulli mints it in Settings → Web Interface); the base URL defaults to the local container. + * Throws when no key is configured — the module then contributes nothing rather than failing. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): TautulliClient { + const url = env.MESH_TAUTULLI_URL ?? `http://127.0.0.1:${env.TAUTULLI_PORT ?? "8181"}`; + const apiKey = env.MESH_TAUTULLI_APIKEY; + if (!apiKey) throw new Error("no Tautulli API key — set MESH_TAUTULLI_APIKEY"); + return new TautulliClient(url, apiKey); + } + + /** Call one Tautulli command and return its unwrapped data, throwing on a non-success result. */ + private async cmd(command: string, params: Record = {}): Promise { + const q = new URLSearchParams({ apikey: this.apiKey, cmd: command, ...params }); + const res = await fetch(`${this.baseUrl}/api/v2?${q.toString()}`); + if (!res.ok) throw new Error(`Tautulli ${command}: ${res.status} ${await res.text()}`); + const body = (await res.json()).response ?? {}; + if (body.result !== "success") throw new Error(`Tautulli ${command}: ${body.message ?? "error"}`); + return body.data; + } + + async getActivity(): Promise<{ streamCount: number; sessions: TautulliSession[] }> { + const data = await this.cmd("get_activity"); + const sessions = ((data?.sessions ?? []) as any[]).map((s) => ({ + user: s.friendly_name ?? s.user ?? "unknown", + title: s.full_title ?? s.title ?? "unknown", + mediaType: s.media_type ?? "unknown", + state: s.state ?? "unknown", + progressPercent: Number(s.progress_percent ?? 0), + player: s.player ?? "unknown", + })); + return { streamCount: Number(data?.stream_count ?? sessions.length), sessions }; + } + + async getHistory(length = 25): Promise { + const data = await this.cmd("get_history", { length: String(length), order_column: "date", order_dir: "desc" }); + return ((data?.data ?? []) as any[]).map((r) => ({ + id: Number(r.row_id ?? r.id ?? r.reference_id ?? 0), + user: r.friendly_name ?? r.user ?? "unknown", + title: r.full_title ?? r.title ?? "unknown", + mediaType: r.media_type ?? "unknown", + watchedStatus: String(r.watched_status ?? ""), + percentComplete: Number(r.percent_complete ?? 0), + date: r.date != null ? Number(r.date) : undefined, + })); + } + + /** The home-page statistics blocks — most-watched shows, most-active users, and so on. */ + async getHomeStats(): Promise { + const data = (await this.cmd("get_home_stats")) as any[]; + return (data ?? []).map((s) => ({ statId: s.stat_id, rows: s.rows ?? [] })); + } +} diff --git a/modules/tautulli/index.ts b/modules/tautulli/index.ts new file mode 100644 index 0000000..faa6c77 --- /dev/null +++ b/modules/tautulli/index.ts @@ -0,0 +1,48 @@ +// tautulli's events. The tool runtime imports this once the broker is bound. It watches Tautulli's +// history and announces each newly recorded watch. +// +// Emits (novox/hq ADR 0046/0047): +// module.tautulli.watch.recorded — a play appeared in Tautulli's history +// +// Diffed on the history row id and primed silently on the first look, so a restart does not +// re-announce the whole existing history as freshly watched. + +import { emit } from "@novox/mesh-sdk/events"; +import { TautulliClient, type TautulliWatch } from "./client.js"; + +// Constructed lazily so an unconfigured node (no API key) loads this entrypoint without crashing +// the events host — it simply watches nothing. +let tautulli: TautulliClient | undefined; +try { + tautulli = TautulliClient.fromEnv(); +} catch (err) { + console.log(`[tautulli] not configured, not watching history: ${err}`); +} + +const seen = new Set(); +let primed = false; + +async function pollHistory(client: TautulliClient): Promise { + const history = await client.getHistory(25); + for (const w of history) { + if (w.id === 0 || seen.has(w.id)) continue; + if (primed) await emitWatch(w); + seen.add(w.id); + } + primed = true; +} + +async function emitWatch(w: TautulliWatch): Promise { + await emit("module.tautulli.watch.recorded", { + title: w.title, user: w.user, mediaType: w.mediaType, + watchedStatus: w.watchedStatus, percentComplete: w.percentComplete, at: w.date, + }); +} + +if (tautulli) { + const client = tautulli; + const run = (): void => void pollHistory(client).catch((err) => console.error(`[tautulli] ${err}`)); + setInterval(run, 60_000); + run(); + console.log("[tautulli] watching watch history"); +} diff --git a/modules/tautulli/module.json b/modules/tautulli/module.json index da05783..d3d3483 100644 --- a/modules/tautulli/module.json +++ b/modules/tautulli/module.json @@ -1,6 +1,12 @@ { "module": "tautulli", "version": "1", + "emits": [ + "module.tautulli.watch.recorded" + ], + "own-secrets": { + "broker": "/var/lib/tautulli/broker" + }, "capabilities": [ "container-runtime" ], diff --git a/modules/tautulli/package.json b/modules/tautulli/package.json new file mode 100644 index 0000000..9cd4335 --- /dev/null +++ b/modules/tautulli/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-tautulli", + "version": "0.1.0", + "description": "tautulli — Plex watch statistics. 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" + } +} diff --git a/modules/tautulli/tools/index.ts b/modules/tautulli/tools/index.ts new file mode 100644 index 0000000..b5697d6 --- /dev/null +++ b/modules/tautulli/tools/index.ts @@ -0,0 +1,44 @@ +// tautulli's tools — importing tautulli's own client (novox/hq ADR 0044). They return structured +// data; the mesh serves them through the sdk's tool harness. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { TautulliClient } from "../client.js"; + +export function getTautulliTools(tautulli: TautulliClient): ToolDefinition[] { + return [ + { + name: "tautulli_activity", + description: "Current Plex activity as Tautulli sees it — who is streaming what, and progress.", + input: {}, + run: async () => tautulli.getActivity(), + }, + { + name: "tautulli_history", + description: "Recent Plex watch history — who watched what, and whether they finished.", + input: { length: { type: "number", description: "how many rows (default 25)" } }, + run: async (args) => { + const history = await tautulli.getHistory(args.length ? Number(args.length) : 25); + return { count: history.length, history }; + }, + }, + { + name: "tautulli_stats", + description: "Tautulli home statistics — most-watched media, most-active users, and platforms.", + input: {}, + run: async () => { + const stats = await tautulli.getHomeStats(); + return { count: stats.length, stats }; + }, + }, + ]; +} + +// The tools exist only when an API key can be resolved; without one, tautulli contributes none +// rather than failing the whole tool runtime. +registerModuleTools("tautulli", (env) => { + try { + return getTautulliTools(TautulliClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/tautulli/tsconfig.json b/modules/tautulli/tsconfig.json new file mode 100644 index 0000000..3677859 --- /dev/null +++ b/modules/tautulli/tsconfig.json @@ -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"] +}