Merge branch 'main' into feat/the-console

This commit is contained in:
2026-09-30 14:47:11 +00:00
22 changed files with 1126 additions and 172 deletions
+47 -22
View File
@@ -1,6 +1,26 @@
{ {
"module": "icecast", "module": "icecast",
"version": "1", "version": "1",
"requires": [
"route",
"secret"
],
"contributes": {
"route": {
"label": "icecast",
"endpoint": "stream"
}
},
"binds": {
"route": "${dir:state}/route.json"
},
"secrets": {
"secret": {
"source": "${dir:state}/source.secret",
"admin": "${dir:state}/admin.secret",
"relay": "${dir:state}/relay.secret"
}
},
"capabilities": [ "capabilities": [
"container-runtime" "container-runtime"
], ],
@@ -17,7 +37,7 @@
"port": 8000, "port": 8000,
"protocol": "tcp", "protocol": "tcp",
"from": "mesh", "from": "mesh",
"why": "streams in from sources and out to listeners" "why": "streams in from sources (HTTP PUT) and out to listeners, plus the status and admin pages; a public name is its route"
} }
], ],
"resources": [ "resources": [
@@ -30,28 +50,43 @@
{ {
"id": "state", "id": "state",
"type": "directory", "type": "directory",
"path": "/var/lib/icecast-module", "mode": "0700",
"mode": "0700" "place": "."
}, },
{ {
"id": "server-env", "id": "logs",
"type": "directory",
"mode": "0700",
"owner": "100:101"
},
{
"id": "server-conf",
"type": "file", "type": "file",
"path": "/var/lib/icecast-module/server.env", "path": "${dir:state}/icecast.xml",
"mode": "0600", "mode": "0600",
"content": "ICECAST_SOURCE_PASSWORD=${secret:source}\nICECAST_ADMIN_PASSWORD=${secret:admin}\nICECAST_RELAY_PASSWORD=${secret:relay}\nICECAST_ADMIN_USERNAME=admin\n" "content": "<icecast>\n <!-- Written by the mesh (modules/icecast). Passwords arrive as secrets rendered into this file,\n never as environment: the image's entrypoint seds ICECAST_* variables into the file only\n when they are set, and none are. -->\n <location>Earth</location>\n <admin>icemaster@localhost</admin>\n <limits>\n <clients>100</clients>\n <sources>2</sources>\n <queue-size>524288</queue-size>\n <client-timeout>30</client-timeout>\n <header-timeout>15</header-timeout>\n <source-timeout>10</source-timeout>\n <burst-on-connect>1</burst-on-connect>\n <burst-size>65535</burst-size>\n </limits>\n <authentication>\n <source-password>${secret:source}</source-password>\n <relay-password>${secret:relay}</relay-password>\n <admin-user>admin</admin-user>\n <admin-password>${secret:admin}</admin-password>\n </authentication>\n <!-- The name icecast writes into playlists (.m3u/.xspf: http://<hostname>:<port>/<mount>) and\n would announce to YP (none configured). A machine's own name belongs to its assignment, and\n an assignment merges only into JSON; this XML cannot take it, so the neutral default stays. -->\n <hostname>localhost</hostname>\n <listen-socket>\n <port>8000</port>\n </listen-socket>\n <http-headers>\n <header name=\"Access-Control-Allow-Origin\" value=\"*\" />\n </http-headers>\n <fileserve>1</fileserve>\n <paths>\n <basedir>/usr/share/icecast</basedir>\n <logdir>/var/log/icecast</logdir>\n <webroot>/usr/share/icecast/web</webroot>\n <adminroot>/usr/share/icecast/admin</adminroot>\n <alias source=\"/\" destination=\"/status.xsl\"/>\n </paths>\n <logging>\n <accesslog>access.log</accesslog>\n <errorlog>error.log</errorlog>\n <loglevel>3</loglevel>\n <logsize>10000</logsize>\n </logging>\n <security>\n <chroot>0</chroot>\n <!-- Starts as root, reads this 0600 root-owned file, then drops to the image's icecast user\n (uid 100, group icecast 101) before serving. -->\n <changeowner>\n <user>icecast</user>\n <group>icecast</group>\n </changeowner>\n </security>\n</icecast>\n"
},
{
"id": "net",
"type": "network",
"name": "icecast"
}, },
{ {
"id": "server", "id": "server",
"type": "container", "type": "container",
"name": "icecast", "name": "icecast",
"image": "infiniteproject/icecast@sha256:cd506cf3dfe31ce05fd37d7e672dbd1213e7255cc93d28ecf5a3b547af4e162c", "image": "infiniteproject/icecast@sha256:cd506cf3dfe31ce05fd37d7e672dbd1213e7255cc93d28ecf5a3b547af4e162c",
"env-file": [ "network": "icecast",
"/var/lib/icecast-module/server.env"
],
"ports": [ "ports": [
"8000" "8000"
], ],
"secrets-in-environment": "the image seds ICECAST_*_PASSWORD into icecast.xml and has no _FILE; convertible by mounting a generated icecast.xml, not yet done" "volumes": [
"${dir:state}/icecast.xml:/etc/icecast.xml:ro",
"${dir:logs}:/var/log/icecast"
],
"restart-on": [
"server-conf"
]
}, },
{ {
"id": "runtime-config", "id": "runtime-config",
@@ -65,14 +100,14 @@
"id": "runtime", "id": "runtime",
"type": "container", "type": "container",
"name": "mesh-icecast", "name": "mesh-icecast",
"network": "host", "network": "icecast",
"volumes": [ "volumes": [
"/var/lib/mesh/icecast/broker:/run/secrets/broker:ro", "/var/lib/mesh/icecast/broker:/run/secrets/broker:ro",
"/var/lib/mesh/icecast/config.json:/run/config/config.json:ro" "/var/lib/mesh/icecast/config.json:/run/config/config.json:ro"
], ],
"env": { "env": {
"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_ICECAST_URL": "http://127.0.0.1:8000", "MESH_ICECAST_URL": "http://icecast:8000",
"MESH_ICECAST_CONFIG_FILE": "/run/config/config.json" "MESH_ICECAST_CONFIG_FILE": "/run/config/config.json"
}, },
"restart-on": [ "restart-on": [
@@ -101,15 +136,5 @@
"from": "Dockerfile" "from": "Dockerfile"
} }
] ]
},
"requires": [
"secret"
],
"secrets": {
"secret": {
"source": "/var/lib/icecast-module/source.secret",
"admin": "/var/lib/icecast-module/admin.secret",
"relay": "/var/lib/icecast-module/relay.secret"
}
} }
} }
+2 -2
View File
@@ -13,7 +13,7 @@ ARG RUNTIME_BASE
FROM ${BUILD_BASE} AS build FROM ${BUILD_BASE} AS build
WORKDIR /app/modules/influxdb WORKDIR /app/modules/influxdb
COPY . . COPY . .
RUN node /app/node_modules/typescript/bin/tsc client.ts tools/index.ts \ RUN node /app/node_modules/typescript/bin/tsc client.ts grants.ts provisioner/index.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE} FROM ${RUNTIME_BASE}
@@ -21,4 +21,4 @@ COPY --from=build /app/modules/influxdb/dist /app/modules/influxdb/dist
# Every serve-time entrypoint, loaded by the runtime in serve mode: tools and events serve, and a # Every serve-time entrypoint, loaded by the runtime in serve mode: tools and events serve, and a
# provider's provisioner runs its reconcile loop in the same process, with the broker connected — # provider's provisioner runs its reconcile loop in the same process, with the broker connected —
# the convention novox/hq issues 060/061 settled. # the convention novox/hq issues 060/061 settled.
ENV MESH_TOOL_MODULES=/app/modules/influxdb/dist/tools/index.js ENV MESH_TOOL_MODULES=/app/modules/influxdb/dist/tools/index.js,/app/modules/influxdb/dist/provisioner/index.js
+118 -3
View File
@@ -17,6 +17,25 @@ export interface InfluxBucket {
retentionSeconds?: number; retentionSeconds?: number;
} }
/** One permission of an authorization, as InfluxDB represents it: an action on a resource type,
* in one org, optionally narrowed to one resource by id (no id = every resource of that type). */
export interface InfluxPermission {
action: "read" | "write";
resource: { type: string; orgID?: string; id?: string; name?: string; org?: string };
}
/** A v1-compatibility ("legacy") authorization: a username (InfluxDB calls it `token`) and a
* password the caller chooses, scoped by permissions. The one credential InfluxDB 2.x lets a
* caller set to a value it did not generate — which is what a mesh-minted password needs. */
export interface LegacyAuthorization {
id: string;
token: string;
orgID: string;
status?: "active" | "inactive";
description?: string;
permissions: InfluxPermission[];
}
/** The settings-merged config the mesh delivers (novox/hq ADR 0046): { url, apiKey, token, password, user, ... }. */ /** The settings-merged config the mesh delivers (novox/hq ADR 0046): { url, apiKey, token, password, user, ... }. */
function meshConfig(file?: string): Record<string, string> { function meshConfig(file?: string): Record<string, string> {
if (!file) return {}; if (!file) return {};
@@ -24,13 +43,20 @@ function meshConfig(file?: string): Record<string, string> {
catch { return {}; } catch { return {}; }
} }
/** A secret delivered as a file, trimmed; undefined when there is none, so the caller can fall back. */
function tokenFromFile(file?: string): string | undefined {
if (!file) return undefined;
try { return readFileSync(file, "utf8").trim() || undefined; }
catch { return undefined; }
}
export class InfluxDBClient { export class InfluxDBClient {
readonly baseUrl: string; readonly baseUrl: string;
constructor( constructor(
url: string, url: string,
private readonly token: string, private readonly token: string,
private readonly org: string, readonly org: string,
) { ) {
this.baseUrl = url.replace(/\/$/, ""); this.baseUrl = url.replace(/\/$/, "");
} }
@@ -43,8 +69,10 @@ export class InfluxDBClient {
static fromEnv(env: NodeJS.ProcessEnv = process.env): InfluxDBClient { static fromEnv(env: NodeJS.ProcessEnv = process.env): InfluxDBClient {
const cfg = meshConfig(env.MESH_INFLUXDB_CONFIG_FILE); const cfg = meshConfig(env.MESH_INFLUXDB_CONFIG_FILE);
const url = cfg.url ?? env.MESH_INFLUXDB_URL ?? `http://127.0.0.1:${env.INFLUXDB_PORT ?? "8086"}`; const url = cfg.url ?? env.MESH_INFLUXDB_URL ?? `http://127.0.0.1:${env.INFLUXDB_PORT ?? "8086"}`;
const token = cfg.token ?? env.MESH_INFLUXDB_TOKEN; // The token reaches the process as a file (novox/hq ADR 0086); the environment variable stays
if (!token) throw new Error("no InfluxDB token — set MESH_INFLUXDB_TOKEN"); // only for a workstation running the tools by hand.
const token = cfg.token ?? tokenFromFile(env.MESH_INFLUXDB_TOKEN_FILE) ?? env.MESH_INFLUXDB_TOKEN;
if (!token) throw new Error("no InfluxDB token — set MESH_INFLUXDB_TOKEN_FILE");
const org = cfg.org ?? env.MESH_INFLUXDB_ORG ?? "mesh"; const org = cfg.org ?? env.MESH_INFLUXDB_ORG ?? "mesh";
return new InfluxDBClient(url, token, org); return new InfluxDBClient(url, token, org);
} }
@@ -61,6 +89,93 @@ export class InfluxDBClient {
return res; return res;
} }
/** Like request, but the answer is returned whatever its status, for the caller to read. */
private async raw(path: string, init?: RequestInit): Promise<Response> {
return fetch(`${this.baseUrl}${path}`, {
...init,
headers: { Authorization: `Token ${this.token}`, ...(init?.headers ?? {}) },
});
}
private async send(path: string, method: string, body?: unknown): Promise<Response> {
return this.request(path, {
method,
headers: { "Content-Type": "application/json" },
body: body === undefined ? undefined : JSON.stringify(body),
});
}
/** The id of the org of this name, or undefined when there is none. */
async orgID(name: string): Promise<string | undefined> {
const res = await this.raw(`/api/v2/orgs?org=${encodeURIComponent(name)}`);
if (res.status === 404) return undefined;
if (!res.ok) throw new Error(`InfluxDB API /api/v2/orgs: ${res.status} ${await res.text()}`);
const body = (await res.json()) as { orgs?: { id: string; name: string }[] };
return body.orgs?.find((o) => o.name === name)?.id;
}
/** The bucket of exactly this name in the org, or undefined. */
async findBucket(orgID: string, name: string): Promise<InfluxBucket | undefined> {
const res = await this.raw(`/api/v2/buckets?orgID=${encodeURIComponent(orgID)}&name=${encodeURIComponent(name)}`);
if (res.status === 404) return undefined;
if (!res.ok) throw new Error(`InfluxDB API /api/v2/buckets: ${res.status} ${await res.text()}`);
const body = (await res.json()) as { buckets?: { id: string; name: string; orgID?: string }[] };
const b = body.buckets?.find((x) => x.name === name);
return b ? { id: b.id, name: b.name, orgID: b.orgID } : undefined;
}
/** Create a bucket that keeps its data for ever — retention is the operator's choice, never the mesh's. */
async createBucket(orgID: string, name: string, description: string): Promise<InfluxBucket> {
const b = (await (await this.send("/api/v2/buckets", "POST", {
orgID, name, description, retentionRules: [],
})).json()) as { id: string; name: string; orgID?: string };
return { id: b.id, name: b.name, orgID: b.orgID };
}
/** The v1 authorization whose username is exactly this, or undefined. */
async findLegacy(username: string): Promise<LegacyAuthorization | undefined> {
const path = `/private/legacy/authorizations?token=${encodeURIComponent(username)}`;
const res = await this.raw(path);
// InfluxDB answers a filter matching nothing with 404, not an empty list.
if (res.status === 404) return undefined;
if (!res.ok) throw new Error(`InfluxDB API ${path}: ${res.status} ${await res.text()}`);
const body = (await res.json()) as { authorizations?: LegacyAuthorization[] };
return body.authorizations?.find((a) => a.token === username);
}
async createLegacy(a: Omit<LegacyAuthorization, "id">): Promise<LegacyAuthorization> {
return (await (await this.send("/private/legacy/authorizations", "POST", a)).json()) as LegacyAuthorization;
}
/** Set a v1 authorization's password. InfluxDB keeps only a hash of it, so it can be set, never read. */
async setLegacyPassword(id: string, password: string): Promise<void> {
await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}/password`, "POST", { password });
}
async updateLegacy(id: string, patch: { status?: "active" | "inactive"; description?: string }): Promise<void> {
await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}`, "PATCH", patch);
}
async deleteLegacy(id: string): Promise<void> {
await this.send(`/private/legacy/authorizations/${encodeURIComponent(id)}`, "DELETE");
}
/**
* Whether this username and password sign in on the v1 API — the consumer's own view. Asked with
* a statement that reads nothing (`SHOW DATABASES` lists only what the credential may read), sent
* with Basic auth so the password is never in a URL. 401 is a wrong password or no such user;
* anything else that is not a server error means InfluxDB knew who was asking.
*/
async legacySignsIn(username: string, password: string): Promise<boolean> {
const res = await fetch(`${this.baseUrl}/query?q=${encodeURIComponent("SHOW DATABASES")}`, {
headers: { Authorization: `Basic ${Buffer.from(`${username}:${password}`).toString("base64")}` },
});
await res.arrayBuffer();
if (res.status === 401) return false;
if (res.status >= 500) throw new Error(`InfluxDB v1 /query: ${res.status}`);
return true;
}
/** Server health — the one endpoint that needs no token, but we send it anyway. */ /** Server health — the one endpoint that needs no token, but we send it anyway. */
async health(): Promise<InfluxHealth> { async health(): Promise<InfluxHealth> {
return (await (await this.request("/health")).json()) as InfluxHealth; return (await (await this.request("/health")).json()) as InfluxHealth;
+186
View File
@@ -0,0 +1,186 @@
// What the `influxdb-api` provision means in InfluxDB: one v1-compatibility authorization per
// consumer, in the org this module serves, under the username and password the mesh gave both ends,
// allowed exactly the access the consumer contributed. The provisioner (provisioner/index.ts) is the
// sdk harness calling these; they are here, apart from it, so they can be exercised against a fake
// InfluxDB without a broker or a contributions file.
//
// **Why a v1 authorization and not a v2 API token.** The mesh mints the consumer's password and
// hands it to both ends (novox/hq ADR 0048); the provider sets it, and never hands one back. An
// InfluxDB 2.x API token is generated by the server — `POST /api/v2/authorizations` ignores a token
// the caller sends — so a token could only ever be the operator's to accept, one per pair, by hand.
// A v1 authorization is a username and a password the caller chooses (8–72 characters; the mesh
// mints 40), stored hashed, and it reads and writes through InfluxQL (`/query`) and line protocol
// (`/write`), which every bucket answers under its own name as a database (InfluxDB maps each
// bucket to a database of the same name by itself). That is what grafana's InfluxDB data source
// speaks, and what Node-RED's influxdb nodes speak in their 1.x mode — so the mesh can make every
// consumer's credential, rotate it and withdraw it, with no person in the loop.
//
// **What a consumer contributes.** `access`: "read" (the default), "write" or "read-write".
// `buckets`: the buckets it may use, by name. A reader that names none may read every bucket of the
// org — a dashboard is pointed at data, it does not own it. A writer must name its buckets: writing
// everywhere, the org's system buckets included, is never what a consumer means. A named bucket
// that does not exist is created, keeping its data for ever; the mesh never deletes a bucket.
//
// **Only what the mesh made is touched.** An authorization this module creates is named with the
// mesh's identity prefix and its description starts with MARK. One with the same username that
// lacks the mark is somebody else's: it is refused, never adopted, never updated, never deleted.
// Every other authorization, token, user and bucket in the instance is left exactly as it was.
import type { InfluxDBClient, InfluxPermission, LegacyAuthorization } from "./client.js";
/** How a description marks an authorization as the mesh's own work. */
export const MARK = "[mesh]";
/** The prefix the mesh gives every consumer identity (novox/hq ADR 0049). */
const IDENTITY_PREFIX = "mesh_";
/** One consumer, as the harness hands it over. */
export interface ApiGrant {
readonly as: string;
readonly password: string;
readonly values: Readonly<Record<string, unknown>>;
readonly consumer?: string;
}
export type Access = "read" | "write" | "read-write";
/** What a contribution asks for, checked. Refused when it cannot be served as asked. */
export function askedFor(values: Readonly<Record<string, unknown>>): { access: Access; buckets: string[] } {
const access = values.access ?? "read";
if (access !== "read" && access !== "write" && access !== "read-write") {
throw new Error(`contributes an access of ${JSON.stringify(access)} — it is "read", "write" or "read-write"`);
}
const raw = values.buckets ?? [];
if (!Array.isArray(raw) || raw.some((b) => typeof b !== "string" || b.trim() === "")) {
throw new Error(`contributes buckets of ${JSON.stringify(raw)} — a list of bucket names`);
}
const buckets = [...new Set((raw as string[]).map((b) => b.trim()))].sort();
if (access !== "read" && buckets.length === 0) {
throw new Error(`asks to write and names no bucket (\`buckets\`) — a writer names what it writes to`);
}
if (buckets.some((b) => b.startsWith("_"))) {
throw new Error(`names a system bucket (${buckets.filter((b) => b.startsWith("_")).join(", ")}) — those are InfluxDB's own`);
}
return { access: access as Access, buckets };
}
/** The permissions a grant resolves to, given each named bucket's id. */
export function permissionsFor(orgID: string, access: Access, bucketIDs: string[]): InfluxPermission[] {
const actions: ("read" | "write")[] = access === "read-write" ? ["read", "write"] : [access];
const out: InfluxPermission[] = [];
for (const action of actions) {
if (bucketIDs.length === 0) {
out.push({ action, resource: { type: "buckets", orgID } });
continue;
}
for (const id of bucketIDs) out.push({ action, resource: { type: "buckets", orgID, id } });
}
return out;
}
/** A permission as a comparable string: what InfluxDB answers carries names and links besides. */
function key(p: InfluxPermission): string {
return `${p.action}:${p.resource.type}:${p.resource.orgID ?? ""}:${p.resource.id ?? "*"}`;
}
function samePermissions(a: readonly InfluxPermission[], b: readonly InfluxPermission[]): boolean {
const x = a.map(key).sort();
const y = b.map(key).sort();
return x.length === y.length && x.every((v, i) => v === y[i]);
}
export function marked(a: Pick<LegacyAuthorization, "token" | "description">): boolean {
return a.token.startsWith(IDENTITY_PREFIX) && (a.description ?? "").startsWith(MARK);
}
function describe(g: ApiGrant): string {
return `${MARK} made by the mesh for ${g.consumer ? `a module on ${g.consumer}` : "a consumer"} — do not edit; it is reset`;
}
export class ApiGrants {
constructor(private readonly influx: InfluxDBClient, readonly org: string) {}
private async orgID(): Promise<string> {
const id = await this.influx.orgID(this.org);
if (!id) throw new Error(`InfluxDB has no org ${JSON.stringify(this.org)} — the org this module serves must exist`);
return id;
}
/** The ids of the named buckets, creating any that are missing when `create` says so. Undefined
* when one is missing and may not be created (a read-only question). */
private async bucketIDs(orgID: string, names: string[], create: ApiGrant | undefined): Promise<string[] | undefined> {
const ids: string[] = [];
for (const name of names) {
let b = await this.influx.findBucket(orgID, name);
if (!b) {
if (!create) return undefined;
b = await this.influx.createBucket(orgID, name, `${MARK} made by the mesh for ${create.as}; the mesh never deletes it`);
}
ids.push(b.id);
}
return ids.sort();
}
/** Create the consumer's authorization, or bring the mesh's existing one back to what the grant
* says. Idempotent: a second apply of the same grant changes nothing beyond re-asserting the
* password, which InfluxDB can be told but never asked. */
async ensure(g: ApiGrant): Promise<"created" | "updated" | "unchanged"> {
if (!g.as.startsWith(IDENTITY_PREFIX)) {
throw new Error(`${g.as} is not a mesh identity — the mesh names every consumer ${IDENTITY_PREFIX}<node>_<module>`);
}
const { access, buckets } = askedFor(g.values);
const orgID = await this.orgID();
const found = await this.influx.findLegacy(g.as);
if (found && !marked(found)) {
throw new Error(
`InfluxDB already has a v1 authorization ${g.as} the mesh did not make — left alone; ` +
`delete it if the mesh should own that name`);
}
const want = permissionsFor(orgID, access, (await this.bucketIDs(orgID, buckets, g))!);
if (found && found.orgID === orgID && samePermissions(found.permissions, want)) {
// Only what differs is written. The password cannot be read back, so it is tried instead.
let changed = false;
if (found.status === "inactive") {
await this.influx.updateLegacy(found.id, { status: "active" });
changed = true;
}
if (!(await this.influx.legacySignsIn(g.as, g.password))) {
await this.influx.setLegacyPassword(found.id, g.password);
changed = true;
}
return changed ? "updated" : "unchanged";
}
// InfluxDB cannot change an authorization's permissions in place, so the mesh's own is made
// again. Only ever one the mesh made: a foreign one was refused above.
if (found) await this.influx.deleteLegacy(found.id);
const made = await this.influx.createLegacy({
token: g.as, orgID, status: "active", description: describe(g), permissions: want,
});
await this.influx.setLegacyPassword(made.id, g.password);
return found ? "updated" : "created";
}
/** Whether InfluxDB still holds this consumer's authorization exactly as the grant says: present,
* the mesh's, active, allowed what was asked and nothing more, and signing in with the mesh's
* password. Reads only — a missing bucket is "not held", never created here. */
async holds(g: ApiGrant): Promise<boolean> {
const { access, buckets } = askedFor(g.values);
const orgID = await this.influx.orgID(this.org);
if (!orgID) return false;
const found = await this.influx.findLegacy(g.as);
if (!found || !marked(found) || found.status === "inactive" || found.orgID !== orgID) return false;
const ids = await this.bucketIDs(orgID, buckets, undefined);
if (!ids || !samePermissions(found.permissions, permissionsFor(orgID, access, ids))) return false;
return this.influx.legacySignsIn(g.as, g.password);
}
/** Withdraw a consumer's authorization — only one the mesh made. Its buckets and their data stay. */
async remove(as: string): Promise<"removed" | "absent" | "not ours"> {
const found = await this.influx.findLegacy(as);
if (!found) return "absent";
if (!marked(found)) return "not ours";
await this.influx.deleteLegacy(found.id);
return "removed";
}
}
+58 -30
View File
@@ -1,11 +1,19 @@
{ {
"module": "influxdb", "module": "influxdb",
"version": "1", "version": "1",
"provides": [
{
"name": "influxdb-api",
"scope": "mesh"
}
],
"capabilities": [ "capabilities": [
"container-runtime" "container-runtime"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/influxdb/broker" "broker": "/var/lib/mesh/influxdb/broker",
"admin": "${dir:state}/admin.secret",
"admin-token": "${dir:state}/admin-token.secret"
}, },
"listens": [ "listens": [
{ {
@@ -13,9 +21,23 @@
"port": 8086, "port": 8086,
"protocol": "tcp", "protocol": "tcp",
"from": "mesh", "from": "mesh",
"why": "queries and writes, over http" "why": "queries, writes and the web UI, over http; consumers granted influxdb-api sign in with the mesh's credential, and a name is a route grant"
} }
], ],
"serves": {
"influxdb-api": {
"scheme": "http",
"port": 8086,
"org": "mesh",
"bucket": "default"
}
},
"receives": {
"influxdb-api": "${dir:grants}/mesh.json"
},
"grants": {
"influxdb-api": "${dir:grants}"
},
"resources": [ "resources": [
{ {
"id": "mesh-state", "id": "mesh-state",
@@ -26,46 +48,50 @@
{ {
"id": "state", "id": "state",
"type": "directory", "type": "directory",
"path": "/var/lib/influxdb-module", "mode": "0700",
"mode": "0700" "place": "."
},
{
"id": "server-env",
"type": "file",
"path": "/var/lib/influxdb-module/server.env",
"mode": "0600",
"content": "DOCKER_INFLUXDB_INIT_MODE=setup\nDOCKER_INFLUXDB_INIT_USERNAME=admin\nDOCKER_INFLUXDB_INIT_PASSWORD=${secret:admin}\nDOCKER_INFLUXDB_INIT_ADMIN_TOKEN=${secret:admin-token}\nDOCKER_INFLUXDB_INIT_ORG=mesh\nDOCKER_INFLUXDB_INIT_BUCKET=default\n"
}, },
{ {
"id": "data", "id": "data",
"type": "directory", "type": "directory",
"path": "/services/influxdb/data",
"mode": "0700", "mode": "0700",
"owner": "1000:1000" "owner": "1000:1000"
}, },
{ {
"id": "config", "id": "config",
"type": "directory", "type": "directory",
"path": "/services/influxdb/config",
"mode": "0700", "mode": "0700",
"owner": "1000:1000" "owner": "1000:1000"
}, },
{
"id": "grants",
"type": "directory",
"mode": "0700"
},
{
"id": "server-env",
"type": "file",
"path": "${dir:state}/server.env",
"mode": "0600",
"content": "DOCKER_INFLUXDB_INIT_MODE=setup\nDOCKER_INFLUXDB_INIT_USERNAME=admin\nDOCKER_INFLUXDB_INIT_PASSWORD_FILE=/run/secrets/admin\nDOCKER_INFLUXDB_INIT_ADMIN_TOKEN_FILE=/run/secrets/admin-token\nDOCKER_INFLUXDB_INIT_ORG=mesh\nDOCKER_INFLUXDB_INIT_BUCKET=default\n"
},
{ {
"id": "server", "id": "server",
"type": "container", "type": "container",
"name": "influxdb", "name": "influxdb",
"image": "influxdb@sha256:f75e48af0598e8aec7986e991a848d19a119101a7d563a2e5db1dfaac9c45daa", "image": "influxdb@sha256:f75e48af0598e8aec7986e991a848d19a119101a7d563a2e5db1dfaac9c45daa",
"env-file": [ "env-file": [
"/var/lib/influxdb-module/server.env" "${dir:state}/server.env"
], ],
"ports": [ "ports": [
"8086" "8086"
], ],
"volumes": [ "volumes": [
"/services/influxdb/data:/var/lib/influxdb2", "${dir:data}:/var/lib/influxdb2",
"/services/influxdb/config:/etc/influxdb2" "${dir:config}:/etc/influxdb2",
], "${dir:state}/admin.secret:/run/secrets/admin:ro",
"secrets-in-environment": "the image honours DOCKER_INFLUXDB_INIT_PASSWORD_FILE and _ADMIN_TOKEN_FILE; convertible, awaiting a bed that proves it" "${dir:state}/admin-token.secret:/run/secrets/admin-token:ro"
]
}, },
{ {
"id": "runtime-config", "id": "runtime-config",
@@ -83,13 +109,15 @@
"volumes": [ "volumes": [
"/var/lib/mesh/influxdb/broker:/run/secrets/broker:ro", "/var/lib/mesh/influxdb/broker:/run/secrets/broker:ro",
"/var/lib/mesh/influxdb/config.json:/run/config/config.json:ro", "/var/lib/mesh/influxdb/config.json:/run/config/config.json:ro",
"/services/influxdb/config:/var/lib/influxdb/config:ro" "${dir:state}/admin-token.secret:/run/secrets/admin-token:ro",
"${dir:grants}:${dir:grants}:ro"
], ],
"env": { "env": {
"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_INFLUXDB_URL": "http://127.0.0.1:8086", "MESH_INFLUXDB_URL": "http://127.0.0.1:${port:8086}",
"MESH_INFLUXDB_CONFIG_FILE": "/run/config/config.json", "MESH_INFLUXDB_CONFIG_FILE": "/run/config/config.json",
"MESH_INFLUXDB_CONFIG_DIR": "/var/lib/influxdb/config" "MESH_INFLUXDB_TOKEN_FILE": "/run/secrets/admin-token",
"MESH_RECEIVES": "${dir:grants}/mesh.json"
}, },
"restart-on": [ "restart-on": [
"runtime-config" "runtime-config"
@@ -97,6 +125,15 @@
"artifact": "runtime" "artifact": "runtime"
} }
], ],
"requires": [
"route"
],
"contributes": {
"route": {
"label": "influxdb",
"endpoint": "api"
}
},
"build": { "build": {
"on": [ "on": [
{ {
@@ -117,14 +154,5 @@
"from": "Dockerfile" "from": "Dockerfile"
} }
] ]
},
"requires": [
"secret"
],
"secrets": {
"secret": {
"admin": "/var/lib/influxdb-module/admin.secret",
"admin-token": "/var/lib/influxdb-module/admin-token.secret"
}
} }
} }
+6 -1
View File
@@ -1,9 +1,14 @@
{ {
"name": "@novox/module-influxdb", "name": "@novox/module-influxdb",
"version": "0.1.0", "version": "0.1.0",
"description": "influxdb — time-series database. Its API client and tools live here (novox/hq ADR 0039).", "description": "influxdb — time-series database; provides the mesh influxdb-api interface. Its API client, provisioner and tools live here (novox/hq ADR 0039).",
"type": "module", "type": "module",
"private": true, "private": true,
"scripts": {
"build": "tsc client.ts grants.ts provisioner/index.ts tools/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist",
"typecheck": "tsc -p tsconfig.json",
"test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'"
},
"dependencies": { "dependencies": {
"@novox/mesh-sdk": "^0.1.0" "@novox/mesh-sdk": "^0.1.0"
}, },
+54
View File
@@ -0,0 +1,54 @@
// influxdb's provisioner — the adapter that makes influxdb a provider of the mesh `influxdb-api`
// interface. The reconcile loop, the contributions file and reading the mesh's minted secret are the
// sdk harness's; this writes only the per-service half: how InfluxDB creates, checks and removes a
// consumer's credential (novox/hq ADR 0039/0040/0048). What that credential is, and why it is a v1
// authorization, is in ../grants.ts.
//
// The `influxdb-api` interface: a consumer reaches `${bound:influxdb-api:scheme}://…:at:…:port`,
// signs in as `${bound:influxdb-api:as}` with the password the mesh minted for the pair, and reads
// or writes the org's buckets as databases of the same name — `${bound:influxdb-api:bucket}` being
// the one this instance serves by default. The org and the default bucket are the assignment's
// settings, which reach both what is served and this module's config.json, so the org a consumer is
// told and the org its credential is made in cannot disagree.
import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner";
import { InfluxDBClient } from "../client.js";
import { ApiGrants } from "../grants.js";
let grants: ApiGrants | undefined;
try {
const influx = InfluxDBClient.fromEnv();
grants = new ApiGrants(influx, influx.org);
} catch (err) {
// No admin token: nothing can be provisioned, and the tools loaded beside this must still serve.
console.error(`[provisioner:influxdb-api] not started: ${err instanceof Error ? err.message : err}`);
}
if (grants) serve(grants);
function serve(grants: ApiGrants): void {
runProvisioner("influxdb-api", {
async create(p: Provision): Promise<void> {
const done = await grants.ensure(p);
if (done !== "unchanged") {
console.log(`[provisioner:influxdb-api] ${done} v1 authorization ${p.as} in org ${grants.org}`);
}
},
async remove(p: { as: string }): Promise<void> {
const done = await grants.remove(p.as);
if (done === "not ours") {
console.error(`[provisioner:influxdb-api] ${p.as}: an authorization of that name exists that the mesh did not make — left alone`);
} else if (done === "removed") {
console.log(`[provisioner:influxdb-api] removed v1 authorization ${p.as}; its buckets and their data stay`);
}
},
// Asked every minute by the harness: whether InfluxDB still holds this consumer's authorization
// exactly as the mesh gave it, so one deleted, disabled or re-passworded behind the mesh's back is
// made whole again (hq issue 120).
async holds(p: Provision): Promise<boolean> {
return grants.holds(p);
},
});
}
+246
View File
@@ -0,0 +1,246 @@
// What holds influxdb to the `influxdb-api` provision (grants.ts): one v1 authorization per consumer,
// under the username and password the mesh gave, allowed only what the consumer contributed; made
// once and brought back on every apply; buckets created when missing and never deleted; and an
// authorization the mesh did not make — same name or not — never adopted, changed or deleted.
//
// InfluxDB is a fake: the routes the module touches, answering with the status codes and shapes
// InfluxDB 2.9 gives (a filter matching nothing is a 404, a password outside 8–72 characters a 400,
// an inactive authorization or a wrong password a 401 on /query). Run against the compiled module
// (npm test builds first), the way the runtime loads it.
import { test, after, beforeEach } from "node:test";
import assert from "node:assert/strict";
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
import { InfluxDBClient } from "../dist/client.js";
import { ApiGrants, MARK, askedFor, marked } from "../dist/grants.js";
type Rec = Record<string, any>;
const ADMIN = "operator-token";
const orgs = new Map<string, string>([["zurag", "org1"]]);
let buckets: Rec[] = [];
let auths: Rec[] = [];
let calls: string[] = [];
let seq = 0;
function body(req: IncomingMessage): Promise<any> {
return new Promise((resolve) => {
let raw = "";
req.on("data", (c) => (raw += c));
req.on("end", () => resolve(raw ? JSON.parse(raw) : undefined));
});
}
function send(res: ServerResponse, status: number, value?: unknown): void {
res.writeHead(status, { "Content-Type": "application/json" });
res.end(value === undefined ? "" : JSON.stringify(value));
}
const server = createServer(async (req, res) => {
const url = new URL(req.url!, "http://fake");
const p = url.pathname;
calls.push(`${req.method} ${p}`);
if (p === "/query") {
const basic = (req.headers.authorization ?? "").replace(/^Basic /, "");
const [u, pw] = Buffer.from(basic, "base64").toString().split(":");
const a = auths.find((x) => x.token === u);
if (!a || a.status !== "active" || a.password === undefined || a.password !== pw) {
return send(res, 401, { code: "unauthorized", message: "Unauthorized" });
}
return send(res, 200, { results: [{ statement_id: 0 }] });
}
if (req.headers.authorization !== `Token ${ADMIN}`) return send(res, 401, { code: "unauthorized" });
if (p === "/api/v2/orgs") {
const id = orgs.get(url.searchParams.get("org") ?? "");
if (!id) return send(res, 404, { code: "not found", message: "organization name not found" });
return send(res, 200, { orgs: [{ id, name: url.searchParams.get("org") }] });
}
if (p === "/api/v2/buckets" && req.method === "GET") {
const found = buckets.filter((b) => b.orgID === url.searchParams.get("orgID") && b.name === url.searchParams.get("name"));
if (found.length === 0) return send(res, 404, { code: "not found", message: "bucket not found" });
return send(res, 200, { buckets: found });
}
if (p === "/api/v2/buckets" && req.method === "POST") {
const b = { ...(await body(req)), id: `b${++seq}` };
buckets.push(b);
return send(res, 201, b);
}
if (p === "/private/legacy/authorizations" && req.method === "GET") {
const found = auths.filter((a) => a.token === url.searchParams.get("token"));
if (found.length === 0) return send(res, 404, { code: "not found", message: "authorization not found" });
// Never answers with the password: InfluxDB keeps only its hash.
return send(res, 200, { authorizations: found.map(({ password, ...a }) => ({ ...a, links: {} })) });
}
if (p === "/private/legacy/authorizations" && req.method === "POST") {
const a = await body(req);
if (auths.some((x) => x.token === a.token)) return send(res, 409, { code: "conflict", message: "token already exists" });
const made = { ...a, id: `a${++seq}`, status: a.status ?? "active" };
auths.push(made);
return send(res, 201, made);
}
const m = /^\/private\/legacy\/authorizations\/([^/]+)(\/password)?$/.exec(p);
const a = m && auths.find((x) => x.id === m[1]);
if (!a) return send(res, 404, { code: "not found" });
if (m![2] && req.method === "POST") {
const { password } = await body(req);
if (typeof password !== "string" || password.length < 8 || password.length > 72) {
return send(res, 400, { code: "invalid", message: "passwords must be between 8 and 72 characters long" });
}
a.password = password;
return send(res, 204);
}
if (req.method === "PATCH") {
Object.assign(a, await body(req));
return send(res, 200, a);
}
if (req.method === "DELETE") {
auths = auths.filter((x) => x !== a);
return send(res, 204);
}
send(res, 405);
});
await new Promise<void>((r) => server.listen(0, "127.0.0.1", r));
after(() => server.close());
const port = (server.address() as { port: number }).port;
const grants = new ApiGrants(new InfluxDBClient(`http://127.0.0.1:${port}`, ADMIN, "zurag"), "zurag");
const PW = "mesh-minted-password-of-forty-characters";
/** Grafana on ace, as the mesh hands it to the provisioner. */
function grafana(password = PW, values: Record<string, unknown> = { access: "read" }) {
return { as: "mesh_ace_grafana", password, consumer: "ace", values };
}
/** Node-RED on ace: writes one bucket. */
function nodered(password = PW, values: Record<string, unknown> = { access: "write", buckets: ["zurag"] }) {
return { as: "mesh_ace_nodered", password, consumer: "ace", values };
}
function only(token: string): Rec {
const found = auths.filter((a) => a.token === token);
assert.equal(found.length, 1, `exactly one authorization ${token}, found ${found.length}`);
return found[0];
}
function perms(a: Rec): string[] {
return a.permissions.map((p: Rec) => `${p.action}:${p.resource.type}:${p.resource.id ?? "*"}`).sort();
}
beforeEach(() => {
buckets = [{ id: "zb", orgID: "org1", name: "zurag" }];
auths = [];
calls = [];
});
test("what a contribution may ask for, and what is refused", () => {
assert.deepEqual(askedFor({}), { access: "read", buckets: [] });
assert.deepEqual(askedFor({ access: "read-write", buckets: ["b", "a", "a"] }), { access: "read-write", buckets: ["a", "b"] });
assert.throws(() => askedFor({ access: "admin" }), /access/);
assert.throws(() => askedFor({ access: "write" }), /names no bucket/);
assert.throws(() => askedFor({ buckets: "zurag" }), /list of bucket names/);
assert.throws(() => askedFor({ access: "write", buckets: ["_monitoring"] }), /system bucket/);
});
test("a reader is given one authorization, reading every bucket of the org, under the mesh's password", async () => {
assert.equal(await grants.ensure(grafana()), "created");
const a = only("mesh_ace_grafana");
assert.equal(a.orgID, "org1");
assert.equal(a.status, "active");
assert.ok(a.description.startsWith(MARK));
assert.deepEqual(perms(a), ["read:buckets:*"]);
assert.equal(a.password, PW);
assert.equal(await grants.holds(grafana()), true);
});
test("a writer is allowed its own buckets only, and a missing one is made — never deleted", async () => {
assert.equal(await grants.ensure(nodered(PW, { access: "write", buckets: ["zurag", "printer"] })), "created");
const made = buckets.find((b) => b.name === "printer");
assert.ok(made, "the missing bucket was created");
assert.deepEqual(made!.retentionRules, [], "kept for ever: retention is the operator's choice");
assert.deepEqual(perms(only("mesh_ace_nodered")), [`write:buckets:${made!.id}`, "write:buckets:zb"]);
assert.equal(await grants.remove("mesh_ace_nodered"), "removed");
assert.equal(buckets.length, 2, "withdrawing the consumer leaves every bucket and its data");
});
test("applying the same grant again writes nothing", async () => {
await grants.ensure(grafana());
calls = [];
assert.equal(await grants.ensure(grafana()), "unchanged");
assert.ok(calls.every((c) => c.startsWith("GET")), `only reads: ${calls.join(", ")}`);
only("mesh_ace_grafana");
});
test("a rotated password is set in place; a changed access remakes only the mesh's own", async () => {
await grants.ensure(nodered());
const id = only("mesh_ace_nodered").id;
assert.equal(await grants.holds(nodered("rotated-password-0123456789")), false);
assert.equal(await grants.ensure(nodered("rotated-password-0123456789")), "updated");
assert.equal(only("mesh_ace_nodered").id, id, "updated, not replaced");
assert.equal(await grants.holds(nodered("rotated-password-0123456789")), true);
await grants.ensure(nodered(PW, { access: "read-write", buckets: ["zurag"] }));
assert.deepEqual(perms(only("mesh_ace_nodered")), ["read:buckets:zb", "write:buckets:zb"]);
assert.equal(await grants.holds(nodered(PW, { access: "read-write", buckets: ["zurag"] })), true);
});
test("an authorization disabled, re-passworded or deleted behind the mesh's back is not held, and is made whole", async () => {
await grants.ensure(grafana());
only("mesh_ace_grafana").status = "inactive";
assert.equal(await grants.holds(grafana()), false);
assert.equal(await grants.ensure(grafana()), "updated");
assert.equal(await grants.holds(grafana()), true);
only("mesh_ace_grafana").password = "somebody-else-set-this";
assert.equal(await grants.holds(grafana()), false);
await grants.ensure(grafana());
assert.equal(await grants.holds(grafana()), true);
auths = [];
assert.equal(await grants.holds(grafana()), false);
assert.equal(await grants.ensure(grafana()), "created");
});
test("holds only reads, and a bucket gone missing is not held rather than made", async () => {
await grants.ensure(nodered());
buckets = [];
calls = [];
assert.equal(await grants.holds(nodered()), false);
assert.ok(calls.every((c) => c.startsWith("GET")), `only reads: ${calls.join(", ")}`);
assert.equal(buckets.length, 0);
});
test("an authorization of the same name the mesh did not make is refused, and left exactly as it was", async () => {
auths = [{ id: "theirs", token: "mesh_ace_grafana", orgID: "org1", status: "active", description: "hand-made",
permissions: [{ action: "write", resource: { type: "buckets", orgID: "org1" } }], password: "their-password" }];
const before = JSON.stringify(auths);
await assert.rejects(grants.ensure(grafana()), /did not make/);
assert.equal(JSON.stringify(auths), before);
assert.ok(calls.every((c) => c.startsWith("GET")), `only reads: ${calls.join(", ")}`);
assert.equal(await grants.holds(grafana()), false);
assert.equal(await grants.remove("mesh_ace_grafana"), "not ours");
assert.equal(auths.length, 1, "never deleted");
});
test("the predecessor's own v1 users and tokens are never touched", async () => {
auths = [{ id: "hal", token: "grafana", orgID: "org1", status: "active", description: "",
permissions: [{ action: "read", resource: { type: "buckets", orgID: "org1" } }], password: "old-password" }];
await grants.ensure(grafana());
assert.equal(auths.find((a) => a.id === "hal")!.password, "old-password");
assert.equal(await grants.remove("grafana"), "not ours");
assert.equal(marked({ token: "grafana", description: `${MARK} x` }), false, "the mark needs the mesh's name too");
});
test("an org the instance does not have, or a non-mesh name, makes nothing", async () => {
const elsewhere = new ApiGrants(new InfluxDBClient(`http://127.0.0.1:${port}`, ADMIN, "nope"), "nope");
await assert.rejects(elsewhere.ensure(grafana()), /no org "nope"/);
await assert.rejects(grants.ensure({ ...grafana(), as: "grafana" }), /not a mesh identity/);
assert.equal(auths.length, 0);
});
test("a withdrawn consumer's authorization is removed, and an absent one is not an error", async () => {
await grants.ensure(grafana());
assert.equal(await grants.remove("mesh_ace_grafana"), "removed");
assert.equal(auths.length, 0);
assert.equal(await grants.remove("mesh_ace_grafana"), "absent");
});
+6 -1
View File
@@ -8,5 +8,10 @@
"skipLibCheck": true, "skipLibCheck": true,
"noEmit": true "noEmit": true
}, },
"include": ["client.ts", "tools/index.ts"] "include": [
"client.ts",
"grants.ts",
"provisioner/index.ts",
"tools/index.ts"
]
} }
+53 -16
View File
@@ -2,7 +2,8 @@
// an indexer proxy: it normalises many torrent trackers behind one Torznab surface. This client // an indexer proxy: it normalises many torrent trackers behind one Torznab surface. This client
// talks its /api/v2.0 REST API, and only jackett's tools import it. // talks its /api/v2.0 REST API, and only jackett's tools import it.
import { readFileSync } from "node:fs"; import { existsSync, readFileSync } from "node:fs";
import { join } from "node:path";
export interface JackettIndexer { export interface JackettIndexer {
id: string; id: string;
@@ -43,18 +44,36 @@ export class JackettClient {
/** /**
* Build from the module's resolved environment. Jackett's REST API is keyed, so both the URL and * Build from the module's resolved environment. Jackett's REST API is keyed, so both the URL and
* the key must be present — without them there is nothing to talk to, so this throws and the * the key must be present. The key is read from the settings-merged config or MESH_JACKETT_API_KEY,
* module contributes no tools rather than failing half-configured. * or, failing those, discovered from Jackett's own ServerConfig.json under MESH_JACKETT_CONFIG_DIR
* — the file Jackett writes it to, as sonarr/radarr read theirs from config.xml — so a running
* server needs no key configured by hand and no secret has to be put in an assignment. Without a
* URL or key there is nothing to talk to, so this throws and the module contributes no tools
* rather than failing half-configured.
*/ */
static fromEnv(env: NodeJS.ProcessEnv = process.env): JackettClient { static fromEnv(env: NodeJS.ProcessEnv = process.env): JackettClient {
const cfg = meshConfig(env.MESH_JACKETT_CONFIG_FILE); const cfg = meshConfig(env.MESH_JACKETT_CONFIG_FILE);
const url = cfg.url ?? env.MESH_JACKETT_URL; const url = cfg.url ?? env.MESH_JACKETT_URL;
const apiKey = cfg.apiKey ?? env.MESH_JACKETT_API_KEY; const apiKey = cfg.apiKey ?? env.MESH_JACKETT_API_KEY
?? JackettClient.detectApiKey(env.MESH_JACKETT_CONFIG_DIR ?? "/config");
if (!url) throw new Error("no Jackett URL — set MESH_JACKETT_URL"); if (!url) throw new Error("no Jackett URL — set MESH_JACKETT_URL");
if (!apiKey) throw new Error("no Jackett API key — set MESH_JACKETT_API_KEY"); if (!apiKey) throw new Error("no Jackett API key — set MESH_JACKETT_API_KEY or make the config dir readable");
return new JackettClient(url, apiKey); return new JackettClient(url, apiKey);
} }
/** Discover the API key from Jackett's ServerConfig.json (the linuxserver image keeps it at
* <config>/Jackett/ServerConfig.json), falling back to null. */
static detectApiKey(configDir: string): string | null {
for (const file of [join(configDir, "Jackett", "ServerConfig.json"), join(configDir, "ServerConfig.json")]) {
if (!existsSync(file)) continue;
try {
const key = (JSON.parse(readFileSync(file, "utf8")) as { APIKey?: unknown }).APIKey;
if (typeof key === "string" && key) return key;
} catch { /* unreadable or mid-write: try the next, then give up */ }
}
return null;
}
private async get(path: string, params: Record<string, string> = {}): Promise<any> { private async get(path: string, params: Record<string, string> = {}): Promise<any> {
const url = new URL(`${this.baseUrl}${path}`); const url = new URL(`${this.baseUrl}${path}`);
url.searchParams.set("apikey", this.apiKey); url.searchParams.set("apikey", this.apiKey);
@@ -64,18 +83,36 @@ export class JackettClient {
return res.json(); return res.json();
} }
/** The configured indexers Jackett proxies. `configured=false` also lists the ones not set up. */ /**
* The configured indexers Jackett proxies. `configured=false` also lists the ones not set up.
* Read from the Torznab `t=indexers` feed, not /api/v2.0/indexers: that one is the web UI's and
* wants a login cookie (it answers an API-key request with a redirect), while the Torznab feed is
* what the key is for. The feed carries no last error, so `lastError` stays unset.
*/
async getIndexers(configuredOnly = true): Promise<JackettIndexer[]> { async getIndexers(configuredOnly = true): Promise<JackettIndexer[]> {
const raw = await this.get("/api/v2.0/indexers", { configured: configuredOnly ? "true" : "false" }); const url = new URL(`${this.baseUrl}/api/v2.0/indexers/all/results/torznab/api`);
const list = Array.isArray(raw) ? raw : []; url.searchParams.set("apikey", this.apiKey);
return list.map((i: any) => ({ url.searchParams.set("t", "indexers");
id: i.id, url.searchParams.set("configured", configuredOnly ? "true" : "false");
name: i.name, const res = await fetch(url.toString(), { headers: { Accept: "application/xml" } });
type: i.type, if (!res.ok) throw new Error(`Jackett API torznab t=indexers: ${res.status} ${await res.text()}`);
configured: i.configured ?? false, const xml = await res.text();
siteLink: i.site_link, // Torznab reports failures (a wrong key among them) as 200 with an <error> body.
lastError: i.last_error || undefined, const err = xml.match(/<error code="(\d+)" description="([^"]*)"/);
})); if (err) throw new Error(`Jackett API torznab t=indexers: error ${err[1]} ${err[2]}`);
const text = (block: string, tag: string) =>
block.match(new RegExp(`<${tag}>([^<]*)</${tag}>`))?.[1];
const out: JackettIndexer[] = [];
for (const m of xml.matchAll(/<indexer id="([^"]+)" configured="([^"]+)">([\s\S]*?)<\/indexer>/g)) {
out.push({
id: m[1],
name: text(m[3], "title") ?? m[1],
type: text(m[3], "type") ?? "unknown",
configured: m[2] === "true",
siteLink: text(m[3], "link"),
});
}
return out;
} }
/** /**
+27 -9
View File
@@ -1,6 +1,19 @@
{ {
"module": "jackett", "module": "jackett",
"version": "1", "version": "1",
"provides": [
{
"name": "jackett-api",
"scope": "mesh"
}
],
"serves": {
"jackett-api": {
"scheme": "http",
"port": 9117,
"url-base": ""
}
},
"capabilities": [ "capabilities": [
"container-runtime" "container-runtime"
], ],
@@ -10,7 +23,7 @@
"port": 9117, "port": 9117,
"protocol": "tcp", "protocol": "tcp",
"from": "mesh", "from": "mesh",
"why": "the indexer proxy" "why": "the indexer proxy: its web UI, and the Torznab feeds the *arr apps search through, which other modules reach as jackett-api"
} }
], ],
"resources": [ "resources": [
@@ -20,10 +33,15 @@
"path": "/var/lib/mesh/jackett", "path": "/var/lib/mesh/jackett",
"mode": "0700" "mode": "0700"
}, },
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
},
{ {
"id": "config", "id": "config",
"type": "directory", "type": "directory",
"path": "/services/jackett/config",
"mode": "0700", "mode": "0700",
"owner": "1000:1000" "owner": "1000:1000"
}, },
@@ -31,7 +49,7 @@
"id": "server", "id": "server",
"type": "container", "type": "container",
"name": "jackett", "name": "jackett",
"image": "lscr.io/linuxserver/jackett@sha256:fd72d42b731ebf750b5de9711127251cf3b3f609419c32083ea8b3b3ee840b77", "image": "lscr.io/linuxserver/jackett@sha256:7b19f4f6ac33d855ca9226600ecbd096ee678f66da28b13a7c09980b035ff583",
"env": { "env": {
"PUID": "1000", "PUID": "1000",
"PGID": "1000", "PGID": "1000",
@@ -41,13 +59,13 @@
"9117" "9117"
], ],
"volumes": [ "volumes": [
"/services/jackett/config:/config" "${dir:config}:/config"
] ]
}, },
{ {
"id": "runtime-config", "id": "runtime-config",
"type": "file", "type": "file",
"path": "/var/lib/mesh/jackett/config.json", "path": "${dir:state}/config.json",
"mode": "0600", "mode": "0600",
"content": "{}\n", "content": "{}\n",
"merge": "json" "merge": "json"
@@ -59,12 +77,12 @@
"network": "host", "network": "host",
"volumes": [ "volumes": [
"/var/lib/mesh/jackett/broker:/run/secrets/broker:ro", "/var/lib/mesh/jackett/broker:/run/secrets/broker:ro",
"/var/lib/mesh/jackett/config.json:/run/config/config.json:ro", "${dir:state}/config.json:/run/config/config.json:ro",
"/services/jackett/config:/var/lib/jackett/config:ro" "${dir:config}:/var/lib/jackett/config:ro"
], ],
"env": { "env": {
"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_JACKETT_URL": "http://127.0.0.1:9117", "MESH_JACKETT_URL": "http://127.0.0.1:${port:9117}",
"MESH_JACKETT_CONFIG_FILE": "/run/config/config.json", "MESH_JACKETT_CONFIG_FILE": "/run/config/config.json",
"MESH_JACKETT_CONFIG_DIR": "/var/lib/jackett/config" "MESH_JACKETT_CONFIG_DIR": "/var/lib/jackett/config"
}, },
@@ -87,7 +105,7 @@
} }
}, },
"binds": { "binds": {
"route": "/var/lib/mesh/jackett/route.json" "route": "${dir:state}/route.json"
}, },
"build": { "build": {
"on": [ "on": [
+1 -1
View File
@@ -9,7 +9,7 @@ export function getJackettTools(jackett: JackettClient): ToolDefinition[] {
return [ return [
{ {
name: "jackett_indexers", name: "jackett_indexers",
description: "List the indexers Jackett proxies, with their type and any last error.", description: "List the indexers Jackett proxies, with their type and site.",
input: { all: { type: "boolean", description: "include indexers not yet configured (default false)" } }, input: { all: { type: "boolean", description: "include indexers not yet configured (default false)" } },
run: async (args) => { run: async (args) => {
const indexers = await jackett.getIndexers(!args.all); const indexers = await jackett.getIndexers(!args.all);
+1 -1
View File
@@ -13,7 +13,7 @@ ARG RUNTIME_BASE
FROM ${BUILD_BASE} AS build FROM ${BUILD_BASE} AS build
WORKDIR /app/modules/mosquitto WORKDIR /app/modules/mosquitto
COPY . . COPY . .
RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts tools/index.ts provisioner/index.ts bootstrap/index.ts \ RUN node /app/node_modules/typescript/bin/tsc topics.ts client.ts index.ts tools/index.ts provisioner/index.ts bootstrap/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE} FROM ${RUNTIME_BASE}
+38 -24
View File
@@ -20,6 +20,8 @@ import { readFileSync } from "node:fs";
import { execFile } from "node:child_process"; import { execFile } from "node:child_process";
import { promisify } from "node:util"; import { promisify } from "node:util";
import { missingAcls, parseRoleAcls, staleAcls, wantedAcls } from "./topics.js";
const run = promisify(execFile); const run = promisify(execFile);
export interface MqttConn { export interface MqttConn {
@@ -141,14 +143,19 @@ export class MosquittoClient {
} }
/** /**
* Create (or reset to a known state) a client scoped to one topic namespace, idempotently. The * Create (or reset to a known state) a client granted exactly these topic filters, idempotently.
* client is confined to `<prefix>/#` by a same-named role: it may publish to, subscribe to and * The grant is a same-named role carrying, for every filter, publish, receive and subscribe — and
* receive on exactly its own subtree and nothing else — the MQTT analog of redis's keyspace-scoped * nothing else: an ACL the role carries that the filters no longer name is removed, so narrowing a
* ACL user. Called again for an existing client, it resets the password and re-asserts the ACLs. * consumer's `topics` narrows what it may do. By default the filters are the consumer's own
* subtree, `<as>/#` (see topics.ts). Called again for an existing client, it resets the password
* and re-asserts the ACLs.
*
* Only the role named for this client is ever changed. A client or role the mesh did not make —
* a device carried from the predecessor's password file, its `legacy-full-access` role — is never
* read, changed or removed here.
*/ */
async createScopedClient(username: string, password: string, topicPrefix: string): Promise<void> { async createScopedClient(username: string, password: string, filters: readonly string[]): Promise<void> {
const role = username; // one role per client, named for it const role = username; // one role per client, named for it
const pattern = `${topicPrefix}/#`;
if (await this.clientExists(username)) { if (await this.clientExists(username)) {
await this.ctl("setClientPassword", username, password); await this.ctl("setClientPassword", username, password);
@@ -161,17 +168,18 @@ export class MosquittoClient {
await this.ctl("createClient", username, "-p", password); await this.ctl("createClient", username, "-p", password);
} }
// A role carrying exactly this client's topic ACLs. createRole, addRoleACL and addClientRole are // createRole and addRoleACL are one-shot: each rejects with an "already exists" when re-run
// all one-shot: each rejects with an "already exists" when re-run against a role/ACL/binding it // against a role/ACL it created on a previous reconcile. That rejection is the intended terminal
// created on a previous reconcile. That rejection is the intended terminal state — the ACL is // state, so it is swallowed.
// deterministic (`<prefix>/#`, allow), so re-adding the identical entry is a no-op — so it is
// swallowed. (Until the exit code was fixed this was invisible: the tool returned 0 and the
// rejection was lost; now it surfaces, and each of these adds must tolerate its own idempotent
// re-run explicitly.)
await ignoreExisting(this.ctl("createRole", role)); await ignoreExisting(this.ctl("createRole", role));
for (const acl of ["publishClientSend", "publishClientReceive", "subscribePattern"]) { const wanted = wantedAcls(filters);
// allow (1) this client to send to, receive on, and subscribe under its own subtree. const current = parseRoleAcls(await this.ctl("getRole", role));
await ignoreExisting(this.ctl("addRoleACL", role, acl, pattern, "allow")); for (const acl of missingAcls(current, wanted)) {
await ignoreExisting(this.ctl("addRoleACL", role, acl.type, acl.topic, "allow"));
}
// What the consumer no longer asks for — added before it narrowed its topics — is taken away.
for (const acl of staleAcls(current, wanted)) {
await ignoreMissing(this.ctl("removeRoleACL", role, acl.type, acl.topic));
} }
// Bind the role only when it is not already bound — addClientRole is the one call whose // Bind the role only when it is not already bound — addClientRole is the one call whose
// idempotent re-run cannot be recognised by message (see clientHasRole). // idempotent re-run cannot be recognised by message (see clientHasRole).
@@ -181,25 +189,31 @@ export class MosquittoClient {
} }
/** /**
* Whether a consumer's client accepts exactly this password and still carries its own role. * Whether a consumer's client accepts exactly this password, still carries its own role, and that
* Read-only. The password is checked the way the consumer is checked, by an MQTT CONNECT as it, * role grants exactly these filters. Read-only. The password is checked the way the consumer is
* and the broker's CONNACK code is the answer: 0 accepted, 4 bad credentials, 5 not authorised. * checked, by an MQTT CONNECT as it, and the broker's CONNACK code is the answer: 0 accepted,
* Nothing rides on argv. An unreachable broker rejects (novox/hq issue 120). * 4 bad credentials, 5 not authorised. Nothing rides on argv. An unreachable broker rejects
* (novox/hq issue 120).
*/ */
async holdsClient(username: string, password: string): Promise<boolean> { async holdsClient(username: string, password: string, filters: readonly string[]): Promise<boolean> {
const code = await mqttConnack(this.conn.host, this.conn.port, username, password); const code = await mqttConnack(this.conn.host, this.conn.port, username, password);
if (code === 4 || code === 5) return false; if (code === 4 || code === 5) return false;
if (code !== 0) throw new Error(`mosquitto refused ${username} with CONNACK ${code}`); if (code !== 0) throw new Error(`mosquitto refused ${username} with CONNACK ${code}`);
// The role, asked directly: only "not found" means absent. Any other failure to ask rejects, // The role, asked directly: only "not found" means absent. Any other failure to ask rejects,
// unlike clientHasRole, which reads every failure as "no role". // unlike clientHasRole, which reads every failure as "no role".
let out: string; let client: string;
let role: string;
try { try {
out = await this.ctl("getClient", username); client = await this.ctl("getClient", username);
role = await this.ctl("getRole", username);
} catch (err) { } catch (err) {
if (/not\s*found|does not exist|no such/i.test(String(err))) return false; if (/not\s*found|does not exist|no such/i.test(String(err))) return false;
throw err; throw err;
} }
return new RegExp(`(^|\\s)${escapeRegExp(username)}\\s+\\(priority`, "m").test(out); if (!new RegExp(`(^|\\s)${escapeRegExp(username)}\\s+\\(priority`, "m").test(client)) return false;
const current = parseRoleAcls(role);
const wanted = wantedAcls(filters);
return missingAcls(current, wanted).length === 0 && staleAcls(current, wanted).length === 0;
} }
/** Remove a client and the per-client role created for it, idempotently. */ /** Remove a client and the per-client role created for it, idempotently. */
+18 -17
View File
@@ -20,16 +20,19 @@
"mosquitto.topic.deprovisioned" "mosquitto.topic.deprovisioned"
], ],
"serves": { "serves": {
"mqtt-topic": {} "mqtt-topic": {
"scheme": "mqtt",
"port": 1883
}
}, },
"receives": { "receives": {
"mqtt-topic": "/var/lib/mosquitto-module/grants/mesh.json" "mqtt-topic": "${dir:grants}/mesh.json"
}, },
"grants": { "grants": {
"mqtt-topic": "/var/lib/mosquitto-module/grants" "mqtt-topic": "${dir:grants}"
}, },
"own-secrets": { "own-secrets": {
"admin": "/var/lib/mosquitto-module/admin.secret", "admin": "/var/lib/mesh/mosquitto/admin",
"broker": "/var/lib/mesh/mosquitto/broker" "broker": "/var/lib/mesh/mosquitto/broker"
}, },
"listens": [ "listens": [
@@ -58,26 +61,24 @@
{ {
"id": "state", "id": "state",
"type": "directory", "type": "directory",
"path": "/var/lib/mosquitto-module", "mode": "0700",
"mode": "0700" "place": "."
}, },
{ {
"id": "grants-dir", "id": "grants",
"type": "directory", "type": "directory",
"path": "/var/lib/mosquitto-module/grants",
"mode": "0700" "mode": "0700"
}, },
{ {
"id": "data", "id": "data",
"type": "directory", "type": "directory",
"path": "/services/mosquitto/data",
"mode": "0700", "mode": "0700",
"owner": "1883:1883" "owner": "1883:1883"
}, },
{ {
"id": "server-conf", "id": "server-conf",
"type": "file", "type": "file",
"path": "/var/lib/mosquitto-module/mosquitto.conf", "path": "${dir:state}/mosquitto.conf",
"mode": "0600", "mode": "0600",
"owner": "1883:1883", "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" "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"
@@ -93,8 +94,8 @@
"name": "mosquitto-bootstrap", "name": "mosquitto-bootstrap",
"run-once": true, "run-once": true,
"volumes": [ "volumes": [
"/services/mosquitto/data:/mosquitto/data", "${dir:data}:/mosquitto/data",
"/var/lib/mosquitto-module/admin.secret:/run/secrets/admin:ro" "/var/lib/mesh/mosquitto/admin:/run/secrets/admin:ro"
], ],
"env": { "env": {
"MESH_PROVISION_MQTT": "mosquitto:1883", "MESH_PROVISION_MQTT": "mosquitto:1883",
@@ -112,15 +113,15 @@
"id": "server", "id": "server",
"type": "container", "type": "container",
"name": "mosquitto", "name": "mosquitto",
"image": "eclipse-mosquitto@sha256:6f8d8a947c506f8a2290ec65cd4bd2bc7cb4d43fb5f6271f861cb013e2ef9797", "image": "eclipse-mosquitto@sha256:38c0da4f2ef84284d47b3b3eeea1cb3bdeabe81ee10caf0cd5c5ff61ee3ea408",
"network": "mosquitto", "network": "mosquitto",
"ports": [ "ports": [
"1883", "1883",
"8081" "8081"
], ],
"volumes": [ "volumes": [
"/services/mosquitto/data:/mosquitto/data", "${dir:data}:/mosquitto/data",
"/var/lib/mosquitto-module/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro" "${dir:state}/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro"
] ]
}, },
{ {
@@ -130,8 +131,8 @@
"network": "mosquitto", "network": "mosquitto",
"volumes": [ "volumes": [
"/var/lib/mesh/mosquitto/broker:/run/secrets/broker:ro", "/var/lib/mesh/mosquitto/broker:/run/secrets/broker:ro",
"/var/lib/mosquitto-module/grants:/var/lib/mosquitto-module/grants:ro", "${dir:grants}:/var/lib/mosquitto-module/grants:ro",
"/var/lib/mosquitto-module/admin.secret:/run/secrets/admin:ro" "/var/lib/mesh/mosquitto/admin:/run/secrets/admin:ro"
], ],
"env": { "env": {
"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_BROKER_FILE": "/run/secrets/broker",
+6 -1
View File
@@ -1,9 +1,14 @@
{ {
"name": "@novox/module-mosquitto", "name": "@novox/module-mosquitto",
"version": "0.1.0", "version": "0.1.0",
"description": "mosquitto — provides the mesh mqtt-topic interface. Its admin client, provisioner, tools and events live here (novox/hq ADR 0039).", "description": "mosquitto \u2014 provides the mesh mqtt-topic interface. Its admin client, provisioner, tools and events live here (novox/hq ADR 0039).",
"type": "module", "type": "module",
"private": true, "private": true,
"scripts": {
"build": "tsc topics.ts client.ts index.ts tools/index.ts provisioner/index.ts bootstrap/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist",
"typecheck": "tsc -p tsconfig.json",
"test": "node --test --experimental-strip-types 'test/*.test.ts'"
},
"dependencies": { "dependencies": {
"@novox/mesh-sdk": "^0.1.1" "@novox/mesh-sdk": "^0.1.1"
}, },
+24 -6
View File
@@ -5,7 +5,14 @@
// //
// The `mqtt-topic` interface: a consumer connects as `as` with the password the mesh minted, and // 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 // publishes and subscribes under `<as>/#`, isolated from every other consumer by a Dynamic Security
// role scoped to exactly that subtree. // role scoped to exactly that subtree — unless it contributed `topics`, the MQTT topic filters its
// work needs (a home-automation hub needs the devices' topics); then the role grants exactly those
// (topics.ts). A list that is not valid topic filters is refused, and the consumer is not created
// or changed until it is fixed.
//
// What a consumer is told (its binding): `at` — the broker's machine — and `port`, the machine port
// of the MQTT listener (the manifest's `serves`); `as` is its login, and its copy of the password is
// the pair credential the mesh delivers to it.
// //
// **The login and password are the mesh's, not the provisioner's (ADR 0048).** The mesh derives the // **The login and password are the mesh's, not the provisioner's (ADR 0048).** The mesh derives the
// login and hands it to both ends so they agree, and mints the password and delivers a copy to each. // login and hands it to both ends so they agree, and mints the password and delivers a copy to each.
@@ -15,6 +22,7 @@
import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner"; import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner";
import { emit } from "@novox/mesh-sdk/events"; import { emit } from "@novox/mesh-sdk/events";
import { MosquittoClient } from "../client.js"; import { MosquittoClient } from "../client.js";
import { topicFilters } from "../topics.js";
const mosquitto = MosquittoClient.fromEnv(); const mosquitto = MosquittoClient.fromEnv();
@@ -29,13 +37,20 @@ async function announce(type: string, body: Record<string, string>): Promise<voi
runProvisioner("mqtt-topic", { runProvisioner("mqtt-topic", {
async create(p: Provision): Promise<void> { async create(p: Provision): Promise<void> {
// The topic subtree is scoped to the consumer's own login, so one cannot read another's topics. // By default the consumer's own subtree, so one cannot read another's topics; what it
const topicPrefix = p.as; // contributed as `topics` otherwise.
await mosquitto.createScopedClient(p.as, p.password, topicPrefix); const granted = topicFilters(p.values, p.as);
if ("problem" in granted) {
// Thrown, so the harness logs it and retries: the consumer stays as it was (or absent) until
// its contribution is valid, rather than being given a grant it did not ask for.
throw new Error(`${p.as}: ${granted.problem}`);
}
await mosquitto.createScopedClient(p.as, p.password, granted.filters);
await announce("topic.provisioned", { await announce("topic.provisioned", {
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
username: p.as, username: p.as,
topicPrefix, topicPrefix: granted.own ? p.as : "",
topics: granted.filters.join(" "),
}); });
}, },
@@ -46,6 +61,9 @@ runProvisioner("mqtt-topic", {
// Asked every minute by the harness: whether the backend still holds this consumer exactly as // Asked every minute by the harness: whether the backend still holds this consumer exactly as
// the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120).
async holds(p: Provision): Promise<boolean> { async holds(p: Provision): Promise<boolean> {
return mosquitto.holdsClient(p.as, p.password); const granted = topicFilters(p.values, p.as);
// An invalid list was never applied; create refuses it again, loudly, on every pass.
if ("problem" in granted) return false;
return mosquitto.holdsClient(p.as, p.password, granted.filters);
}, },
}); });
+72
View File
@@ -0,0 +1,72 @@
// What a consumer of mqtt-topic is granted (topics.ts): its own subtree unless it contributed
// `topics`; a contributed list is granted exactly, refused whole when it is not topic filters; and
// the role is brought to exactly the wanted ACLs — missing ones added, stale ones removed — read from
// `mosquitto_ctrl dynsec getRole` as eclipse-mosquitto 2.1.2 prints it.
import { test } from "node:test";
import assert from "node:assert/strict";
import { filterProblem, missingAcls, parseRoleAcls, staleAcls, topicFilters, wantedAcls } from "../topics.ts";
test("a consumer that contributed nothing gets its own subtree", () => {
assert.deepEqual(topicFilters({}, "mesh_ace_hass"), { ok: true, filters: ["mesh_ace_hass/#"], own: true });
assert.deepEqual(topicFilters(undefined, "x"), { ok: true, filters: ["x/#"], own: true });
// Settings merge into every contribution: keys that are not `topics` change nothing.
assert.deepEqual(topicFilters({ endpoints: { web: {} } }, "x"), { ok: true, filters: ["x/#"], own: true });
});
test("a contributed list is granted exactly, duplicates once", () => {
assert.deepEqual(topicFilters({ topics: ["#"] }, "x"), { ok: true, filters: ["#"], own: false });
assert.deepEqual(topicFilters({ topics: ["stat/+/POWER", "tele/#", "tele/#", "/octoprint/x"] }, "x"), {
ok: true,
filters: ["stat/+/POWER", "tele/#", "/octoprint/x"],
own: false,
});
});
test("a list that is not topic filters is refused whole", () => {
for (const topics of [[], "#", [""], ["a/#/b"], ["a#"], ["a/b+"], [42], ["a\u0000b"], {}]) {
const out = topicFilters({ topics } as Record<string, unknown>, "x");
assert.equal(out.ok, false, JSON.stringify(topics));
}
assert.equal(filterProblem("+/+/#"), undefined);
assert.equal(filterProblem("#"), undefined);
});
const GET_ROLE = `Warning: You are running mosquitto_ctrl without encryption.
This means all of the configuration changes you are making are visible on the network, including passwords.
Rolename: u1
ACLs: publishClientSend : allow : # (priority: 0)
subscribePattern : allow : u1/# (priority: 0)
publishClientReceive : deny : secret topic/with space (priority: -1)
`;
test("getRole's ACL lines are read, the warning and headings are not", () => {
assert.deepEqual(parseRoleAcls(GET_ROLE), [
{ type: "publishClientSend", allow: true, topic: "#" },
{ type: "subscribePattern", allow: true, topic: "u1/#" },
{ type: "publishClientReceive", allow: false, topic: "secret topic/with space" },
]);
assert.deepEqual(parseRoleAcls("Rolename: empty\nACLs:\n"), []);
});
test("the role is brought to exactly the wanted ACLs", () => {
const current = parseRoleAcls(GET_ROLE);
const wanted = wantedAcls(["u1/#"]);
assert.deepEqual(wanted, [
{ type: "publishClientSend", allow: true, topic: "u1/#" },
{ type: "publishClientReceive", allow: true, topic: "u1/#" },
{ type: "subscribePattern", allow: true, topic: "u1/#" },
]);
assert.deepEqual(missingAcls(current, wanted), [
{ type: "publishClientSend", allow: true, topic: "u1/#" },
{ type: "publishClientReceive", allow: true, topic: "u1/#" },
]);
assert.deepEqual(staleAcls(current, wanted), [
{ type: "publishClientSend", allow: true, topic: "#" },
{ type: "publishClientReceive", allow: false, topic: "secret topic/with space" },
]);
assert.deepEqual(staleAcls(wanted, wanted), []);
assert.deepEqual(missingAcls(wanted, wanted), []);
});
+107
View File
@@ -0,0 +1,107 @@
// Which topics a consumer of `mqtt-topic` may use — the one choice a consumer makes about its grant.
//
// **By default, its own subtree and nothing else.** A consumer connects as the login the mesh derived
// (`as`) and may publish, receive and subscribe under `<as>/#` — isolated from every other consumer,
// which is the point of a per-consumer client (novox/hq ADR 0039/0048).
//
// **A consumer whose work IS the shared topic space says so.** Home Assistant discovers devices
// under `homeassistant/#` and `tasmota/discovery/#` and follows whatever state topics they announce;
// Node-RED's flows subscribe to the topics devices publish on (`stat/<device>/POWER`, …). Confined
// to `<as>/#` neither could do its job. So a consumer contributes `topics` to its `mqtt-topic`
// requirement — a list of MQTT topic filters — and the provisioner grants exactly those, both ways.
// Because assignment settings merge into every contribution, an operator narrows (or widens) the
// list per machine with the same key, without editing a manifest.
//
// Pure, so it is tested without a broker (test/topics.test.ts).
/** The dynsec ACL types a granted filter carries: send to it, receive from it, subscribe to it. */
export const GRANTED_ACL_TYPES = ["publishClientSend", "publishClientReceive", "subscribePattern"] as const;
/** One ACL on a role, as `mosquitto_ctrl dynsec getRole` reports it. */
export interface Acl {
type: string;
allow: boolean;
topic: string;
}
export type Filters = { ok: true; filters: string[]; own: boolean } | { ok: false; problem: string };
/**
* The topic filters a consumer is granted: what it contributed as `topics`, or its own subtree when
* it contributed nothing. Refused — never silently narrowed or widened — when the list is not a
* list of valid MQTT topic filters: a grant that quietly differs from what was asked is a consumer
* that fails somewhere far from the cause.
*/
export function topicFilters(values: Readonly<Record<string, unknown>> | undefined, as: string): Filters {
const given = values?.topics;
if (given === undefined || given === null) {
return { ok: true, filters: [`${as}/#`], own: true };
}
if (!Array.isArray(given) || given.length === 0) {
return { ok: false, problem: `topics must be a non-empty list of MQTT topic filters, not ${JSON.stringify(given)}` };
}
const out: string[] = [];
for (const f of given) {
if (typeof f !== "string") {
return { ok: false, problem: `topics holds ${JSON.stringify(f)}, which is not a topic filter` };
}
const problem = filterProblem(f);
if (problem) return { ok: false, problem: `topic filter ${JSON.stringify(f)}: ${problem}` };
if (!out.includes(f)) out.push(f);
}
return { ok: true, filters: out, own: out.length === 1 && out[0] === `${as}/#` };
}
/** Why a string is not a valid MQTT topic filter (MQTT 3.1.1 §4.7), or undefined when it is one. */
export function filterProblem(filter: string): string | undefined {
if (filter.length === 0) return "it is empty";
if (Buffer.byteLength(filter, "utf8") > 65535) return "it is longer than MQTT allows";
if (filter.includes("\u0000")) return "it contains a NUL character";
const levels = filter.split("/");
for (let i = 0; i < levels.length; i++) {
const level = levels[i];
if (level.includes("#") && (level !== "#" || i !== levels.length - 1)) {
return "'#' must be a whole level, and the last one";
}
if (level.includes("+") && level !== "+") return "'+' must be a whole level";
}
return undefined;
}
/** The ACLs a role must carry to grant these filters: every granted type, allowed, on every filter. */
export function wantedAcls(filters: readonly string[]): Acl[] {
const out: Acl[] = [];
for (const topic of filters) {
for (const type of GRANTED_ACL_TYPES) out.push({ type, allow: true, topic });
}
return out;
}
/**
* The ACLs `mosquitto_ctrl dynsec getRole` lists, one per line under its "ACLs:" heading:
* `ACLs: publishClientSend : allow : # (priority: 0)`
* ` subscribePattern : allow : u1/# (priority: 0)`
*/
export function parseRoleAcls(output: string): Acl[] {
const out: Acl[] = [];
const line = /^(?:ACLs:)?\s*([A-Za-z]+)\s*:\s*(allow|deny)\s*:\s*(.*?)\s+\(priority:\s*-?\d+\)\s*$/;
for (const raw of output.split(/\r?\n/)) {
const m = raw.match(line);
if (m) out.push({ type: m[1], allow: m[2] === "allow", topic: m[3] });
}
return out;
}
const key = (a: Acl): string => `${a.type}\u0000${a.allow ? "allow" : "deny"}\u0000${a.topic}`;
/** ACLs a role carries that it should not: in `current` and not in `wanted`. */
export function staleAcls(current: readonly Acl[], wanted: readonly Acl[]): Acl[] {
const want = new Set(wanted.map(key));
return current.filter((a) => !want.has(key(a)));
}
/** ACLs a role should carry and does not. */
export function missingAcls(current: readonly Acl[], wanted: readonly Acl[]): Acl[] {
const have = new Set(current.map(key));
return wanted.filter((a) => !have.has(key(a)));
}
+1 -1
View File
@@ -8,5 +8,5 @@
"skipLibCheck": true, "skipLibCheck": true,
"noEmit": true "noEmit": true
}, },
"include": ["client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts", "bootstrap/index.ts"] "include": ["topics.ts", "client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts", "bootstrap/index.ts"]
} }
+14 -16
View File
@@ -27,13 +27,13 @@
} }
}, },
"receives": { "receives": {
"redis-cache": "/var/lib/redis-module/grants/mesh.json" "redis-cache": "${dir:grants}/mesh.json"
}, },
"grants": { "grants": {
"redis-cache": "/var/lib/redis-module/grants" "redis-cache": "${dir:grants}"
}, },
"secrets": { "secrets": {
"secret": "/var/lib/redis-module/default.secret" "secret": "${dir:state}/default.secret"
}, },
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/redis/broker" "broker": "/var/lib/mesh/redis/broker"
@@ -57,29 +57,27 @@
{ {
"id": "state", "id": "state",
"type": "directory", "type": "directory",
"path": "/var/lib/redis-module", "mode": "0700",
"mode": "0700" "place": "."
}, },
{ {
"id": "grants-dir", "id": "grants",
"type": "directory", "type": "directory",
"path": "/var/lib/redis-module/grants",
"mode": "0700" "mode": "0700"
}, },
{ {
"id": "data", "id": "data",
"type": "directory", "type": "directory",
"path": "/services/redis/data",
"mode": "0700", "mode": "0700",
"owner": "999:999" "owner": "999:1000"
}, },
{ {
"id": "server-conf", "id": "server-conf",
"type": "file", "type": "file",
"path": "/var/lib/redis-module/redis.conf", "path": "${dir:state}/redis.conf",
"mode": "0600", "mode": "0600",
"content": "requirepass ${secret:secret}\nappendonly yes\ndir /data\n", "content": "requirepass ${secret:secret}\nappendonly yes\ndir /data\n",
"owner": "999:999" "owner": "999:1000"
}, },
{ {
"id": "net", "id": "net",
@@ -90,14 +88,14 @@
"id": "server", "id": "server",
"type": "container", "type": "container",
"name": "redis", "name": "redis",
"image": "redis@sha256:ff02b58f971e7d7d156a1267e283fcbbeee91773b6aa36c49dac28ecfe28eadf", "image": "redis@sha256:520775a41a63e77e06c73e35d2fd9cc15921a609516818796b4ecbb813078bc7",
"network": "redis", "network": "redis",
"ports": [ "ports": [
"6379" "6379"
], ],
"volumes": [ "volumes": [
"/services/redis/data:/data", "${dir:data}:/data",
"/var/lib/redis-module/redis.conf:/etc/redis/redis.conf:ro" "${dir:state}/redis.conf:/etc/redis/redis.conf:ro"
], ],
"args": [ "args": [
"/etc/redis/redis.conf" "/etc/redis/redis.conf"
@@ -113,8 +111,8 @@
"network": "redis", "network": "redis",
"volumes": [ "volumes": [
"/var/lib/mesh/redis/broker:/run/secrets/broker:ro", "/var/lib/mesh/redis/broker:/run/secrets/broker:ro",
"/var/lib/redis-module/grants:/var/lib/redis-module/grants:ro", "${dir:grants}:/var/lib/redis-module/grants:ro",
"/var/lib/redis-module/default.secret:/run/secrets/default:ro" "${dir:state}/default.secret:/run/secrets/default:ro"
], ],
"env": { "env": {
"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_BROKER_FILE": "/run/secrets/broker",
+41 -21
View File
@@ -10,35 +10,35 @@
"port": 8443, "port": 8443,
"protocol": "tcp", "protocol": "tcp",
"from": "mesh", "from": "mesh",
"why": "the controller web UI, over its own self-signed tls; reaching it from outside is a route grant later" "why": "the controller web UI and API, over its own self-signed tls; named through the proxy as an https route"
}, },
{ {
"name": "inform", "name": "inform",
"port": 8080, "port": 8080,
"protocol": "tcp", "protocol": "tcp",
"from": "mesh", "from": "mesh",
"why": "device inform \u2014 how APs and switches check in and are adopted" "why": "device inform, how APs and switches check in and are adopted; the controller tells devices this number, so the machine must publish it on the same one"
}, },
{ {
"name": "stun", "name": "stun",
"port": 3478, "port": 3478,
"protocol": "udp", "protocol": "udp",
"from": "mesh", "from": "mesh",
"why": "STUN, so managed devices can find the controller through NAT" "why": "STUN for managed devices; the controller tells devices this number, so the machine must publish it on the same one"
}, },
{ {
"name": "discovery", "name": "discovery",
"port": 10001, "port": 10001,
"protocol": "udp", "protocol": "udp",
"from": "mesh", "from": "mesh",
"why": "device discovery \u2014 the controller finds unadopted devices on the network" "why": "device discovery broadcasts from unadopted devices and the UniFi apps"
}, },
{ {
"name": "discovery-l2", "name": "discovery-l2",
"port": 1902, "port": 1900,
"protocol": "udp", "protocol": "udp",
"from": "mesh", "from": "mesh",
"why": "layer-2 (UBNT) discovery broadcasts; published on 1902, the container listens on 1900" "why": "make-controller-discoverable-on-L2 (SSDP); the software listens on 1900, which machines commonly have taken by another SSDP speaker"
}, },
{ {
"name": "portal-tls", "name": "portal-tls",
@@ -76,10 +76,15 @@
"path": "/var/lib/mesh/unifi", "path": "/var/lib/mesh/unifi",
"mode": "0700" "mode": "0700"
}, },
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
},
{ {
"id": "data", "id": "data",
"type": "directory", "type": "directory",
"path": "/services/unifi/data",
"mode": "0700", "mode": "0700",
"owner": "1000:1000" "owner": "1000:1000"
}, },
@@ -87,17 +92,17 @@
"id": "server", "id": "server",
"type": "container", "type": "container",
"name": "unifi-controller", "name": "unifi-controller",
"image": "lscr.io/linuxserver/unifi-controller@sha256:fcd5d8b13a77a588c79c1b49e5fc9ad08115aa3bb1a3576c589c64908a68845f", "image": "lscr.io/linuxserver/unifi-controller@sha256:0ae315a3a45635e443899e30e86bd507c2c48922cb27f4bc7241777885f4650e",
"ports": [ "ports": [
"8443:8443", "8443",
"8080:8080", "8080",
"3478:3478/udp", "3478/udp",
"10001:10001/udp", "10001/udp",
"1902:1900/udp", "1900/udp",
"8843:8843", "8843",
"8880:8880", "8880",
"6789:6789", "6789",
"5514:5514/udp" "5514/udp"
], ],
"env": { "env": {
"PUID": "1000", "PUID": "1000",
@@ -107,7 +112,7 @@
"MEM_STARTUP": "1024" "MEM_STARTUP": "1024"
}, },
"volumes": [ "volumes": [
"/services/unifi/data:/config" "${dir:data}:/config"
] ]
}, },
{ {
@@ -115,7 +120,7 @@
"type": "file", "type": "file",
"path": "/var/lib/mesh/unifi/config.json", "path": "/var/lib/mesh/unifi/config.json",
"mode": "0600", "mode": "0600",
"content": "{}\n", "content": "{\n \"site\": \"default\",\n \"password\": \"${secret:controller}\"\n}\n",
"merge": "json" "merge": "json"
}, },
{ {
@@ -129,7 +134,7 @@
], ],
"env": { "env": {
"MESH_BROKER_FILE": "/run/secrets/broker", "MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_UNIFI_URL": "https://127.0.0.1:8443", "MESH_UNIFI_URL": "https://127.0.0.1:${port:8443}",
"MESH_UNIFI_CONFIG_FILE": "/run/config/config.json" "MESH_UNIFI_CONFIG_FILE": "/run/config/config.json"
}, },
"restart-on": [ "restart-on": [
@@ -138,8 +143,23 @@
"artifact": "runtime" "artifact": "runtime"
} }
], ],
"requires": [
"route"
],
"contributes": {
"route": {
"label": "unifi",
"endpoint": "web",
"scheme": "https",
"insecure": true
}
},
"binds": {
"route": "${dir:state}/route.json"
},
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/unifi/broker" "broker": "/var/lib/mesh/unifi/broker",
"controller": "/var/lib/mesh/unifi/controller"
}, },
"build": { "build": {
"on": [ "on": [