From c82b3ff7066db65791ae77502b9a2378cc0a3822 Mon Sep 17 00:00:00 2001 From: jochen Date: Sat, 5 Sep 2026 04:17:07 +0200 Subject: [PATCH] Convert four hal modules: lidarr, mongodb, mssql, mosquitto Mirrors the proven catalog patterns field-for-field: - lidarr -> the Servarr twin of radarr/sonarr (API v1, artist content); no provisioner (it is a consumer app). - mongodb -> postgres shape: mongodb-database provider, provisioner mints a per-consumer db+user (ADR 0053), client shells to mongosh (no npm driver, the psql convention). - mssql -> postgres shape: mssql-database provider, sqlcmd client. - mosquitto -> redis shape: mqtt-topic provider via the Dynamic Security plugin, deliberately avoiding hal's password_file (that file is nox issue 011 exactly); provisioner mints a per-consumer MQTT client+role. All four typecheck (strict, NodeNext) against the built @novox/mesh-sdk, and their service images are digest-pinned to resolved registry digests. The mesh-runtime- images keep the all-zeros placeholder the pipeline pins, as postgres/redis do, and must bundle each module's CLI (mongosh/sqlcmd/ mosquitto_ctrl) as mesh-runtime-postgres bundles psql. Not yet lab-verified: each module lists in-code what an integration test must prove (auth model, provisioner reconcile, mosquitto dynsec bootstrap ordering). Claude-Session: https://claude.ai/code/session_01LrgweAeERJYBg88c5cKDzF --- modules/lidarr/client.ts | 144 +++++++++++++++ modules/lidarr/index.ts | 55 ++++++ modules/lidarr/module.json | 87 +++++++++ modules/lidarr/package.json | 14 ++ modules/lidarr/tools/index.ts | 78 ++++++++ modules/lidarr/tsconfig.json | 12 ++ modules/mongodb/client.ts | 180 ++++++++++++++++++ modules/mongodb/index.ts | 25 +++ modules/mongodb/module.json | 117 ++++++++++++ modules/mongodb/package.json | 14 ++ modules/mongodb/provisioner/index.ts | 48 +++++ modules/mongodb/tools/index.ts | 51 ++++++ modules/mongodb/tsconfig.json | 12 ++ modules/mosquitto/client.ts | 184 +++++++++++++++++++ modules/mosquitto/index.ts | 25 +++ modules/mosquitto/module.json | 122 +++++++++++++ modules/mosquitto/package.json | 14 ++ modules/mosquitto/provisioner/index.ts | 46 +++++ modules/mosquitto/tools/index.ts | 54 ++++++ modules/mosquitto/tsconfig.json | 12 ++ modules/mssql/client.ts | 242 +++++++++++++++++++++++++ modules/mssql/index.ts | 25 +++ modules/mssql/module.json | 115 ++++++++++++ modules/mssql/package.json | 14 ++ modules/mssql/provisioner/index.ts | 47 +++++ modules/mssql/tools/index.ts | 44 +++++ modules/mssql/tsconfig.json | 12 ++ 27 files changed, 1793 insertions(+) create mode 100644 modules/lidarr/client.ts create mode 100644 modules/lidarr/index.ts create mode 100644 modules/lidarr/module.json create mode 100644 modules/lidarr/package.json create mode 100644 modules/lidarr/tools/index.ts create mode 100644 modules/lidarr/tsconfig.json create mode 100644 modules/mongodb/client.ts create mode 100644 modules/mongodb/index.ts create mode 100644 modules/mongodb/module.json create mode 100644 modules/mongodb/package.json create mode 100644 modules/mongodb/provisioner/index.ts create mode 100644 modules/mongodb/tools/index.ts create mode 100644 modules/mongodb/tsconfig.json create mode 100644 modules/mosquitto/client.ts create mode 100644 modules/mosquitto/index.ts create mode 100644 modules/mosquitto/module.json create mode 100644 modules/mosquitto/package.json create mode 100644 modules/mosquitto/provisioner/index.ts create mode 100644 modules/mosquitto/tools/index.ts create mode 100644 modules/mosquitto/tsconfig.json create mode 100644 modules/mssql/client.ts create mode 100644 modules/mssql/index.ts create mode 100644 modules/mssql/module.json create mode 100644 modules/mssql/package.json create mode 100644 modules/mssql/provisioner/index.ts create mode 100644 modules/mssql/tools/index.ts create mode 100644 modules/mssql/tsconfig.json diff --git a/modules/lidarr/client.ts b/modules/lidarr/client.ts new file mode 100644 index 0000000..fba07e1 --- /dev/null +++ b/modules/lidarr/client.ts @@ -0,0 +1,144 @@ +// The Lidarr API client — lidarr's own code, living in the module (novox/hq ADR 0044). Ported from +// the shared hal `arr` client, but self-contained: in nox each Servarr app owns its own copy, so a +// change to Lidarr's API rebuilds only lidarr and nothing else. Both this module's tools and its +// events entrypoint import it, and nothing outside lidarr does. + +import { existsSync, readFileSync } from "node:fs"; +import { join } from "node:path"; + +// Lidarr speaks the v1 API (Radarr/Sonarr are v3); its content is the "artist". +const API_VERSION = "v1"; +const CONTENT_ENDPOINT = "artist"; +const APP_NAME = "Lidarr"; + +export interface LidarrQueueItem { + /** The queue record id — stable while the item is in the queue, so events can diff on it. */ + id: number; + title: string; + status: string; + size: string; + sizeleft: string; + timeleft?: string; +} + +export interface LidarrCalendarItem { + title: string; + date: string; + overview?: string; +} + +export interface LidarrContentItem { + title: string; + status?: string; + monitored: boolean; +} + +export class LidarrClient { + readonly baseUrl: string; + + constructor( + url: string, + private readonly apiKey: string, + ) { + this.baseUrl = url.replace(/\/$/, ""); + } + + /** + * Build from the module's resolved environment. The URL defaults to the server on this node (the + * runtime shares its network), and the API key is read from MESH_LIDARR_API_KEY or, failing that, + * discovered from the server's own config.xml under MESH_LIDARR_CONFIG_DIR — the same file Lidarr + * writes it to, so a running server needs nothing configured by hand. Throws when no key can be + * found, so the tools/events simply do not load (the harness treats the throw as "exposes + * nothing"). + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): LidarrClient { + const url = env.MESH_LIDARR_URL ?? `http://127.0.0.1:${env.MESH_LIDARR_PORT ?? "8686"}`; + const configDir = env.MESH_LIDARR_CONFIG_DIR ?? "/config"; + const apiKey = env.MESH_LIDARR_API_KEY ?? LidarrClient.detectApiKey(configDir); + if (!apiKey) { + throw new Error("Lidarr not configured — set MESH_LIDARR_API_KEY or make the config dir readable"); + } + return new LidarrClient(url, apiKey); + } + + /** Discover the API key from the server's config.xml, falling back to null. Every Servarr app + * writes into config.xml at the root of its config directory. */ + static detectApiKey(configDir: string): string | null { + const config = join(configDir, "config.xml"); + if (existsSync(config)) { + const match = readFileSync(config, "utf8").match(/([^<]+)<\/ApiKey>/); + if (match) return match[1]; + } + return null; + } + + private async get(endpoint: string, params?: Record): Promise { + const url = new URL(`${this.baseUrl}/api/${API_VERSION}/${endpoint}`); + if (params) { + for (const [k, v] of Object.entries(params)) url.searchParams.set(k, v); + } + const res = await fetch(url.toString(), { headers: { "X-Api-Key": this.apiKey } }); + if (!res.ok) throw new Error(`${APP_NAME} API /${endpoint}: ${res.status} ${await res.text()}`); + return res.json(); + } + + async getStatus(): Promise<{ appName: string; version: string }> { + const data = (await this.get("system/status")) as { appName?: string; version?: string }; + return { appName: data.appName || APP_NAME, version: data.version ?? "unknown" }; + } + + async getContent(limit?: number): Promise { + const data = await this.get(CONTENT_ENDPOINT); + const items: any[] = Array.isArray(data) ? data : ((data as any)?.records ?? []); + const mapped = items.map((item) => ({ + // Lidarr's content is an artist; its display name is artistName, not title. + title: item.artistName ?? item.title ?? "Unknown", + status: item.status, + monitored: item.monitored ?? true, + })); + return limit ? mapped.slice(0, limit) : mapped; + } + + /** Library search is a filter over existing content, not an indexer lookup — same as hal's. */ + async searchContent(term: string): Promise { + const all = await this.getContent(); + const lower = term.toLowerCase(); + return all.filter((item) => item.title.toLowerCase().includes(lower)); + } + + async getQueue(): Promise<{ totalRecords: number; items: LidarrQueueItem[] }> { + const data = (await this.get("queue", { pageSize: "50" })) as { totalRecords?: number; records?: any[] }; + const records = data.records ?? []; + return { + totalRecords: data.totalRecords ?? records.length, + items: records.map((r) => ({ + id: r.id, + title: r.title ?? r.artist?.artistName ?? r.album?.title ?? "Unknown", + status: r.status ?? "unknown", + size: formatBytes(r.size ?? 0), + sizeleft: formatBytes(r.sizeleft ?? 0), + timeleft: r.timeleft, + })), + }; + } + + async getCalendar(days = 7): Promise { + const start = new Date().toISOString().split("T")[0]; + const end = new Date(Date.now() + days * 86400000).toISOString().split("T")[0]; + const data = await this.get("calendar", { start, end }); + const items: any[] = Array.isArray(data) ? data : []; + return items.map((item) => ({ + // A Lidarr calendar entry is an album release. + title: item.title ?? item.artist?.artistName ?? "Unknown", + date: item.releaseDate ?? "", + overview: item.overview?.slice(0, 150), + })); + } +} + +function formatBytes(bytes: number): string { + if (bytes === 0) return "0 B"; + const units = ["B", "KB", "MB", "GB", "TB"]; + const i = Math.floor(Math.log(bytes) / Math.log(1024)); + return `${(bytes / Math.pow(1024, i)).toFixed(1)} ${units[i]}`; +} diff --git a/modules/lidarr/index.ts b/modules/lidarr/index.ts new file mode 100644 index 0000000..53e63db --- /dev/null +++ b/modules/lidarr/index.ts @@ -0,0 +1,55 @@ +// lidarr's events. The tool runtime imports this once the broker is bound. It watches the download +// queue and turns its comings and goings into mesh events. +// +// Emits (novox/hq ADR 0046/0047): +// module.lidarr.album.grabbed — a release entered the queue (Lidarr grabbed it) +// module.lidarr.download.completed — a release left the queue, imported. The download.completed +// routing key matches what a media consumer subscribes to +// (module.*.download.completed) to rescan its library. +// Consumes: none. +// +// The queue is polled and diffed, primed silently on the first look (like plex's index.ts) so a +// restart mid-download does not re-announce everything already in flight as freshly grabbed. + +import { emit } from "@novox/mesh-sdk/events"; +import { LidarrClient, type LidarrQueueItem } from "./client.js"; + +const lidarr = LidarrClient.fromEnv(); + +// Lidarr removes an item from the queue once it has been imported; a "warning"/"failed" status is +// how a stuck or broken grab shows itself, so we do not call those a completion when they vanish. +const FAILED_STATUSES = new Set(["failed", "warning"]); + +const inQueue = new Map(); +let primed = false; + +async function pollQueue(): Promise { + const { items } = await lidarr.getQueue(); + const now = new Map(items.map((i) => [i.id, i])); + + if (primed) { + // Entered the queue since last look — Lidarr grabbed a release. + for (const [id, item] of now) { + if (!inQueue.has(id)) await emit("module.lidarr.album.grabbed", { title: item.title, status: item.status }); + } + // Left the queue — imported and done, unless it was last seen failing. + for (const [id, item] of inQueue) { + if (!now.has(id) && !FAILED_STATUSES.has(item.status)) { + await emit("module.lidarr.download.completed", { title: item.title }); + } + } + } + + inQueue.clear(); + for (const [id, item] of now) inQueue.set(id, item); + primed = true; +} + +const tick = (fn: () => Promise, everyMs: number): void => { + const run = (): void => void fn().catch((err) => console.error(`[lidarr] ${err}`)); + setInterval(run, everyMs); + run(); +}; +tick(pollQueue, 30_000); + +console.log("[lidarr] watching the download queue, emitting grabs and completions"); diff --git a/modules/lidarr/module.json b/modules/lidarr/module.json new file mode 100644 index 0000000..5172bba --- /dev/null +++ b/modules/lidarr/module.json @@ -0,0 +1,87 @@ +{ + "module": "lidarr", + "version": "1", + "capabilities": [ + "container-runtime" + ], + "emits": [ + "module.lidarr.album.grabbed", + "module.lidarr.download.completed" + ], + "consumes": [], + "own-secrets": { + "broker": "/var/lib/mesh/lidarr/broker" + }, + "listens": [ + { + "port": 8686, + "protocol": "tcp", + "from": "mesh", + "why": "managing music" + } + ], + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/lidarr", + "mode": "0700" + }, + { + "id": "config", + "type": "directory", + "path": "/services/lidarr/config", + "mode": "0700", + "owner": "1000:1000" + }, + { + "id": "media-music", + "type": "directory", + "path": "/services/media/music", + "mode": "0755", + "owner": "1000:1000" + }, + { + "id": "media-downloads", + "type": "directory", + "path": "/services/media/downloads", + "mode": "0755", + "owner": "1000:1000" + }, + { + "id": "server", + "type": "container", + "name": "lidarr", + "image": "lscr.io/linuxserver/lidarr@sha256:6b38dd330b0c653351c2e23c8b962ea51c95683dd7acace9d106c922baf85f75", + "env": { + "PUID": "1000", + "PGID": "1000", + "TZ": "Etc/UTC" + }, + "ports": [ + "8686" + ], + "volumes": [ + "/services/lidarr/config:/config", + "/services/media/music:/music", + "/services/media/downloads:/downloads" + ] + }, + { + "id": "runtime", + "type": "container", + "name": "mesh-lidarr", + "image": "mesh-runtime-lidarr@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "host", + "volumes": [ + "/var/lib/mesh/lidarr/broker:/run/secrets/broker:ro", + "/services/lidarr/config:/var/lib/lidarr/config:ro" + ], + "env": { + "MESH_BROKER_FILE": "/run/secrets/broker", + "MESH_LIDARR_URL": "http://127.0.0.1:8686", + "MESH_LIDARR_CONFIG_DIR": "/var/lib/lidarr/config" + } + } + ] +} diff --git a/modules/lidarr/package.json b/modules/lidarr/package.json new file mode 100644 index 0000000..c2a671d --- /dev/null +++ b/modules/lidarr/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-lidarr", + "version": "0.1.0", + "description": "lidarr — music management. 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/lidarr/tools/index.ts b/modules/lidarr/tools/index.ts new file mode 100644 index 0000000..85aee09 --- /dev/null +++ b/modules/lidarr/tools/index.ts @@ -0,0 +1,78 @@ +// lidarr's tools — ported from the shared hal sdk (novox/hq ADR 0044), importing lidarr's own +// client. They return structured data (not pre-formatted text as hal did); the mesh serves them +// through the sdk's tool harness. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { LidarrClient } from "../client.js"; + +export function getLidarrTools(lidarr: LidarrClient): ToolDefinition[] { + return [ + { + name: "lidarr_status", + description: "Lidarr status overview: version, artist count, monitored count, queue size.", + input: {}, + run: async () => { + const [status, content, queue] = await Promise.all([ + lidarr.getStatus(), + lidarr.getContent(), + lidarr.getQueue(), + ]); + return { + app: status.appName, + version: status.version, + artists: content.length, + monitored: content.filter((c) => c.monitored).length, + queue: queue.totalRecords, + }; + }, + }, + { + name: "lidarr_library", + description: "List artists from the Lidarr library.", + input: { limit: { type: "number", description: "max items to return (default 50)" } }, + run: async (args) => { + const items = await lidarr.getContent(args.limit ? Number(args.limit) : 50); + return { count: items.length, artists: items }; + }, + }, + { + name: "lidarr_search", + description: "Search the Lidarr library for artists by name (filters existing content, not indexers).", + input: { query: { type: "string", description: "the search term" } }, + run: async (args) => { + const query = String(args.query); + return { query, results: await lidarr.searchContent(query) }; + }, + }, + { + name: "lidarr_queue", + description: "Show the Lidarr download queue — what is downloading and how far along.", + input: {}, + run: async () => { + const queue = await lidarr.getQueue(); + return { count: queue.totalRecords, items: queue.items }; + }, + }, + { + name: "lidarr_calendar", + description: "Upcoming album releases from the Lidarr calendar.", + input: { days: { type: "number", description: "how many days to look ahead (default 7)" } }, + run: async (args) => { + const days = args.days ? Number(args.days) : 7; + const items = await lidarr.getCalendar(days); + items.sort((a, b) => a.date.localeCompare(b.date)); + return { days, count: items.length, items }; + }, + }, + ]; +} + +// The tools exist only when Lidarr is configured; without a URL and key, lidarr contributes none +// rather than failing the whole runtime. +registerModuleTools("lidarr", (env) => { + try { + return getLidarrTools(LidarrClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/lidarr/tsconfig.json b/modules/lidarr/tsconfig.json new file mode 100644 index 0000000..3677859 --- /dev/null +++ b/modules/lidarr/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/mongodb/client.ts b/modules/mongodb/client.ts new file mode 100644 index 0000000..18b5ffe --- /dev/null +++ b/modules/mongodb/client.ts @@ -0,0 +1,180 @@ +// mongodb's admin client — mongodb's own code, living in the module (novox/hq ADR 0044). Both this +// module's tools and its provisioner import it, and nothing outside mongodb does. +// +// Commands run through `mongosh`, not a wire-protocol driver: the module may take NO npm dependency +// beyond @novox/mesh-sdk, and hand-rolling the MongoDB wire protocol + SCRAM auth is more surface +// than this should carry — so it shells out to the shell the mongodb image ships, the same way +// postgres drives itself through `psql`, minio through `mc` and mailu through doveadm. One boundary, +// `evalJs()`, and every method is built on it: a snippet of JavaScript is evaluated server-side and +// its result comes back as EJSON on stdout. + +import { randomBytes } from "node:crypto"; +import { readFileSync } from "node:fs"; +import { execFile } from "node:child_process"; +import { promisify } from "node:util"; + +const run = promisify(execFile); + +export interface DatabaseInfo { + readonly name: string; + readonly sizeBytes: number; +} + +export interface MongoConn { + readonly host: string; + readonly port: number; + readonly user: string; + readonly password: string; + /** The database the admin user authenticates against — `admin` for the root user. */ + readonly authSource: string; +} + +export class MongoClient { + constructor(private readonly conn: MongoConn) {} + + /** + * Build from the module's resolved environment. Reads MESH_MONGODB_* first (the documented names), + * falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner container. + * Throws if it cannot find a host and an admin password. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): MongoClient { + const url = env.MESH_PROVISION_MONGODB ? safeUrl(env.MESH_PROVISION_MONGODB) : undefined; + const host = env.MESH_MONGODB_HOST ?? url?.hostname; + const port = Number(env.MESH_MONGODB_PORT ?? url?.port ?? "27017") || 27017; + const user = env.MESH_MONGODB_USER ?? (url?.username ? decodeURIComponent(url.username) : "root"); + const authSource = + env.MESH_MONGODB_AUTHSOURCE ?? url?.searchParams.get("authSource") ?? "admin"; + const password = env.MESH_MONGODB_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); + if (!host || !password) { + throw new Error("mongodb host or admin password is not set — mongodb's own code cannot reach the server"); + } + return new MongoClient({ host, port, user, password, authSource }); + } + + get host(): string { + return this.conn.host; + } + + get port(): number { + return this.conn.port; + } + + /** The admin connection URI mongosh authenticates with, credentials percent-encoded. */ + private uri(): string { + const u = encodeURIComponent(this.conn.user); + const p = encodeURIComponent(this.conn.password); + const a = encodeURIComponent(this.conn.authSource); + return `mongodb://${u}:${p}@${this.conn.host}:${this.conn.port}/?authSource=${a}`; + } + + /** + * Evaluate a JavaScript snippet server-side through `mongosh` and parse the JSON it prints (see + * header). The snippet MUST `print()` exactly one JSON document as its only stdout — every method + * below ends in `print(EJSON.stringify(...))`. `--quiet` suppresses the shell banner so stdout is + * the JSON alone; a non-zero exit (auth failure, bad command) rejects here rather than returning + * a partial success. + */ + async evalJs(js: string): Promise { + const { stdout } = await run( + "mongosh", + [this.uri(), "--quiet", "--eval", js], + { maxBuffer: 16 << 20 }, + ); + const text = stdout.trim(); + if (text.length === 0) { + throw new Error("mongosh returned no output — the eval printed nothing"); + } + return JSON.parse(text) as T; + } + + /** + * Create a login user and the database it owns, idempotently. The user is created inside the + * target database with the `dbOwner` role scoped to that database, so the consumer owns exactly + * its own and authenticates with the target database as its authSource. Re-running updates the + * password and roles, so a rotated credential converges. + */ + async createDatabaseAndUser(database: string, user: string, password: string): Promise { + const js = ` +const target = db.getSiblingDB(${lit(database)}); +let existing = null; +try { existing = target.getUser(${lit(user)}); } catch (e) { existing = null; } +const roles = [{ role: "dbOwner", db: ${lit(database)} }]; +if (existing) { + target.updateUser(${lit(user)}, { pwd: ${lit(password)}, roles: roles }); +} else { + target.createUser({ user: ${lit(user)}, pwd: ${lit(password)}, roles: roles }); +} +print(EJSON.stringify({ ok: 1 })); +`; + await this.evalJs<{ ok: number }>(js); + } + + /** Drop a database and its owning user, idempotently. Dropping the database evicts its data; the + * user is removed first so a re-grant of the same login starts clean. */ + async dropDatabaseAndUser(database: string, user: string): Promise { + const js = ` +const target = db.getSiblingDB(${lit(database)}); +try { target.dropUser(${lit(user)}); } catch (e) {} +target.dropDatabase(); +print(EJSON.stringify({ ok: 1 })); +`; + await this.evalJs<{ ok: number }>(js); + } + + /** List the databases on the server, with on-disk size, for the mongodb_list_databases tool. */ + async listDatabases(): Promise { + const res = await this.evalJs<{ databases: { name: string; sizeOnDisk?: number }[] }>( + `print(EJSON.stringify(db.adminCommand({ listDatabases: 1 })));`, + ); + return (res.databases ?? []) + .map((d) => ({ name: String(d.name), sizeBytes: Number(d.sizeOnDisk ?? 0) })) + .sort((a, b) => a.name.localeCompare(b.name)); + } + + /** + * Run a read-only `find` against a collection in a named database, for the mongodb_query tool. + * `find` mutates nothing; the limit is capped so a tool call cannot stream an unbounded result. + */ + async find( + database: string, + collection: string, + filter: Readonly>, + limit: number, + ): Promise[]> { + const capped = Math.max(1, Math.min(limit, 1000)); + const js = + `print(EJSON.stringify(` + + `db.getSiblingDB(${lit(database)}).getCollection(${lit(collection)})` + + `.find(${JSON.stringify(filter)}).limit(${capped}).toArray()` + + `));`; + return this.evalJs[]>(js); + } +} + +/** Generate a URL-safe password. */ +export function generatePassword(): string { + return randomBytes(24).toString("base64url"); +} + +/** Embed a value as a JavaScript literal inside a mongosh snippet — JSON.stringify escapes quotes, + * backslashes and control characters, so a string cannot break out of the snippet. */ +function lit(val: unknown): string { + return JSON.stringify(val); +} + +function readSecretFile(path: string | undefined): string | undefined { + if (!path) return undefined; + try { + return readFileSync(path, "utf8").trim(); + } catch { + return undefined; + } +} + +function safeUrl(raw: string): URL | undefined { + try { + return new URL(raw); + } catch { + return undefined; + } +} diff --git a/modules/mongodb/index.ts b/modules/mongodb/index.ts new file mode 100644 index 0000000..1a8bf8b --- /dev/null +++ b/modules/mongodb/index.ts @@ -0,0 +1,25 @@ +// mongodb's events entrypoint, loaded by the per-node tool host (the provisioner container runs +// ./provisioner separately). The database lifecycle events are EMITTED from the provisioner, where +// the lifecycle actually happens (novox/hq ADR 0046/0047): +// module.mongodb.database.provisioned — a consumer's database + owning user was created +// module.mongodb.database.deprovisioned — that database was removed +// Here in the tool host we react to them, keeping a lightweight audit trail of who was granted a +// database and who lost one — observability the provider itself is best placed to log. + +import { on } from "@novox/mesh-sdk/events"; + +interface DatabaseEvent { + consumer: string; + database: string; + user?: string; +} + +await on("module.mongodb.database.provisioned", async (e) => { + console.log(`[mongodb] database provisioned for ${e.body.consumer} (db ${e.body.database})`); +}); + +await on("module.mongodb.database.deprovisioned", async (e) => { + console.log(`[mongodb] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`); +}); + +console.log("[mongodb] auditing database lifecycle events"); diff --git a/modules/mongodb/module.json b/modules/mongodb/module.json new file mode 100644 index 0000000..0c81adb --- /dev/null +++ b/modules/mongodb/module.json @@ -0,0 +1,117 @@ +{ + "module": "mongodb", + "version": "1", + "provides": [ + { + "name": "mongodb-database", + "scope": "mesh" + } + ], + "capabilities": [ + "container-runtime" + ], + "emits": [ + "module.mongodb.database.provisioned", + "module.mongodb.database.deprovisioned" + ], + "consumes": [ + "module.mongodb.database.provisioned", + "module.mongodb.database.deprovisioned" + ], + "listens": [ + { + "port": 27017, + "protocol": "tcp", + "from": "mesh", + "why": "modules on any machine that were granted a database" + } + ], + "serves": { + "mongodb-database": {} + }, + "receives": { + "mongodb-database": "/var/lib/mongodb/grants/mesh.json" + }, + "grants": { + "mongodb-database": "/var/lib/mongodb/grants" + }, + "own-secrets": { + "root": "/var/lib/mongodb/root.secret", + "broker": "/var/lib/mesh/mongodb/broker" + }, + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/mongodb", + "mode": "0700" + }, + { + "id": "state", + "type": "directory", + "path": "/var/lib/mongodb", + "mode": "0700" + }, + { + "id": "grants", + "type": "directory", + "path": "/var/lib/mongodb/grants", + "mode": "0700" + }, + { + "id": "root-env", + "type": "file", + "path": "/var/lib/mongodb/root.env", + "mode": "0600", + "content": "MONGO_INITDB_ROOT_PASSWORD=${secret:root}\n" + }, + { + "id": "data", + "type": "directory", + "path": "/services/mongodb/db-data", + "mode": "0700" + }, + { + "id": "net", + "type": "network", + "name": "mongodb" + }, + { + "id": "server", + "type": "container", + "name": "mongo", + "image": "mongo@sha256:e3fa459b4f4b72f3257c67a23c145e250b8b5700f033860392c68539b998bbe3", + "network": "mongodb", + "env": { + "MONGO_INITDB_ROOT_USERNAME": "root" + }, + "env-file": [ + "/var/lib/mongodb/root.env" + ], + "ports": [ + "27017" + ], + "volumes": [ + "/services/mongodb/db-data:/data/db" + ] + }, + { + "id": "runtime", + "type": "container", + "name": "mesh-mongodb", + "image": "mesh-runtime-mongodb@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "mongodb", + "volumes": [ + "/var/lib/mesh/mongodb/broker:/run/secrets/broker:ro", + "/var/lib/mongodb/grants:/var/lib/mongodb/grants:ro", + "/var/lib/mongodb/root.secret:/run/secrets/root:ro" + ], + "env": { + "MESH_PROVISION_MONGODB": "mongodb://root@mongo:27017/admin?authSource=admin", + "MESH_PROVISION_PASSWORD_FILE": "/run/secrets/root", + "MESH_BROKER_FILE": "/run/secrets/broker", + "MESH_RECEIVES": "/var/lib/mongodb/grants/mesh.json" + } + } + ] +} diff --git a/modules/mongodb/package.json b/modules/mongodb/package.json new file mode 100644 index 0000000..509b818 --- /dev/null +++ b/modules/mongodb/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-mongodb", + "version": "0.1.0", + "description": "mongodb — provides the mesh mongodb-database interface. Its client, provisioner, 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/mongodb/provisioner/index.ts b/modules/mongodb/provisioner/index.ts new file mode 100644 index 0000000..27a069c --- /dev/null +++ b/modules/mongodb/provisioner/index.ts @@ -0,0 +1,48 @@ +// mongodb's provisioner — the adapter that makes mongodb a provider of the mesh `mongodb-database` +// interface. The reconcile loop, the contributions file, and reading the mesh's minted password are +// the sdk harness's; this writes only the per-service half: how mongodb creates and removes a +// consumer's database + owning user (novox/hq ADR 0044/0045/0053). +// +// The `mongodb-database` interface: a consumer connects to a database it alone owns, as `as` with the +// password the mesh minted, authenticating against that same database. +// +// **The user name and password are the mesh's, not the provisioner's (ADR 0053).** The mesh derives +// the login and hands it to both ends, and mints the password. mongodb creates a user and a +// same-named database under exactly that login — a name the consumer cannot learn is a database it +// cannot reach. +// +// The commands run through MongoClient.evalJs(), which is the module's one execution boundary (see +// client.ts). + +import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner"; +import { emit } from "@novox/mesh-sdk/events"; +import { MongoClient } from "../client.js"; + +const mongo = MongoClient.fromEnv(); + +/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */ +async function announce(type: string, body: Record): Promise { + try { + await emit(type, body); + } catch (err) { + console.error(`[provisioner:mongodb-database] emit ${type} failed: ${err}`); + } +} + +runProvisioner("mongodb-database", { + async create(p: Provision): Promise { + // Database and owning user share the consumer's login, so the consumer owns exactly its own. + const database = p.as; + await mongo.createDatabaseAndUser(database, p.as, p.password); + await announce("module.mongodb.database.provisioned", { + consumer: p.consumer ?? "", + database, + user: p.as, + }); + }, + + async remove(p: { as: string }): Promise { + await mongo.dropDatabaseAndUser(p.as, p.as); + await announce("module.mongodb.database.deprovisioned", { database: p.as }); + }, +}); diff --git a/modules/mongodb/tools/index.ts b/modules/mongodb/tools/index.ts new file mode 100644 index 0000000..5c14acd --- /dev/null +++ b/modules/mongodb/tools/index.ts @@ -0,0 +1,51 @@ +// mongodb's tools — mongodb's own code (novox/hq ADR 0044), importing mongodb's own client. They +// return structured data; the mesh serves them through the sdk's tool harness. Both call through +// MongoClient.evalJs(), the module's one execution boundary (see client.ts). + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { MongoClient } from "../client.js"; + +export function getMongoTools(mongo: MongoClient): ToolDefinition[] { + return [ + { + name: "mongodb_list_databases", + description: "List the databases on the mongodb server, with their on-disk size.", + input: {}, + run: async () => ({ databases: await mongo.listDatabases() }), + }, + { + name: "mongodb_query", + description: "Run a read-only find against a collection in a named database and return the matching documents.", + input: { + database: { type: "string", description: "the database to query" }, + collection: { type: "string", description: "the collection to read from" }, + filter: { type: "object", description: "the MongoDB query filter (defaults to {} — all documents)" }, + limit: { type: "number", description: "maximum documents to return (default 100, capped at 1000)" }, + }, + run: async (args) => { + const database = String(args.database ?? ""); + const collection = String(args.collection ?? ""); + if (!database) throw new Error("mongodb_query: database is required"); + if (!collection) throw new Error("mongodb_query: collection is required"); + const filter = isObject(args.filter) ? args.filter : {}; + const limit = Number(args.limit ?? 100) || 100; + const documents = await mongo.find(database, collection, filter, limit); + return { database, collection, documents }; + }, + }, + ]; +} + +function isObject(v: unknown): v is Record { + return typeof v === "object" && v !== null && !Array.isArray(v); +} + +// The tools exist only when the server can be reached from the environment; without it, mongodb +// contributes none rather than failing the whole tool runtime. +registerModuleTools("mongodb", (env) => { + try { + return getMongoTools(MongoClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/mongodb/tsconfig.json b/modules/mongodb/tsconfig.json new file mode 100644 index 0000000..51f4046 --- /dev/null +++ b/modules/mongodb/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", "provisioner/index.ts", "tools/index.ts"] +} diff --git a/modules/mosquitto/client.ts b/modules/mosquitto/client.ts new file mode 100644 index 0000000..ca70939 --- /dev/null +++ b/modules/mosquitto/client.ts @@ -0,0 +1,184 @@ +// mosquitto's admin client — mosquitto's own code, living in the module (novox/hq ADR 0044). Both +// this module's tools and its provisioner import it, and nothing outside mosquitto does. +// +// Client, role and ACL administration is driven through `mosquitto_ctrl dynsec`, not a hand-rolled +// MQTT stack: the module may take NO npm dependency beyond @novox/mesh-sdk, and mosquitto ships the +// exact admin client for its Dynamic Security plugin — so it shells out to it, the same way postgres +// drives itself through `psql`, minio through `mc` and mailu through `doveadm`. One boundary, +// `ctl()`, and every method is built on it. +// +// Why the Dynamic Security plugin and not `password_file`: dynsec creates and revokes clients while +// the broker runs, over an admin connection, with no broker restart and no file the host rewrites — +// the true analog of redis's runtime ACL users. A `password_file` would have to be re-read on a +// SIGHUP the runtime cannot cleanly send across containers, and — declared as a managed file — would +// be rewritten by the host on every reconcile, wiping every provisioned user (novox/nox issue 011). +// The one cost dynsec carries is the bootstrap file; see initBootstrapFile() and the module README. + +import { randomBytes } from "node:crypto"; +import { readFileSync } from "node:fs"; +import { execFile } from "node:child_process"; +import { promisify } from "node:util"; + +const run = promisify(execFile); + +export interface MqttConn { + readonly host: string; + readonly port: number; + /** The Dynamic Security admin client the runtime authenticates as. */ + readonly adminUser: string; + readonly adminPassword: string; +} + +export class MosquittoClient { + constructor(private readonly conn: MqttConn) {} + + /** + * Build from the module's resolved environment. Reads MESH_MQTT_* first (the documented names), + * falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner + * container. Throws if it cannot find a host and an admin password — the right failure, because + * without them nothing it does can work. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): MosquittoClient { + const endpoint = env.MESH_PROVISION_MQTT ?? ""; // "host:port" + const host = env.MESH_MQTT_HOST ?? (endpoint ? endpoint.split(":")[0] : undefined); + const port = + Number(env.MESH_MQTT_PORT ?? (endpoint.includes(":") ? endpoint.split(":")[1] : "") ?? "1883") || 1883; + const adminUser = env.MESH_MQTT_ADMIN_USER ?? env.MESH_PROVISION_ADMIN_USER ?? "mesh-admin"; + const adminPassword = env.MESH_MQTT_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); + if (!host || !adminPassword) { + throw new Error( + "mosquitto host or admin password is not set — mosquitto's own code cannot reach the broker", + ); + } + return new MosquittoClient({ host, port, adminUser, adminPassword: adminPassword ?? "" }); + } + + get host(): string { + return this.conn.host; + } + + get port(): number { + return this.conn.port; + } + + /** + * Run one `mosquitto_ctrl dynsec ` command against the broker as the admin client and return + * its stdout. Connects over MQTT with the verified connect flags `-h`/`-p`/`-u`/`-P`. A non-zero + * exit rejects — a failed command is an error here, not a success with a warning. + * + * The admin password rides on argv (`-P`): mosquitto_ctrl 2.x exposes no password env var and no + * password file for a broker connection — its only non-interactive mechanism is `-P`, its only + * other mechanism an interactive prompt. This is a real mosquitto limitation, not a choice; unlike + * psql's PGPASSWORD there is nothing cleaner to reach for. The exposure is momentary and confined + * to this single-purpose runtime container; see the module README. + */ + async ctl(...args: string[]): Promise { + const base = [ + "-h", this.conn.host, + "-p", String(this.conn.port), + "-u", this.conn.adminUser, + "-P", this.conn.adminPassword, + ]; + const { stdout } = await run("mosquitto_ctrl", [...base, "dynsec", ...args], { + maxBuffer: 16 << 20, + }); + return stdout; + } + + /** Whether a dynsec client with this username already exists. */ + async clientExists(username: string): Promise { + try { + await this.ctl("getClient", username); + return true; + } catch { + return false; + } + } + + /** + * Create (or reset to a known state) a client scoped to one topic namespace, idempotently. The + * client is confined to `/#` by a same-named role: it may publish to, subscribe to and + * receive on exactly its own subtree and nothing else — the MQTT analog of redis's keyspace-scoped + * ACL user. Called again for an existing client, it resets the password and re-asserts the ACLs. + */ + async createScopedClient(username: string, password: string, topicPrefix: string): Promise { + const role = username; // one role per client, named for it + const pattern = `${topicPrefix}/#`; + + if (await this.clientExists(username)) { + await this.ctl("setClientPassword", username, password); + } else { + await this.ctl("createClient", username, "-p", password); + } + + // A role carrying exactly this client's topic ACLs. createRole fails if it already exists; that + // is fine — the setRoleACL calls below assert the intended state either way. + await ignoreExisting(this.ctl("createRole", role)); + for (const acl of ["publishClientSend", "publishClientReceive", "subscribePattern"]) { + // allow (1) this client to send to, receive on, and subscribe under its own subtree. + await this.ctl("addRoleACL", role, acl, pattern, "allow"); + } + await ignoreExisting(this.ctl("addClientRole", username, role)); + } + + /** Remove a client and the per-client role created for it, idempotently. */ + async deleteScopedClient(username: string): Promise { + await ignoreMissing(this.ctl("deleteClient", username)); + await ignoreMissing(this.ctl("deleteRole", username)); + } + + /** The dynsec client list, parsed from `listClients`. */ + async listClients(): Promise { + const out = await this.ctl("listClients"); + return out + .split(/\r?\n/) + .map((l) => l.trim()) + .filter((l) => l.length > 0); + } + + /** + * Write the Dynamic Security bootstrap file offline, creating the admin client the plugin loads at + * broker startup. This is a one-time seed, NOT part of the reconcile loop: run once before the + * broker first starts, against the same path the broker's `plugin_opt_config_file` names. It must + * never be a host-reconciled managed file — see the module README and novox/nox issue 011. + */ + async initBootstrapFile(configFile: string): Promise { + // `dynsec init [admin-password]` is an offline file operation — it does + // not connect to the broker. The password is a positional argument (omitting it prompts). + await run("mosquitto_ctrl", ["dynsec", "init", configFile, this.conn.adminUser, this.conn.adminPassword], { + maxBuffer: 16 << 20, + }); + } +} + +/** Generate a URL-safe password with no argv- or MQTT-hostile characters. */ +export function generatePassword(): string { + return randomBytes(24).toString("base64url"); +} + +/** Swallow a "already exists" failure so create paths are idempotent; rethrow anything else. */ +async function ignoreExisting(p: Promise): Promise { + try { + await p; + } catch (err) { + if (!/exist/i.test(String(err))) throw err; + } +} + +/** Swallow a "not found" failure so delete paths are idempotent; rethrow anything else. */ +async function ignoreMissing(p: Promise): Promise { + try { + await p; + } catch (err) { + if (!/not\s*found|does not exist|no such/i.test(String(err))) throw err; + } +} + +function readSecretFile(path: string | undefined): string | undefined { + if (!path) return undefined; + try { + return readFileSync(path, "utf8").trim(); + } catch { + return undefined; + } +} diff --git a/modules/mosquitto/index.ts b/modules/mosquitto/index.ts new file mode 100644 index 0000000..e3e3170 --- /dev/null +++ b/modules/mosquitto/index.ts @@ -0,0 +1,25 @@ +// mosquitto's events entrypoint, loaded by the per-node tool host (the provisioner container runs +// ./provisioner separately). The topic lifecycle events are EMITTED from the provisioner, where the +// lifecycle actually happens (novox/hq ADR 0046/0047): +// module.mosquitto.topic.provisioned — a consumer's client + scoped role was created +// module.mosquitto.topic.deprovisioned — that client was removed +// Here in the tool host we react to them, keeping a lightweight audit trail of who was granted a +// topic namespace and who lost one — observability the provider itself is best placed to log. + +import { on } from "@novox/mesh-sdk/events"; + +interface TopicEvent { + consumer: string; + username: string; + topicPrefix?: string; +} + +await on("module.mosquitto.topic.provisioned", async (e) => { + console.log(`[mosquitto] topic provisioned for ${e.body.consumer} (client ${e.body.username})`); +}); + +await on("module.mosquitto.topic.deprovisioned", async (e) => { + console.log(`[mosquitto] topic deprovisioned for ${e.body.consumer} (client ${e.body.username})`); +}); + +console.log("[mosquitto] auditing topic lifecycle events"); diff --git a/modules/mosquitto/module.json b/modules/mosquitto/module.json new file mode 100644 index 0000000..b804b82 --- /dev/null +++ b/modules/mosquitto/module.json @@ -0,0 +1,122 @@ +{ + "module": "mosquitto", + "version": "1", + "provides": [ + { + "name": "mqtt-topic", + "scope": "mesh" + } + ], + "capabilities": [ + "container-runtime" + ], + "emits": [ + "module.mosquitto.topic.provisioned", + "module.mosquitto.topic.deprovisioned" + ], + "consumes": [ + "module.mosquitto.topic.provisioned", + "module.mosquitto.topic.deprovisioned" + ], + "serves": { + "mqtt-topic": {} + }, + "receives": { + "mqtt-topic": "/var/lib/mosquitto-module/grants/mesh.json" + }, + "grants": { + "mqtt-topic": "/var/lib/mosquitto-module/grants" + }, + "own-secrets": { + "admin": "/var/lib/mosquitto-module/admin.secret", + "broker": "/var/lib/mesh/mosquitto/broker" + }, + "listens": [ + { + "port": 1883, + "protocol": "tcp", + "from": "mesh", + "why": "modules on any machine that were granted a topic namespace" + }, + { + "port": 8081, + "protocol": "tcp", + "from": "mesh", + "why": "the same broker over MQTT-on-WebSockets, for browser clients" + } + ], + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/mosquitto", + "mode": "0700" + }, + { + "id": "state", + "type": "directory", + "path": "/var/lib/mosquitto-module", + "mode": "0700" + }, + { + "id": "grants-dir", + "type": "directory", + "path": "/var/lib/mosquitto-module/grants", + "mode": "0700" + }, + { + "id": "data", + "type": "directory", + "path": "/services/mosquitto/data", + "mode": "0700", + "owner": "1883:1883" + }, + { + "id": "server-conf", + "type": "file", + "path": "/var/lib/mosquitto-module/mosquitto.conf", + "mode": "0600", + "owner": "1883:1883", + "content": "persistence true\npersistence_location /mosquitto/data\n\nlog_dest stdout\nlog_type warning\nlog_type error\nlog_type notice\n\n# Every client authenticates; identities and their per-topic ACLs are managed\n# at runtime by the dynamic security plugin, whose store the plugin itself owns.\nallow_anonymous false\nplugin /usr/lib/mosquitto_dynamic_security.so\nplugin_opt_config_file /mosquitto/data/dynamic-security.json\n\n# MQTT listener\nlistener 1883\n\n# MQTT-over-WebSockets listener\nlistener 8081\nprotocol websockets\n" + }, + { + "id": "net", + "type": "network", + "name": "mosquitto" + }, + { + "id": "server", + "type": "container", + "name": "mosquitto", + "image": "eclipse-mosquitto@sha256:6f8d8a947c506f8a2290ec65cd4bd2bc7cb4d43fb5f6271f861cb013e2ef9797", + "network": "mosquitto", + "ports": [ + "1883", + "8081" + ], + "volumes": [ + "/services/mosquitto/data:/mosquitto/data", + "/var/lib/mosquitto-module/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro" + ] + }, + { + "id": "runtime", + "type": "container", + "name": "mesh-mosquitto", + "image": "mesh-runtime-mosquitto@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "mosquitto", + "volumes": [ + "/var/lib/mesh/mosquitto/broker:/run/secrets/broker:ro", + "/var/lib/mosquitto-module/grants:/var/lib/mosquitto-module/grants:ro", + "/var/lib/mosquitto-module/admin.secret:/run/secrets/admin:ro" + ], + "env": { + "MESH_BROKER_FILE": "/run/secrets/broker", + "MESH_RECEIVES": "/var/lib/mosquitto-module/grants/mesh.json", + "MESH_PROVISION_MQTT": "mosquitto:1883", + "MESH_PROVISION_ADMIN_USER": "mesh-admin", + "MESH_PROVISION_PASSWORD_FILE": "/run/secrets/admin" + } + } + ] +} diff --git a/modules/mosquitto/package.json b/modules/mosquitto/package.json new file mode 100644 index 0000000..5b89552 --- /dev/null +++ b/modules/mosquitto/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-mosquitto", + "version": "0.1.0", + "description": "mosquitto — provides the mesh mqtt-topic interface. Its admin client, provisioner, 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/mosquitto/provisioner/index.ts b/modules/mosquitto/provisioner/index.ts new file mode 100644 index 0000000..24c3bc7 --- /dev/null +++ b/modules/mosquitto/provisioner/index.ts @@ -0,0 +1,46 @@ +// mosquitto's provisioner — the adapter that makes mosquitto a provider of the mesh `mqtt-topic` +// interface. The reconcile loop, the contributions file, and reading the mesh's minted password are +// the sdk harness's; this writes only the per-service half: how mosquitto creates and removes a +// per-consumer MQTT client (novox/hq ADR 0044/0045/0053). +// +// The `mqtt-topic` interface: a consumer connects as `as` with the password the mesh minted, and +// publishes and subscribes under `/#`, isolated from every other consumer by a Dynamic Security +// role scoped to exactly that subtree. +// +// **The login and password are the mesh's, not the provisioner's (ADR 0053).** The mesh derives the +// login and hands it to both ends so they agree, and mints the password and delivers a copy to each. +// mosquitto creates exactly that client with exactly that password — a name or password the +// provisioner invented is one the consumer could never present. + +import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner"; +import { emit } from "@novox/mesh-sdk/events"; +import { MosquittoClient } from "../client.js"; + +const mosquitto = MosquittoClient.fromEnv(); + +/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */ +async function announce(type: string, body: Record): Promise { + try { + await emit(type, body); + } catch (err) { + console.error(`[provisioner:mqtt-topic] emit ${type} failed: ${err}`); + } +} + +runProvisioner("mqtt-topic", { + async create(p: Provision): Promise { + // The topic subtree is scoped to the consumer's own login, so one cannot read another's topics. + const topicPrefix = p.as; + await mosquitto.createScopedClient(p.as, p.password, topicPrefix); + await announce("module.mosquitto.topic.provisioned", { + consumer: p.consumer ?? "", + username: p.as, + topicPrefix, + }); + }, + + async remove(p: { as: string }): Promise { + await mosquitto.deleteScopedClient(p.as); + await announce("module.mosquitto.topic.deprovisioned", { username: p.as }); + }, +}); diff --git a/modules/mosquitto/tools/index.ts b/modules/mosquitto/tools/index.ts new file mode 100644 index 0000000..ffbbcb0 --- /dev/null +++ b/modules/mosquitto/tools/index.ts @@ -0,0 +1,54 @@ +// mosquitto's tools — mosquitto's own code (novox/hq ADR 0044), importing mosquitto's own admin +// 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 { MosquittoClient } from "../client.js"; + +export function getMosquittoTools(mosquitto: MosquittoClient): ToolDefinition[] { + return [ + { + name: "mqtt_list_clients", + description: "List the Dynamic Security clients registered on the mosquitto broker.", + input: {}, + run: async () => ({ clients: await mosquitto.listClients() }), + }, + { + name: "mqtt_get_client", + description: "Show one Dynamic Security client — its roles and enabled state.", + input: { username: { type: "string", description: "the client's username" } }, + run: async (args) => { + const username = String(args.username ?? ""); + if (!username) throw new Error("mqtt_get_client: username is required"); + return { username, detail: await mosquitto.ctl("getClient", username) }; + }, + }, + { + name: "mqtt_ctrl", + description: + "Run an arbitrary 'mosquitto_ctrl dynsec' subcommand, e.g. 'listRoles', 'getRole myrole'. Admin surface.", + input: { command: { type: "string", description: "the dynsec subcommand and its arguments, space-separated" } }, + run: async (args) => { + const parts = tokenize(String(args.command ?? "")); + if (parts.length === 0) throw new Error("mqtt_ctrl: empty command"); + const output = await mosquitto.ctl(...parts); + return { command: parts.join(" "), output }; + }, + }, + ]; +} + +/** Split a command line into arguments, honouring double-quoted spans. */ +function tokenize(command: string): string[] { + const matches = command.match(/(?:[^\s"]+|"[^"]*")+/g) ?? []; + return matches.map((p) => p.replace(/^"|"$/g, "")); +} + +// The tools exist only when the broker can be reached from the environment; without it, mosquitto +// contributes none rather than failing the whole tool runtime. +registerModuleTools("mosquitto", (env) => { + try { + return getMosquittoTools(MosquittoClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/mosquitto/tsconfig.json b/modules/mosquitto/tsconfig.json new file mode 100644 index 0000000..51f4046 --- /dev/null +++ b/modules/mosquitto/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", "provisioner/index.ts", "tools/index.ts"] +} diff --git a/modules/mssql/client.ts b/modules/mssql/client.ts new file mode 100644 index 0000000..308bce2 --- /dev/null +++ b/modules/mssql/client.ts @@ -0,0 +1,242 @@ +// mssql's admin client — mssql's own code, living in the module (novox/hq ADR 0044). Both this +// module's tools and its provisioner import it, and nothing outside mssql does. +// +// SQL is executed through `sqlcmd`, not a wire-protocol driver: the module may take NO npm +// dependency beyond @novox/mesh-sdk, and hand-rolling the TDS handshake, pre-login and query +// protocol is more surface than this should carry — so it shells out to the client the mssql +// tools ship, the same way postgres drives itself through `psql`, minio through `mc`, and mailu +// through doveadm. One boundary, `run()`, and every method is built on it. +// +// Structured rows come back as JSON: SQL Server itself renders the result with `FOR JSON`, and +// this parses the single JSON document sqlcmd prints — far more robust than parsing sqlcmd's +// column-aligned text, since SQL Server owns the quoting and typing. + +import { randomBytes } from "node:crypto"; +import { readFileSync } from "node:fs"; +import { execFile } from "node:child_process"; +import { promisify } from "node:util"; + +const run = promisify(execFile); + +export interface QueryResult { + /** The leading keyword of the statement, e.g. "SELECT", "CREATE". */ + readonly command: string; + readonly rows: Record[]; +} + +export interface MssqlConn { + readonly host: string; + readonly port: number; + readonly user: string; + readonly password: string; +} + +export class MssqlClient { + constructor(private readonly conn: MssqlConn) {} + + /** + * Build from the module's resolved environment. Reads MESH_MSSQL_* first (the documented + * names), falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner + * container. Throws if it cannot find a host and an admin password. + */ + static fromEnv(env: NodeJS.ProcessEnv = process.env): MssqlClient { + const url = env.MESH_PROVISION_MSSQL ? safeUrl(env.MESH_PROVISION_MSSQL) : undefined; + const host = env.MESH_MSSQL_HOST ?? url?.hostname; + const port = Number(env.MESH_MSSQL_PORT ?? url?.port ?? "1433") || 1433; + const user = env.MESH_MSSQL_USER ?? url?.username ?? "sa"; + const password = env.MESH_MSSQL_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE); + if (!host || !password) { + throw new Error("mssql host or admin password is not set — mssql's own code cannot reach the server"); + } + return new MssqlClient({ host, port, user, password }); + } + + get host(): string { + return this.conn.host; + } + + get port(): number { + return this.conn.port; + } + + /** + * Execute a batch that returns no rows (DDL and the like), through `sqlcmd`. The password is + * passed by SQLCMDPASSWORD, never on argv, the way postgres passes PGPASSWORD; `-b` makes a + * failed statement an error here rather than a success with a warning, and `-C` trusts the + * server's self-signed certificate the mssql image ships with. + */ + async exec(sql: string, database = "master"): Promise { + await this.sqlcmd(sql, database); + } + + /** + * Run a SELECT and return its rows as objects. The caller's SQL must be a single SELECT; it is + * wrapped so SQL Server renders the result with `FOR JSON PATH`, and the JSON document sqlcmd + * prints (split across output lines for a large result, and reassembled here) is parsed. An + * empty result yields no output at all — an empty array. + */ + async query(select: string, database = "master"): Promise[]> { + const wrapped = `SET NOCOUNT ON;\n${stripTrailingSemis(select)}\nFOR JSON PATH, INCLUDE_NULL_VALUES;`; + const stdout = await this.sqlcmd(wrapped, database); + return parseJsonRows(stdout); + } + + /** The one execution boundary: invoke `sqlcmd` and return its concatenated stdout. */ + private async sqlcmd(sql: string, database: string): Promise { + // `-h -1` drops the column-header rule; `-y 0`/`-Y 0` lift the display-width cap so a long + // JSON document is not truncated; `-W` trims trailing whitespace so the JSON chunks rejoin + // cleanly. sqlcmd from the mssql-tools ships in the runtime container, the way `psql` ships + // with postgres's — the module owns its own code (ADR 0044) and shells out to it. + const { stdout } = await run( + "sqlcmd", + [ + "-S", `${this.conn.host},${this.conn.port}`, + "-U", this.conn.user, + "-d", database, + "-C", + "-b", + "-h", "-1", + "-y", "0", + "-Y", "0", + "-W", + "-Q", sql, + ], + { env: { ...process.env, SQLCMDPASSWORD: this.conn.password }, maxBuffer: 16 << 20 }, + ); + return stdout; + } + + /** + * Create a login and a database it owns (mapped as a db_owner user), idempotently. The login, + * the database and the user all carry the consumer's minted name, so the consumer owns exactly + * its own database — a name it cannot learn is a database it cannot reach (ADR 0053). + */ + async createDatabaseAndLogin(database: string, login: string, password: string): Promise { + const logins = await this.query( + `SELECT 1 AS ok FROM sys.server_principals WHERE name = ${literal(login)}`, + ); + if (logins.length === 0) { + await this.exec( + `CREATE LOGIN ${ident(login)} WITH PASSWORD = ${literal(password)}, CHECK_POLICY = OFF`, + ); + } else { + await this.exec(`ALTER LOGIN ${ident(login)} WITH PASSWORD = ${literal(password)}`); + } + + const dbs = await this.query( + `SELECT 1 AS ok FROM sys.databases WHERE name = ${literal(database)}`, + ); + if (dbs.length === 0) { + // CREATE DATABASE must stand alone in its batch; it runs as its own sqlcmd invocation. + await this.exec(`CREATE DATABASE ${ident(database)}`); + } + + // Map the login to a db_owner user inside the database it owns. + const users = await this.query( + `SELECT 1 AS ok FROM sys.database_principals WHERE name = ${literal(login)}`, + database, + ); + if (users.length === 0) { + await this.exec(`CREATE USER ${ident(login)} FOR LOGIN ${ident(login)}`, database); + } + await this.exec(`ALTER ROLE db_owner ADD MEMBER ${ident(login)}`, database); + } + + /** Drop a database and its login, idempotently, after evicting live connections. */ + async dropDatabaseAndLogin(database: string, login: string): Promise { + const dbs = await this.query( + `SELECT 1 AS ok FROM sys.databases WHERE name = ${literal(database)}`, + ); + if (dbs.length > 0) { + // SINGLE_USER WITH ROLLBACK IMMEDIATE evicts every other session before the drop. + await this.exec(`ALTER DATABASE ${ident(database)} SET SINGLE_USER WITH ROLLBACK IMMEDIATE`); + await this.exec(`DROP DATABASE ${ident(database)}`); + } + const logins = await this.query( + `SELECT 1 AS ok FROM sys.server_principals WHERE name = ${literal(login)}`, + ); + if (logins.length > 0) { + await this.exec(`DROP LOGIN ${ident(login)}`); + } + } + + /** List the user databases (database_id > 4 excludes the system four), with size, for the tool. */ + async listDatabases(): Promise<{ name: string; sizeBytes: number; state: string }[]> { + const rows = await this.query( + "SELECT d.name AS name, d.state_desc AS state, " + + "SUM(CAST(f.size AS bigint)) * 8 * 1024 AS size_bytes " + + "FROM sys.databases d JOIN sys.master_files f ON d.database_id = f.database_id " + + "WHERE d.database_id > 4 GROUP BY d.name, d.state_desc ORDER BY d.name", + ); + return rows.map((r) => ({ + name: String(r.name), + sizeBytes: Number(r.size_bytes ?? 0), + state: String(r.state ?? ""), + })); + } + + /** Run a read-only SELECT against a named database, for the mssql_query tool. */ + async readOnlyQuery(database: string, sql: string): Promise { + // The read-only guarantee is a wrapping transaction that is always rolled back: any write the + // statement attempts is undone. The rows are rendered by FOR JSON inside query(). + const rows = await this.query( + `BEGIN TRANSACTION;\n${stripTrailingSemis(sql)}\nFOR JSON PATH, INCLUDE_NULL_VALUES;\nROLLBACK;`, + database, + ); + return { command: sql.trimStart().split(/\s+/)[0]?.toUpperCase() ?? "", rows }; + } +} + +/** Generate a URL-safe password. */ +export function generatePassword(): string { + return randomBytes(24).toString("base64url"); +} + +/** Quote a T-SQL identifier (square brackets, doubled internal `]`). */ +export function ident(id: string): string { + return "[" + id.replace(/]/g, "]]") + "]"; +} + +/** Quote a T-SQL string literal (single quotes, doubled internal quotes). */ +export function literal(val: string): string { + return "'" + val.replace(/'/g, "''") + "'"; +} + +/** Strip trailing semicolons and whitespace so FOR JSON can be appended to a caller's SELECT. */ +function stripTrailingSemis(sql: string): string { + return sql.replace(/[\s;]+$/, ""); +} + +function readSecretFile(path: string | undefined): string | undefined { + if (!path) return undefined; + try { + return readFileSync(path, "utf8").trim(); + } catch { + return undefined; + } +} + +function safeUrl(raw: string): URL | undefined { + try { + return new URL(raw); + } catch { + return undefined; + } +} + +/** + * Parse the JSON a FOR JSON query prints through sqlcmd. SQL Server splits a large FOR JSON result + * into ~2033-character chunks, one per output row; with `-h -1 -W` each lands on its own line, so + * the document is reassembled by concatenating the non-empty lines. No output (an empty result, or + * a pure DDL batch) means no rows. + */ +function parseJsonRows(stdout: string): Record[] { + const joined = stdout + .split(/\r?\n/) + .map((l) => l.trimEnd()) + .filter((l) => l.length > 0) + .join(""); + if (joined.length === 0) return []; + const parsed = JSON.parse(joined); + return Array.isArray(parsed) ? (parsed as Record[]) : [parsed as Record]; +} diff --git a/modules/mssql/index.ts b/modules/mssql/index.ts new file mode 100644 index 0000000..5547a8e --- /dev/null +++ b/modules/mssql/index.ts @@ -0,0 +1,25 @@ +// mssql's events entrypoint, loaded by the per-node tool host (the provisioner container runs +// ./provisioner separately). The database lifecycle events are EMITTED from the provisioner, where +// the lifecycle actually happens (novox/hq ADR 0046/0047): +// module.mssql.database.provisioned — a consumer's database + login/user was created +// module.mssql.database.deprovisioned — that database was removed +// Here in the tool host we react to them, keeping a lightweight audit trail of who was granted a +// database and who lost one — observability the provider itself is best placed to log. + +import { on } from "@novox/mesh-sdk/events"; + +interface DatabaseEvent { + consumer: string; + database: string; + user?: string; +} + +await on("module.mssql.database.provisioned", async (e) => { + console.log(`[mssql] database provisioned for ${e.body.consumer} (db ${e.body.database})`); +}); + +await on("module.mssql.database.deprovisioned", async (e) => { + console.log(`[mssql] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`); +}); + +console.log("[mssql] auditing database lifecycle events"); diff --git a/modules/mssql/module.json b/modules/mssql/module.json new file mode 100644 index 0000000..87d4708 --- /dev/null +++ b/modules/mssql/module.json @@ -0,0 +1,115 @@ +{ + "module": "mssql", + "version": "1", + "provides": [ + { + "name": "mssql-database", + "scope": "mesh" + } + ], + "capabilities": [ + "container-runtime" + ], + "emits": [ + "module.mssql.database.provisioned", + "module.mssql.database.deprovisioned" + ], + "consumes": [ + "module.mssql.database.provisioned", + "module.mssql.database.deprovisioned" + ], + "listens": [ + { + "port": 1433, + "protocol": "tcp", + "from": "mesh", + "why": "modules on any machine that were granted a database" + } + ], + "serves": { + "mssql-database": {} + }, + "receives": { + "mssql-database": "/var/lib/mssql/grants/mesh.json" + }, + "grants": { + "mssql-database": "/var/lib/mssql/grants" + }, + "own-secrets": { + "sa": "/var/lib/mssql/sa.secret", + "broker": "/var/lib/mesh/mssql/broker" + }, + "resources": [ + { + "id": "mesh-state", + "type": "directory", + "path": "/var/lib/mesh/mssql", + "mode": "0700" + }, + { + "id": "state", + "type": "directory", + "path": "/var/lib/mssql", + "mode": "0700" + }, + { + "id": "grants", + "type": "directory", + "path": "/var/lib/mssql/grants", + "mode": "0700" + }, + { + "id": "sa-env", + "type": "file", + "path": "/var/lib/mssql/sa.env", + "mode": "0600", + "content": "ACCEPT_EULA=Y\nMSSQL_SA_PASSWORD=${secret:sa}\n" + }, + { + "id": "data", + "type": "directory", + "path": "/services/mssql/db-data", + "mode": "0700", + "owner": "10001:0" + }, + { + "id": "net", + "type": "network", + "name": "mssql" + }, + { + "id": "server", + "type": "container", + "name": "mssql", + "image": "mcr.microsoft.com/mssql/server@sha256:ba4c8329f48fb8f02e1416be6a930ebfd71268caee78aa985f3af4315e457c89", + "network": "mssql", + "env-file": [ + "/var/lib/mssql/sa.env" + ], + "ports": [ + "1433" + ], + "volumes": [ + "/services/mssql/db-data:/var/opt/mssql" + ] + }, + { + "id": "runtime", + "type": "container", + "name": "mesh-mssql", + "image": "mesh-runtime-mssql@sha256:0000000000000000000000000000000000000000000000000000000000000000", + "network": "mssql", + "volumes": [ + "/var/lib/mesh/mssql/broker:/run/secrets/broker:ro", + "/var/lib/mssql/grants:/var/lib/mssql/grants:ro", + "/var/lib/mssql/sa.secret:/run/secrets/sa:ro" + ], + "env": { + "MESH_PROVISION_MSSQL": "mssql://sa@mssql:1433/master", + "MESH_PROVISION_PASSWORD_FILE": "/run/secrets/sa", + "MESH_BROKER_FILE": "/run/secrets/broker", + "MESH_RECEIVES": "/var/lib/mssql/grants/mesh.json" + } + } + ] +} diff --git a/modules/mssql/package.json b/modules/mssql/package.json new file mode 100644 index 0000000..d1a54ab --- /dev/null +++ b/modules/mssql/package.json @@ -0,0 +1,14 @@ +{ + "name": "@novox/module-mssql", + "version": "0.1.0", + "description": "mssql — provides the mesh mssql-database interface. Its client, provisioner, 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/mssql/provisioner/index.ts b/modules/mssql/provisioner/index.ts new file mode 100644 index 0000000..5827dc1 --- /dev/null +++ b/modules/mssql/provisioner/index.ts @@ -0,0 +1,47 @@ +// mssql's provisioner — the adapter that makes mssql a provider of the mesh `mssql-database` +// interface. The reconcile loop, the contributions file, and reading the mesh's minted password +// are the sdk harness's; this writes only the per-service half: how mssql creates and removes a +// consumer's database + login/user with the mesh-minted credential (novox/hq ADR 0044/0045/0053). +// +// The `mssql-database` interface: a consumer connects to a database it alone owns, as `as` with +// the password the mesh minted. +// +// **The login name and password are the mesh's, not the provisioner's (ADR 0053).** The mesh +// derives the login and hands it to both ends, and mints the password. mssql creates a login, a +// same-named database, and a db_owner user under exactly that login — a name the consumer cannot +// learn is a database it cannot reach. +// +// The DDL runs through MssqlClient, which is the module's one pending boundary (see client.ts). + +import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner"; +import { emit } from "@novox/mesh-sdk/events"; +import { MssqlClient } from "../client.js"; + +const mssql = MssqlClient.fromEnv(); + +/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */ +async function announce(type: string, body: Record): Promise { + try { + await emit(type, body); + } catch (err) { + console.error(`[provisioner:mssql-database] emit ${type} failed: ${err}`); + } +} + +runProvisioner("mssql-database", { + async create(p: Provision): Promise { + // Database, login and user share the consumer's name, so the consumer owns exactly its own. + const database = p.as; + await mssql.createDatabaseAndLogin(database, p.as, p.password); + await announce("module.mssql.database.provisioned", { + consumer: p.consumer ?? "", + database, + user: p.as, + }); + }, + + async remove(p: { as: string }): Promise { + await mssql.dropDatabaseAndLogin(p.as, p.as); + await announce("module.mssql.database.deprovisioned", { database: p.as }); + }, +}); diff --git a/modules/mssql/tools/index.ts b/modules/mssql/tools/index.ts new file mode 100644 index 0000000..2ec40f3 --- /dev/null +++ b/modules/mssql/tools/index.ts @@ -0,0 +1,44 @@ +// mssql's tools — mssql's own code (novox/hq ADR 0044), importing mssql's own client. They return +// structured data; the mesh serves them through the sdk's tool harness. Both call through +// MssqlClient, the module's one pending execution boundary (see client.ts): the tool shapes are +// fixed and correct, and surface the work honestly through that boundary. + +import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools"; +import { MssqlClient } from "../client.js"; + +export function getMssqlTools(mssql: MssqlClient): ToolDefinition[] { + return [ + { + name: "mssql_list_databases", + description: "List the user databases on the mssql server, with their on-disk size and state.", + input: {}, + run: async () => ({ databases: await mssql.listDatabases() }), + }, + { + name: "mssql_query", + description: "Run a read-only SELECT against a named database (wrapped in a rolled-back transaction).", + input: { + database: { type: "string", description: "the database to query" }, + sql: { type: "string", description: "a single SELECT statement" }, + }, + run: async (args) => { + const database = String(args.database ?? ""); + const sql = String(args.sql ?? ""); + if (!database) throw new Error("mssql_query: database is required"); + if (!sql) throw new Error("mssql_query: sql is required"); + const result = await mssql.readOnlyQuery(database, sql); + return { database, command: result.command, rows: result.rows }; + }, + }, + ]; +} + +// The tools exist only when the server can be reached from the environment; without it, mssql +// contributes none rather than failing the whole tool runtime. +registerModuleTools("mssql", (env) => { + try { + return getMssqlTools(MssqlClient.fromEnv(env)); + } catch { + return []; + } +}); diff --git a/modules/mssql/tsconfig.json b/modules/mssql/tsconfig.json new file mode 100644 index 0000000..51f4046 --- /dev/null +++ b/modules/mssql/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", "provisioner/index.ts", "tools/index.ts"] +}