Compare commits
2
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1080f45012 | ||
|
|
a32394ec22 |
@@ -94,7 +94,7 @@
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_BAZARR_URL": "http://127.0.0.1:${port:6767}",
|
||||
"MESH_BAZARR_URL": "http://127.0.0.1:6767",
|
||||
"MESH_BAZARR_API_KEY_FILE": "/run/secrets/api-key",
|
||||
"MESH_BAZARR_CONFIG_FILE": "/run/config/config.json",
|
||||
"MESH_BAZARR_CONFIG_DIR": "/var/lib/bazarr/config"
|
||||
|
||||
@@ -76,7 +76,7 @@
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_BOOKSHELF_URL": "http://127.0.0.1:${port:8787}",
|
||||
"MESH_BOOKSHELF_URL": "http://127.0.0.1:8787",
|
||||
"MESH_BOOKSHELF_CONFIG_DIR": "/var/lib/bookshelf/config"
|
||||
},
|
||||
"artifact": "runtime"
|
||||
|
||||
@@ -75,7 +75,7 @@
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_LIDARR_URL": "http://127.0.0.1:${port:8686}",
|
||||
"MESH_LIDARR_URL": "http://127.0.0.1:8686",
|
||||
"MESH_LIDARR_CONFIG_DIR": "/var/lib/lidarr/config"
|
||||
},
|
||||
"artifact": "runtime"
|
||||
|
||||
@@ -13,7 +13,7 @@ ARG RUNTIME_BASE
|
||||
FROM ${BUILD_BASE} AS build
|
||||
WORKDIR /app/modules/mosquitto
|
||||
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
|
||||
|
||||
FROM ${RUNTIME_BASE}
|
||||
|
||||
+38
-24
@@ -20,6 +20,8 @@ import { readFileSync } from "node:fs";
|
||||
import { execFile } from "node:child_process";
|
||||
import { promisify } from "node:util";
|
||||
|
||||
import { missingAcls, parseRoleAcls, staleAcls, wantedAcls } from "./topics.js";
|
||||
|
||||
const run = promisify(execFile);
|
||||
|
||||
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
|
||||
* 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.
|
||||
* Create (or reset to a known state) a client granted exactly these topic filters, idempotently.
|
||||
* The grant is a same-named role carrying, for every filter, publish, receive and subscribe — and
|
||||
* nothing else: an ACL the role carries that the filters no longer name is removed, so narrowing a
|
||||
* 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 pattern = `${topicPrefix}/#`;
|
||||
|
||||
if (await this.clientExists(username)) {
|
||||
await this.ctl("setClientPassword", username, password);
|
||||
@@ -161,17 +168,18 @@ export class MosquittoClient {
|
||||
await this.ctl("createClient", username, "-p", password);
|
||||
}
|
||||
|
||||
// A role carrying exactly this client's topic ACLs. createRole, addRoleACL and addClientRole are
|
||||
// all one-shot: each rejects with an "already exists" when re-run against a role/ACL/binding it
|
||||
// created on a previous reconcile. That rejection is the intended terminal state — the ACL is
|
||||
// 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.)
|
||||
// createRole and addRoleACL are one-shot: each rejects with an "already exists" when re-run
|
||||
// against a role/ACL it created on a previous reconcile. That rejection is the intended terminal
|
||||
// state, so it is swallowed.
|
||||
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 ignoreExisting(this.ctl("addRoleACL", role, acl, pattern, "allow"));
|
||||
const wanted = wantedAcls(filters);
|
||||
const current = parseRoleAcls(await this.ctl("getRole", role));
|
||||
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
|
||||
// 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.
|
||||
* Read-only. The password is checked the way the consumer is checked, by an MQTT CONNECT as it,
|
||||
* and the broker's CONNACK code is the answer: 0 accepted, 4 bad credentials, 5 not authorised.
|
||||
* Nothing rides on argv. An unreachable broker rejects (novox/hq issue 120).
|
||||
* Whether a consumer's client accepts exactly this password, still carries its own role, and that
|
||||
* role grants exactly these filters. Read-only. The password is checked the way the consumer is
|
||||
* checked, by an MQTT CONNECT as it, and the broker's CONNACK code is the answer: 0 accepted,
|
||||
* 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);
|
||||
if (code === 4 || code === 5) return false;
|
||||
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,
|
||||
// unlike clientHasRole, which reads every failure as "no role".
|
||||
let out: string;
|
||||
let client: string;
|
||||
let role: string;
|
||||
try {
|
||||
out = await this.ctl("getClient", username);
|
||||
client = await this.ctl("getClient", username);
|
||||
role = await this.ctl("getRole", username);
|
||||
} catch (err) {
|
||||
if (/not\s*found|does not exist|no such/i.test(String(err))) return false;
|
||||
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. */
|
||||
|
||||
@@ -20,16 +20,19 @@
|
||||
"mosquitto.topic.deprovisioned"
|
||||
],
|
||||
"serves": {
|
||||
"mqtt-topic": {}
|
||||
"mqtt-topic": {
|
||||
"scheme": "mqtt",
|
||||
"port": 1883
|
||||
}
|
||||
},
|
||||
"receives": {
|
||||
"mqtt-topic": "/var/lib/mosquitto-module/grants/mesh.json"
|
||||
"mqtt-topic": "${dir:grants}/mesh.json"
|
||||
},
|
||||
"grants": {
|
||||
"mqtt-topic": "/var/lib/mosquitto-module/grants"
|
||||
"mqtt-topic": "${dir:grants}"
|
||||
},
|
||||
"own-secrets": {
|
||||
"admin": "/var/lib/mosquitto-module/admin.secret",
|
||||
"admin": "/var/lib/mesh/mosquitto/admin",
|
||||
"broker": "/var/lib/mesh/mosquitto/broker"
|
||||
},
|
||||
"listens": [
|
||||
@@ -58,26 +61,24 @@
|
||||
{
|
||||
"id": "state",
|
||||
"type": "directory",
|
||||
"path": "/var/lib/mosquitto-module",
|
||||
"mode": "0700"
|
||||
"mode": "0700",
|
||||
"place": "."
|
||||
},
|
||||
{
|
||||
"id": "grants-dir",
|
||||
"id": "grants",
|
||||
"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",
|
||||
"path": "${dir:state}/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"
|
||||
@@ -93,8 +94,8 @@
|
||||
"name": "mosquitto-bootstrap",
|
||||
"run-once": true,
|
||||
"volumes": [
|
||||
"/services/mosquitto/data:/mosquitto/data",
|
||||
"/var/lib/mosquitto-module/admin.secret:/run/secrets/admin:ro"
|
||||
"${dir:data}:/mosquitto/data",
|
||||
"/var/lib/mesh/mosquitto/admin:/run/secrets/admin:ro"
|
||||
],
|
||||
"env": {
|
||||
"MESH_PROVISION_MQTT": "mosquitto:1883",
|
||||
@@ -112,15 +113,15 @@
|
||||
"id": "server",
|
||||
"type": "container",
|
||||
"name": "mosquitto",
|
||||
"image": "eclipse-mosquitto@sha256:6f8d8a947c506f8a2290ec65cd4bd2bc7cb4d43fb5f6271f861cb013e2ef9797",
|
||||
"image": "eclipse-mosquitto@sha256:38c0da4f2ef84284d47b3b3eeea1cb3bdeabe81ee10caf0cd5c5ff61ee3ea408",
|
||||
"network": "mosquitto",
|
||||
"ports": [
|
||||
"1883",
|
||||
"8081"
|
||||
],
|
||||
"volumes": [
|
||||
"/services/mosquitto/data:/mosquitto/data",
|
||||
"/var/lib/mosquitto-module/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro"
|
||||
"${dir:data}:/mosquitto/data",
|
||||
"${dir:state}/mosquitto.conf:/mosquitto/config/mosquitto.conf:ro"
|
||||
]
|
||||
},
|
||||
{
|
||||
@@ -130,8 +131,8 @@
|
||||
"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"
|
||||
"${dir:grants}:/var/lib/mosquitto-module/grants:ro",
|
||||
"/var/lib/mesh/mosquitto/admin:/run/secrets/admin:ro"
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
|
||||
@@ -1,9 +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 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",
|
||||
"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": {
|
||||
"@novox/mesh-sdk": "^0.1.1"
|
||||
},
|
||||
|
||||
@@ -5,7 +5,14 @@
|
||||
//
|
||||
// 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.
|
||||
// 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
|
||||
// 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 { emit } from "@novox/mesh-sdk/events";
|
||||
import { MosquittoClient } from "../client.js";
|
||||
import { topicFilters } from "../topics.js";
|
||||
|
||||
const mosquitto = MosquittoClient.fromEnv();
|
||||
|
||||
@@ -29,13 +37,20 @@ async function announce(type: string, body: Record<string, string>): Promise<voi
|
||||
|
||||
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);
|
||||
// By default the consumer's own subtree, so one cannot read another's topics; what it
|
||||
// contributed as `topics` otherwise.
|
||||
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", {
|
||||
consumer: p.consumer ?? "",
|
||||
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
|
||||
// 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> {
|
||||
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);
|
||||
},
|
||||
});
|
||||
|
||||
@@ -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), []);
|
||||
});
|
||||
@@ -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)));
|
||||
}
|
||||
@@ -8,5 +8,5 @@
|
||||
"skipLibCheck": 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"]
|
||||
}
|
||||
|
||||
@@ -81,7 +81,7 @@
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_NZBGET_URL": "http://127.0.0.1:${port:6789}",
|
||||
"MESH_NZBGET_URL": "http://127.0.0.1:6789",
|
||||
"MESH_NZBGET_PASSWORD_FILE": "/run/secrets/password",
|
||||
"MESH_NZBGET_CONFIG_FILE": "/run/config/config.json",
|
||||
"MESH_NZBGET_CONFIG_DIR": "/var/lib/nzbget/config"
|
||||
|
||||
@@ -3,10 +3,8 @@
|
||||
// qbittorrent. Both this module's tools and its events entrypoint import it, and nothing outside
|
||||
// qbittorrent does.
|
||||
//
|
||||
// The WebUI authenticates with a session cookie obtained by POSTing credentials, and guards
|
||||
// against CSRF by checking the Referer header. The cookie was `SID` before qBittorrent 5.2 and is
|
||||
// `QBT_SID_<port>` since, and a successful login answers 200 "Ok." before and 204 with no body
|
||||
// since — both are accepted. Node's fetch keeps no cookie jar, so the cookie is
|
||||
// The WebUI authenticates with a session cookie (SID) obtained by POSTing credentials, and guards
|
||||
// against CSRF by checking the Referer header. Node's fetch keeps no cookie jar, so the SID is
|
||||
// captured on login and carried by hand on every later call, with a single re-login on expiry.
|
||||
|
||||
import { readFileSync } from "node:fs";
|
||||
@@ -41,21 +39,6 @@ function meshConfig(file?: string): Record<string, string> {
|
||||
catch { return {}; }
|
||||
}
|
||||
|
||||
/** The WebUI username from qBittorrent's own qBittorrent.conf, in the config directory the mesh
|
||||
* mounts read-only (MESH_QBITTORRENT_CONFIG_DIR, which is the container's /config). The software's
|
||||
* file is the truth about who may log in, so the tools ask it rather than a setting that could
|
||||
* disagree. Absent, unreadable or unset yields undefined. */
|
||||
function confUsername(dir: string | undefined): string | undefined {
|
||||
if (!dir) return undefined;
|
||||
try {
|
||||
const line = readFileSync(`${dir.replace(/\/$/, "")}/qBittorrent/qBittorrent.conf`, "utf8")
|
||||
.split(/\r?\n/)
|
||||
.find((l) => l.startsWith("WebUI\\Username="));
|
||||
const value = line?.slice("WebUI\\Username=".length).trim();
|
||||
return value ? value : undefined;
|
||||
} catch { return undefined; }
|
||||
}
|
||||
|
||||
/** Read a secret the mesh mounted at a file path (an own-secret delivered by `secret accept`);
|
||||
* absent or unreadable yields undefined so callers fall back rather than crash. */
|
||||
function readSecret(file?: string): string | undefined {
|
||||
@@ -66,8 +49,7 @@ function readSecret(file?: string): string | undefined {
|
||||
|
||||
export class QbittorrentClient {
|
||||
readonly baseUrl: string;
|
||||
/** The session cookie as `name=value`, sent back exactly as it was set. */
|
||||
private session: string | null = null;
|
||||
private sid: string | null = null;
|
||||
|
||||
constructor(
|
||||
baseUrl: string,
|
||||
@@ -81,9 +63,7 @@ export class QbittorrentClient {
|
||||
* Build from the module's resolved environment. URL and password are read from
|
||||
* MESH_QBITTORRENT_URL and MESH_QBITTORRENT_PASSWORD; both must be present — an unconfigured
|
||||
* qBittorrent throws rather than pretend to be reachable, so the tools/events simply do not load
|
||||
* (the harness treats the throw as "exposes nothing"). The user is read from qBittorrent.conf,
|
||||
* falling back to "admin", the image's default. The password cannot be read there — qBittorrent
|
||||
* keeps only a PBKDF2 hash — so it is the own-secret the operator accepts.
|
||||
* (the harness treats the throw as "exposes nothing"). The user defaults to "admin".
|
||||
*/
|
||||
static fromEnv(env: NodeJS.ProcessEnv = process.env): QbittorrentClient {
|
||||
const cfg = meshConfig(env.MESH_QBITTORRENT_CONFIG_FILE);
|
||||
@@ -92,9 +72,7 @@ export class QbittorrentClient {
|
||||
if (!url || !password) {
|
||||
throw new Error("qBittorrent not configured — set MESH_QBITTORRENT_URL and MESH_QBITTORRENT_PASSWORD");
|
||||
}
|
||||
// The WebUI username: a setting or the environment if one says so, else whatever
|
||||
// qBittorrent.conf holds (an adopted machine keeps its own), else the image's "admin".
|
||||
const user = cfg.user ?? env.MESH_QBITTORRENT_USER ?? confUsername(env.MESH_QBITTORRENT_CONFIG_DIR) ?? "admin";
|
||||
const user = cfg.user ?? env.MESH_QBITTORRENT_USER ?? "admin";
|
||||
return new QbittorrentClient(url, user, password);
|
||||
}
|
||||
|
||||
@@ -104,22 +82,19 @@ export class QbittorrentClient {
|
||||
headers: { "Content-Type": "application/x-www-form-urlencoded", Referer: this.baseUrl },
|
||||
body: new URLSearchParams({ username: this.user, password: this.password }),
|
||||
});
|
||||
if (res.status === 401) throw new Error("qBittorrent login rejected — check credentials");
|
||||
if (!res.ok) throw new Error(`qBittorrent login: ${res.status} ${await res.text()}`);
|
||||
// 4.x/5.0/5.1 answer 200 "Ok." or 200 "Fails."; 5.2 answers 204 with no body, or 401.
|
||||
const body = (await res.text()).trim();
|
||||
if (res.status !== 204 && body !== "Ok.") {
|
||||
if ((await res.text()).trim() !== "Ok.") {
|
||||
throw new Error("qBittorrent login rejected — check credentials");
|
||||
}
|
||||
const match = res.headers.get("set-cookie")?.match(/((?:QBT_)?SID(?:_\d+)?)=([^;]+)/);
|
||||
if (!match) throw new Error("qBittorrent login returned no session cookie");
|
||||
this.session = `${match[1]}=${match[2]}`;
|
||||
const match = res.headers.get("set-cookie")?.match(/SID=([^;]+)/);
|
||||
if (!match) throw new Error("qBittorrent login returned no SID cookie");
|
||||
this.sid = match[1];
|
||||
}
|
||||
|
||||
private async call(method: "GET" | "POST", path: string, form?: Record<string, string>): Promise<Response> {
|
||||
if (!this.session) await this.login();
|
||||
if (!this.sid) await this.login();
|
||||
const doFetch = (): Promise<Response> => {
|
||||
const headers: Record<string, string> = { Referer: this.baseUrl, Cookie: this.session ?? "" };
|
||||
const headers: Record<string, string> = { Referer: this.baseUrl, Cookie: `SID=${this.sid}` };
|
||||
const init: RequestInit = { method, headers };
|
||||
if (form) {
|
||||
headers["Content-Type"] = "application/x-www-form-urlencoded";
|
||||
@@ -128,8 +103,8 @@ export class QbittorrentClient {
|
||||
return fetch(`${this.baseUrl}/api/v2/${path}`, init);
|
||||
};
|
||||
let res = await doFetch();
|
||||
if (res.status === 403 || res.status === 401) {
|
||||
// The session expired — re-authenticate once and retry, rather than fail a routine call.
|
||||
if (res.status === 403) {
|
||||
// The SID expired — re-authenticate once and retry, rather than fail a routine call.
|
||||
await this.login();
|
||||
res = await doFetch();
|
||||
}
|
||||
|
||||
@@ -17,24 +17,10 @@
|
||||
"listens": [
|
||||
{
|
||||
"name": "web",
|
||||
"port": 8112,
|
||||
"port": 8080,
|
||||
"protocol": "tcp",
|
||||
"from": "mesh",
|
||||
"why": "the download client's pages, and the WebUI API its consumers and its own tools call. qBittorrent refuses a request whose Host names a port other than the one it listens on, so it listens on the machine port itself (host network, WEBUI_PORT=${port:8112}) and the two cannot differ on any machine"
|
||||
},
|
||||
{
|
||||
"name": "peers",
|
||||
"port": 6881,
|
||||
"protocol": "tcp",
|
||||
"from": "mesh",
|
||||
"why": "incoming BitTorrent peer connections; announced to trackers and peers, so qBittorrent listens on the machine port itself (TORRENTING_PORT=${port:6881})"
|
||||
},
|
||||
{
|
||||
"name": "peers-udp",
|
||||
"port": 6881,
|
||||
"protocol": "udp",
|
||||
"from": "mesh",
|
||||
"why": "DHT and uTP on the same number as the peer port"
|
||||
"why": "the download client's pages"
|
||||
}
|
||||
],
|
||||
"accesses": [
|
||||
@@ -50,15 +36,10 @@
|
||||
"path": "/var/lib/mesh/qbittorrent",
|
||||
"mode": "0700"
|
||||
},
|
||||
{
|
||||
"id": "state",
|
||||
"type": "directory",
|
||||
"mode": "0700",
|
||||
"place": "."
|
||||
},
|
||||
{
|
||||
"id": "config",
|
||||
"type": "directory",
|
||||
"path": "/services/qbittorrent/config",
|
||||
"mode": "0700",
|
||||
"owner": "1000:1000"
|
||||
},
|
||||
@@ -66,24 +47,24 @@
|
||||
"id": "server",
|
||||
"type": "container",
|
||||
"name": "qbittorrent",
|
||||
"image": "lscr.io/linuxserver/qbittorrent@sha256:457e4eec2ee3f5e4ef59f237ad51f6143deba9f7445ab48bb5204a98888ef9aa",
|
||||
"image": "lscr.io/linuxserver/qbittorrent@sha256:a00b6a597a3832a1814cde0ef60abc55c94644f3f80902c3432f6af6de8d4a96",
|
||||
"env": {
|
||||
"PUID": "1000",
|
||||
"PGID": "1000",
|
||||
"TZ": "Etc/UTC",
|
||||
"WEBUI_PORT": "${port:8112}",
|
||||
"TORRENTING_PORT": "${port:6881}"
|
||||
"TZ": "Etc/UTC"
|
||||
},
|
||||
"volumes": [
|
||||
"${dir:config}:/config",
|
||||
"/services/media/downloads:/downloads"
|
||||
"ports": [
|
||||
"8080"
|
||||
],
|
||||
"network": "host"
|
||||
"volumes": [
|
||||
"/services/qbittorrent/config:/config",
|
||||
"/services/media/downloads:/downloads"
|
||||
]
|
||||
},
|
||||
{
|
||||
"id": "runtime-config",
|
||||
"type": "file",
|
||||
"path": "${dir:state}/config.json",
|
||||
"path": "/var/lib/mesh/qbittorrent/config.json",
|
||||
"mode": "0600",
|
||||
"content": "{}\n",
|
||||
"merge": "json"
|
||||
@@ -96,46 +77,22 @@
|
||||
"volumes": [
|
||||
"/var/lib/mesh/qbittorrent/broker:/run/secrets/broker:ro",
|
||||
"/var/lib/mesh/qbittorrent/password:/run/secrets/password:ro",
|
||||
"${dir:state}/config.json:/run/config/config.json:ro",
|
||||
"${dir:config}:/var/lib/qbittorrent/config:ro"
|
||||
"/var/lib/mesh/qbittorrent/config.json:/run/config/config.json:ro",
|
||||
"/services/qbittorrent/config:/var/lib/qbittorrent/config:ro"
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_QBITTORRENT_URL": "http://127.0.0.1:${port:8112}",
|
||||
"MESH_QBITTORRENT_URL": "http://127.0.0.1:8080",
|
||||
"MESH_QBITTORRENT_PASSWORD_FILE": "/run/secrets/password",
|
||||
"MESH_QBITTORRENT_CONFIG_FILE": "/run/config/config.json",
|
||||
"MESH_QBITTORRENT_CONFIG_DIR": "/var/lib/qbittorrent/config"
|
||||
},
|
||||
"restart-on": [
|
||||
"runtime-config",
|
||||
"needs-password"
|
||||
"runtime-config"
|
||||
],
|
||||
"artifact": "runtime"
|
||||
}
|
||||
],
|
||||
"provides": [
|
||||
"qbittorrent-api"
|
||||
],
|
||||
"serves": {
|
||||
"qbittorrent-api": {
|
||||
"scheme": "http",
|
||||
"port": 8112,
|
||||
"url-base": "",
|
||||
"username": "admin"
|
||||
}
|
||||
},
|
||||
"requires": [
|
||||
"route"
|
||||
],
|
||||
"contributes": {
|
||||
"route": {
|
||||
"label": "qbittorrent",
|
||||
"endpoint": "web"
|
||||
}
|
||||
},
|
||||
"binds": {
|
||||
"route": "${dir:state}/route.json"
|
||||
},
|
||||
"build": {
|
||||
"on": [
|
||||
{
|
||||
|
||||
@@ -75,7 +75,7 @@
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_RADARR_URL": "http://127.0.0.1:${port:7878}",
|
||||
"MESH_RADARR_URL": "http://127.0.0.1:7878",
|
||||
"MESH_RADARR_CONFIG_DIR": "/var/lib/radarr/config"
|
||||
},
|
||||
"artifact": "runtime"
|
||||
|
||||
@@ -100,7 +100,7 @@
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_SEARXNG_URL": "http://127.0.0.1:${port:8080}",
|
||||
"MESH_SEARXNG_URL": "http://127.0.0.1:8080",
|
||||
"MESH_SEARXNG_CONFIG_FILE": "/run/config/config.json"
|
||||
},
|
||||
"restart-on": [
|
||||
|
||||
@@ -80,7 +80,7 @@
|
||||
],
|
||||
"env": {
|
||||
"MESH_BROKER_FILE": "/run/secrets/broker",
|
||||
"MESH_SONARR_URL": "http://127.0.0.1:${port:8989}",
|
||||
"MESH_SONARR_URL": "http://127.0.0.1:8989",
|
||||
"MESH_SONARR_CONFIG_DIR": "/var/lib/sonarr/config"
|
||||
},
|
||||
"artifact": "runtime"
|
||||
|
||||
Reference in New Issue
Block a user