grafana, tautulli, nextcloud, nodered: full nox modules (ADR 0044/0046)
grafana: status/datasources/dashboards/alerts tools, emits alert.firing. tautulli: activity/history/stats tools, emits watch.recorded. nextcloud: users/shares/apps/occ tools (occ via docker exec, shares over OCS), emits user.created/share.created. nodered: flows/nodes/deploy tools, emits flows.deployed inline from the deploy tool. All typecheck; manifests parse.
This commit is contained in:
@@ -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<string, string>;
|
||||
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<any> {
|
||||
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<GrafanaHealth> {
|
||||
const h = await this.get("/api/health");
|
||||
return { database: h.database ?? "unknown", version: h.version ?? "unknown", commit: h.commit ?? "unknown" };
|
||||
}
|
||||
|
||||
async listDatasources(): Promise<GrafanaDatasource[]> {
|
||||
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<GrafanaDashboard[]> {
|
||||
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<GrafanaAlert[]> {
|
||||
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,
|
||||
}));
|
||||
}
|
||||
}
|
||||
@@ -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<string>();
|
||||
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<void> {
|
||||
const now = new Set<string>();
|
||||
const byKey = new Map<string, GrafanaAlert>();
|
||||
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");
|
||||
}
|
||||
@@ -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"
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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"]
|
||||
}
|
||||
@@ -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<string, string>;
|
||||
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<string, unknown>; disabled?: Record<string, unknown> };
|
||||
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<NextcloudShare[]> {
|
||||
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",
|
||||
}));
|
||||
}
|
||||
}
|
||||
@@ -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<string>();
|
||||
let usersPrimed = false;
|
||||
async function pollUsers(client: NextcloudClient): Promise<void> {
|
||||
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<string>();
|
||||
let sharesPrimed = false;
|
||||
async function pollShares(client: NextcloudClient): Promise<void> {
|
||||
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>): 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");
|
||||
}
|
||||
@@ -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"
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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"]
|
||||
}
|
||||
@@ -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<string, string> = {}): Record<string, string> {
|
||||
return { Accept: "application/json", ...(this.token ? { Authorization: `Bearer ${this.token}` } : {}), ...extra };
|
||||
}
|
||||
|
||||
private async req(path: string, init: RequestInit = {}): Promise<any> {
|
||||
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<any[]> {
|
||||
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<NodeRedNodeModule[]> {
|
||||
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 };
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,12 @@
|
||||
{
|
||||
"module": "nodered",
|
||||
"version": "1",
|
||||
"emits": [
|
||||
"module.nodered.flows.deployed"
|
||||
],
|
||||
"own-secrets": {
|
||||
"broker": "/var/lib/nodered/broker"
|
||||
},
|
||||
"capabilities": [
|
||||
"container-runtime"
|
||||
],
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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"]
|
||||
}
|
||||
@@ -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=…&<params>, 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<Record<string, unknown>>;
|
||||
}
|
||||
|
||||
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<string, string> = {}): Promise<any> {
|
||||
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<TautulliWatch[]> {
|
||||
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<TautulliHomeStat[]> {
|
||||
const data = (await this.cmd("get_home_stats")) as any[];
|
||||
return (data ?? []).map((s) => ({ statId: s.stat_id, rows: s.rows ?? [] }));
|
||||
}
|
||||
}
|
||||
@@ -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<number>();
|
||||
let primed = false;
|
||||
|
||||
async function pollHistory(client: TautulliClient): Promise<void> {
|
||||
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<void> {
|
||||
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");
|
||||
}
|
||||
@@ -1,6 +1,12 @@
|
||||
{
|
||||
"module": "tautulli",
|
||||
"version": "1",
|
||||
"emits": [
|
||||
"module.tautulli.watch.recorded"
|
||||
],
|
||||
"own-secrets": {
|
||||
"broker": "/var/lib/tautulli/broker"
|
||||
},
|
||||
"capabilities": [
|
||||
"container-runtime"
|
||||
],
|
||||
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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