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-<mod> 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
This commit is contained in:
@@ -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 <ApiKey> 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>([^<]+)<\/ApiKey>/);
|
||||
if (match) return match[1];
|
||||
}
|
||||
return null;
|
||||
}
|
||||
|
||||
private async get(endpoint: string, params?: Record<string, string>): Promise<unknown> {
|
||||
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<LidarrContentItem[]> {
|
||||
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<LidarrContentItem[]> {
|
||||
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<LidarrCalendarItem[]> {
|
||||
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]}`;
|
||||
}
|
||||
@@ -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<number, LidarrQueueItem>();
|
||||
let primed = false;
|
||||
|
||||
async function pollQueue(): Promise<void> {
|
||||
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<void>, 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");
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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,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<T>(js: string): Promise<T> {
|
||||
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<void> {
|
||||
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<void> {
|
||||
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<DatabaseInfo[]> {
|
||||
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<Record<string, unknown>>,
|
||||
limit: number,
|
||||
): Promise<Record<string, unknown>[]> {
|
||||
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<Record<string, unknown>[]>(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;
|
||||
}
|
||||
}
|
||||
@@ -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<DatabaseEvent>("module.mongodb.database.provisioned", async (e) => {
|
||||
console.log(`[mongodb] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
|
||||
});
|
||||
|
||||
await on<DatabaseEvent>("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");
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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<string, string>): Promise<void> {
|
||||
try {
|
||||
await emit(type, body);
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:mongodb-database] emit ${type} failed: ${err}`);
|
||||
}
|
||||
}
|
||||
|
||||
runProvisioner("mongodb-database", {
|
||||
async create(p: Provision): Promise<void> {
|
||||
// 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<void> {
|
||||
await mongo.dropDatabaseAndUser(p.as, p.as);
|
||||
await announce("module.mongodb.database.deprovisioned", { database: p.as });
|
||||
},
|
||||
});
|
||||
@@ -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<string, unknown> {
|
||||
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 [];
|
||||
}
|
||||
});
|
||||
@@ -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"]
|
||||
}
|
||||
@@ -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 <args>` 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<string> {
|
||||
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<boolean> {
|
||||
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 `<prefix>/#` 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<void> {
|
||||
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<void> {
|
||||
await ignoreMissing(this.ctl("deleteClient", username));
|
||||
await ignoreMissing(this.ctl("deleteRole", username));
|
||||
}
|
||||
|
||||
/** The dynsec client list, parsed from `listClients`. */
|
||||
async listClients(): Promise<string[]> {
|
||||
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<void> {
|
||||
// `dynsec init <file> <admin-username> [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<string>): Promise<void> {
|
||||
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<string>): Promise<void> {
|
||||
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;
|
||||
}
|
||||
}
|
||||
@@ -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<TopicEvent>("module.mosquitto.topic.provisioned", async (e) => {
|
||||
console.log(`[mosquitto] topic provisioned for ${e.body.consumer} (client ${e.body.username})`);
|
||||
});
|
||||
|
||||
await on<TopicEvent>("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");
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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 `<as>/#`, 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<string, string>): Promise<void> {
|
||||
try {
|
||||
await emit(type, body);
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:mqtt-topic] emit ${type} failed: ${err}`);
|
||||
}
|
||||
}
|
||||
|
||||
runProvisioner("mqtt-topic", {
|
||||
async create(p: Provision): Promise<void> {
|
||||
// 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<void> {
|
||||
await mosquitto.deleteScopedClient(p.as);
|
||||
await announce("module.mosquitto.topic.deprovisioned", { username: p.as });
|
||||
},
|
||||
});
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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"]
|
||||
}
|
||||
@@ -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<string, unknown>[];
|
||||
}
|
||||
|
||||
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<void> {
|
||||
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<Record<string, unknown>[]> {
|
||||
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<string> {
|
||||
// `-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<void> {
|
||||
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<void> {
|
||||
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<QueryResult> {
|
||||
// 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<string, unknown>[] {
|
||||
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<string, unknown>[]) : [parsed as Record<string, unknown>];
|
||||
}
|
||||
@@ -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<DatabaseEvent>("module.mssql.database.provisioned", async (e) => {
|
||||
console.log(`[mssql] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
|
||||
});
|
||||
|
||||
await on<DatabaseEvent>("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");
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
]
|
||||
}
|
||||
@@ -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"
|
||||
}
|
||||
}
|
||||
@@ -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<string, string>): Promise<void> {
|
||||
try {
|
||||
await emit(type, body);
|
||||
} catch (err) {
|
||||
console.error(`[provisioner:mssql-database] emit ${type} failed: ${err}`);
|
||||
}
|
||||
}
|
||||
|
||||
runProvisioner("mssql-database", {
|
||||
async create(p: Provision): Promise<void> {
|
||||
// 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<void> {
|
||||
await mssql.dropDatabaseAndLogin(p.as, p.as);
|
||||
await announce("module.mssql.database.deprovisioned", { database: p.as });
|
||||
},
|
||||
});
|
||||
@@ -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 [];
|
||||
}
|
||||
});
|
||||
@@ -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"]
|
||||
}
|
||||
Reference in New Issue
Block a user