Compare commits
10
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
1080f45012 | ||
|
|
a32394ec22 | ||
|
|
8064e5da8f | ||
|
|
63a255c5cb | ||
|
|
7ad1fbd5c6 | ||
|
|
67f5f4cffd | ||
|
|
8797335fbc | ||
|
|
53dc108603 | ||
|
|
ebf5ba2d4c | ||
|
|
bbac08a7d2 |
@@ -0,0 +1,56 @@
|
|||||||
|
{
|
||||||
|
"module": "ca-trust",
|
||||||
|
"version": "1",
|
||||||
|
"slug": "catrust",
|
||||||
|
"capabilities": [
|
||||||
|
"service-manager"
|
||||||
|
],
|
||||||
|
"requires": [
|
||||||
|
"internal-acme-ca"
|
||||||
|
],
|
||||||
|
"seats": [
|
||||||
|
{
|
||||||
|
"name": "the-mesh-trust-anchor",
|
||||||
|
"scope": "node"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"claims": [
|
||||||
|
{
|
||||||
|
"name": "the-mesh-trust-anchor",
|
||||||
|
"scope": "node"
|
||||||
|
}
|
||||||
|
],
|
||||||
|
"resources": [
|
||||||
|
{
|
||||||
|
"id": "state",
|
||||||
|
"type": "directory",
|
||||||
|
"mode": "0700",
|
||||||
|
"place": "."
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "anchor",
|
||||||
|
"type": "file",
|
||||||
|
"path": "${dir:state}/anchor",
|
||||||
|
"mode": "0755",
|
||||||
|
"content": "#!/bin/sh\n# The mesh's internal certificate authority, trusted by this machine.\n#\n# Written by the mesh from the ca-trust module's manifest (novox/hq ADR 0147).\n# Editing it here lasts until the next apply.\n#\n# There is no prior trust to verify the fetch against \u2014 this is the thing that\n# establishes it \u2014 so it is made over the mesh's own private network, which is\n# what authenticates it (novox/hq ADR 0098, the same reasoning that lets the\n# route proxy fetch this root for itself). What comes back is checked here: a\n# body that is not a certificate is refused now, rather than believed and then\n# failed by whatever reads the trust store next.\nset -eu\n\nROOTS='https://${bound:internal-acme-ca:at}:${bound:internal-acme-ca:port}${bound:internal-acme-ca:roots}'\nANCHORS=/etc/ca-certificates/trust-source/anchors\nANCHOR=\"$ANCHORS/mesh-internal-ca.crt\"\n\n# Arch's layout, said out loud rather than assumed: a machine that keeps its\n# anchors elsewhere fails here, visibly, instead of writing a file nothing\n# reads. That failure is the signal that this belongs in the host, where one\n# operating system's difference lives (novox/hq ADR 0147, option 2).\n[ -d \"$ANCHORS\" ] || {\n\techo \"this machine keeps no trust anchors in $ANCHORS; ca-trust is written for that layout\" >&2\n\texit 1\n}\n\ncase \"${1:-}\" in\ninstall)\n\ttmp=$(mktemp)\n\ttrap 'rm -f \"$tmp\"' EXIT\n\t# The authority may still be starting, or this machine may have come up\n\t# before it: two minutes of asking, then an honest failure.\n\tn=0\n\twhile [ \"$n\" -lt 60 ]; do\n\t\tif curl --fail --silent --show-error --insecure --max-time 10 \\\n\t\t\t--output \"$tmp\" \"$ROOTS\" &&\n\t\t\tgrep -q 'BEGIN CERTIFICATE' \"$tmp\"; then\n\t\t\tinstall -m 0644 \"$tmp\" \"$ANCHOR\"\n\t\t\tupdate-ca-trust\n\t\t\texit 0\n\t\tfi\n\t\tn=$((n + 1))\n\t\tsleep 2\n\tdone\n\techo \"the authority at $ROOTS did not serve a certificate within two minutes\" >&2\n\texit 1\n\t;;\nremove)\n\t# What stopping the unit does, and therefore what being unassigned does.\n\trm -f \"$ANCHOR\"\n\tupdate-ca-trust\n\t;;\n*)\n\techo \"usage: $(basename \"$0\") install|remove\" >&2\n\texit 2\n\t;;\nesac\n"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "unit",
|
||||||
|
"type": "file",
|
||||||
|
"path": "/etc/systemd/system/mesh-ca-trust.service",
|
||||||
|
"mode": "0644",
|
||||||
|
"content": "[Unit]\nDescription=The mesh's internal certificate authority, trusted by this machine\n# novox/hq ADR 0147. Starting this unit places the mesh's root among this\n# machine's trust anchors; stopping it takes the root away again, which is what\n# the host does when the module is no longer assigned here.\nWants=network-online.target\nAfter=network-online.target\n\n[Service]\nType=oneshot\nRemainAfterExit=yes\nExecStart=${dir:state}/anchor install\nExecStop=${dir:state}/anchor remove\n\n[Install]\nWantedBy=multi-user.target\n"
|
||||||
|
},
|
||||||
|
{
|
||||||
|
"id": "trust",
|
||||||
|
"type": "service",
|
||||||
|
"unit": "mesh-ca-trust.service",
|
||||||
|
"state": "running",
|
||||||
|
"boot": "enabled",
|
||||||
|
"restart-on": [
|
||||||
|
"anchor",
|
||||||
|
"unit"
|
||||||
|
]
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
@@ -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
@@ -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. */
|
||||||
|
|||||||
@@ -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",
|
||||||
|
|||||||
@@ -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"
|
||||||
},
|
},
|
||||||
|
|||||||
@@ -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);
|
||||||
},
|
},
|
||||||
});
|
});
|
||||||
|
|||||||
@@ -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,
|
"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"]
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,84 +0,0 @@
|
|||||||
// Dial everything the mesh claims is reachable, and say what was found (novox/hq ADR 0145).
|
|
||||||
//
|
|
||||||
// Runs on a cadence, from this machine, in this module's own container — the same position every other
|
|
||||||
// module on the machine calls from. That is the whole point: a check run by the host or by the control
|
|
||||||
// plane reaches these addresses by a path no ordinary caller uses, and would have passed throughout the
|
|
||||||
// outage that produced this module (novox/hq 04-ISSUES/145).
|
|
||||||
//
|
|
||||||
// It reports and does nothing else. A checker that repaired things would be a second control plane.
|
|
||||||
|
|
||||||
import { readFileSync, writeFileSync, mkdirSync, renameSync } from "node:fs";
|
|
||||||
import { dirname, join } from "node:path";
|
|
||||||
|
|
||||||
import { dial, tally, targetsFor, type Counts, type Result, type Roster } from "../reach.js";
|
|
||||||
|
|
||||||
/** Where the mesh renders this machine's view of the others, and where the counts are kept between runs. */
|
|
||||||
const rosterFile = process.env.MESH_NETWORK_CHECKER_ROSTER ?? "/run/config/roster.json";
|
|
||||||
const stateDir = process.env.MESH_NETWORK_CHECKER_STATE ?? "/run/state";
|
|
||||||
const probePort = Number(process.env.MESH_NETWORK_CHECKER_PORT ?? "9876");
|
|
||||||
const publicPort = process.env.MESH_NETWORK_CHECKER_PUBLIC_PORT
|
|
||||||
? Number(process.env.MESH_NETWORK_CHECKER_PUBLIC_PORT)
|
|
||||||
: undefined;
|
|
||||||
const timeoutMs = Number(process.env.MESH_NETWORK_CHECKER_TIMEOUT_MS ?? "4000");
|
|
||||||
const threshold = Number(process.env.MESH_NETWORK_CHECKER_THRESHOLD ?? "2");
|
|
||||||
|
|
||||||
/** read is a JSON file or a stated failure — never a silent default, which is how a checker comes to
|
|
||||||
* report that everything is fine because it read nothing. */
|
|
||||||
function read<T>(path: string, whenMissing: T | null): T {
|
|
||||||
try {
|
|
||||||
return JSON.parse(readFileSync(path, "utf8")) as T;
|
|
||||||
} catch (err) {
|
|
||||||
if (whenMissing !== null) return whenMissing;
|
|
||||||
console.error(`network-checker: cannot read ${path}: ${(err as Error).message}`);
|
|
||||||
process.exit(1);
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
function writeAtomically(path: string, body: string): void {
|
|
||||||
mkdirSync(dirname(path), { recursive: true });
|
|
||||||
const temp = `${path}.writing`;
|
|
||||||
writeFileSync(temp, body);
|
|
||||||
renameSync(temp, path);
|
|
||||||
}
|
|
||||||
|
|
||||||
async function main(): Promise<void> {
|
|
||||||
const roster = read<Roster>(rosterFile, null);
|
|
||||||
if (!roster.machines?.length) {
|
|
||||||
console.error("network-checker: the roster names no machines; nothing to check");
|
|
||||||
process.exit(1);
|
|
||||||
}
|
|
||||||
|
|
||||||
const targets = targetsFor(roster, probePort, publicPort);
|
|
||||||
// In parallel, because a machine that is away should not delay the rest: a run that takes
|
|
||||||
// machines × timeout would outlast its own cadence on a mesh of any size.
|
|
||||||
const results: Result[] = await Promise.all(targets.map((t) => dial(t, timeoutMs)));
|
|
||||||
|
|
||||||
const countsFile = join(stateDir, "consecutive.json");
|
|
||||||
const { counts, broken } = tally(results, read<Counts>(countsFile, {}), threshold);
|
|
||||||
writeAtomically(countsFile, JSON.stringify(counts, null, 1));
|
|
||||||
|
|
||||||
// Written whole, every run: a reader asking "what does this machine reach" gets an answer about now
|
|
||||||
// rather than the last time something changed.
|
|
||||||
writeAtomically(join(stateDir, "reach.json"), JSON.stringify({
|
|
||||||
node: roster.node,
|
|
||||||
at: new Date().toISOString(),
|
|
||||||
checked: results.length,
|
|
||||||
broken: broken.length,
|
|
||||||
results,
|
|
||||||
}, null, 1));
|
|
||||||
|
|
||||||
for (const b of broken) {
|
|
||||||
console.error(
|
|
||||||
`network-checker: ${roster.node} cannot reach ${b.machine} (${b.claim}) at ${b.at}:${b.port} — ` +
|
|
||||||
`${b.failed} failed${b.detail ? `: ${b.detail}` : ""}, ${b.consecutive} run(s) running`);
|
|
||||||
}
|
|
||||||
if (broken.length === 0) {
|
|
||||||
console.log(`network-checker: ${roster.node} reaches all ${results.length} checked path(s)`);
|
|
||||||
}
|
|
||||||
|
|
||||||
// A broken path is not this process failing. It did its job; exiting non-zero would make the mesh
|
|
||||||
// read the checker as the fault, and a scheduled step that fails is retried rather than believed.
|
|
||||||
process.exit(0);
|
|
||||||
}
|
|
||||||
|
|
||||||
void main();
|
|
||||||
@@ -1,61 +0,0 @@
|
|||||||
{
|
|
||||||
"module": "network-checker",
|
|
||||||
"version": "1",
|
|
||||||
"slug": "netcheck",
|
|
||||||
"listens": [
|
|
||||||
{
|
|
||||||
"name": "probe",
|
|
||||||
"port": 9876,
|
|
||||||
"protocol": "tcp",
|
|
||||||
"from": "mesh",
|
|
||||||
"why": "what the other machines' checkers dial. Deliberately this module's own endpoint and nothing else's: it is admitted by exactly the rule that governs every internally-exposed service, so it fails when that rule is wrong. A probe on a port that is never closed — ssh, say — would have passed throughout the outage this module exists to catch (novox/hq ADR 0145)"
|
|
||||||
}
|
|
||||||
],
|
|
||||||
"facts": {
|
|
||||||
"roster": {
|
|
||||||
"path": "/var/lib/network-checker/roster.json",
|
|
||||||
"template": "{\n \"generated\": \"by the mesh — do not edit; replaced whenever a machine joins or leaves\",\n \"node\": \"{{.Node}}\",\n \"machines\": [{{range $i, $m := .Machines}}{{if $i}},{{end}}\n { \"name\": \"{{$m.Name}}\", \"fqdn\": \"{{$m.FQDN}}\", \"address\": \"{{$m.Address}}\" }{{end}}\n ]\n}\n"
|
|
||||||
}
|
|
||||||
},
|
|
||||||
"build": {
|
|
||||||
"artifacts": [
|
|
||||||
{
|
|
||||||
"name": "code",
|
|
||||||
"kind": "bundle",
|
|
||||||
"language": "typescript",
|
|
||||||
"entrypoints": ["probe/index.js", "check/index.js"]
|
|
||||||
}
|
|
||||||
]
|
|
||||||
},
|
|
||||||
"resources": [
|
|
||||||
{
|
|
||||||
"id": "state",
|
|
||||||
"type": "directory",
|
|
||||||
"path": "/var/lib/network-checker",
|
|
||||||
"mode": "0700"
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"id": "probe",
|
|
||||||
"type": "container",
|
|
||||||
"name": "mesh-network-checker-probe",
|
|
||||||
"args": ["run", "/app/modules/network-checker/dist/probe/index.js"],
|
|
||||||
"ports": ["9876"],
|
|
||||||
"state": "running",
|
|
||||||
"env": { "MESH_NETWORK_CHECKER_PORT": "9876" }
|
|
||||||
},
|
|
||||||
{
|
|
||||||
"id": "check",
|
|
||||||
"type": "container",
|
|
||||||
"name": "mesh-network-checker-check",
|
|
||||||
"network": "host",
|
|
||||||
"schedule": "*/5 * * * *",
|
|
||||||
"args": ["run", "/app/modules/network-checker/dist/check/index.js"],
|
|
||||||
"volumes": ["/var/lib/network-checker:/run/state"],
|
|
||||||
"env": {
|
|
||||||
"MESH_NETWORK_CHECKER_ROSTER": "/run/state/roster.json",
|
|
||||||
"MESH_NETWORK_CHECKER_STATE": "/run/state",
|
|
||||||
"MESH_NETWORK_CHECKER_PORT": "${port:9876}"
|
|
||||||
}
|
|
||||||
}
|
|
||||||
]
|
|
||||||
}
|
|
||||||
@@ -1,8 +0,0 @@
|
|||||||
{
|
|
||||||
"name": "@novox/module-network-checker",
|
|
||||||
"version": "0.1.0",
|
|
||||||
"description": "network-checker — dials what the mesh claims is reachable, from where the callers are, and says what it found.",
|
|
||||||
"type": "module",
|
|
||||||
"private": true,
|
|
||||||
"devDependencies": { "@types/node": "^22.0.0", "typescript": "^5.6.0" }
|
|
||||||
}
|
|
||||||
@@ -1,31 +0,0 @@
|
|||||||
// The endpoint the other machines' checkers dial (novox/hq ADR 0145).
|
|
||||||
//
|
|
||||||
// **This module's own endpoint is the instrument.** It is declared reachable over the private network
|
|
||||||
// like any other service, so it is admitted by exactly the rule that governs every internally-exposed
|
|
||||||
// service and it fails when that rule is wrong. A probe on a port that is never closed — ssh, say —
|
|
||||||
// would have passed throughout the outage this module exists to catch.
|
|
||||||
//
|
|
||||||
// It accepts a connection and closes it. Answering anything would make this a protocol, and then the
|
|
||||||
// question would be whether the protocol worked rather than whether the path did.
|
|
||||||
|
|
||||||
import { createServer } from "node:net";
|
|
||||||
|
|
||||||
const port = Number(process.env.MESH_NETWORK_CHECKER_PORT ?? "9876");
|
|
||||||
|
|
||||||
const server = createServer((socket) => {
|
|
||||||
// Written before closing so a person dialling it by hand sees something, and so a half-open
|
|
||||||
// connection is not mistaken for a working path by a client that only checks the handshake.
|
|
||||||
socket.end("mesh network-checker\n");
|
|
||||||
});
|
|
||||||
|
|
||||||
server.on("error", (err: Error) => {
|
|
||||||
// Said and fatal: a probe that cannot listen must not look like a probe that nothing dialled.
|
|
||||||
console.error(`network-checker: cannot serve the probe on ${port}: ${err.message}`);
|
|
||||||
process.exit(1);
|
|
||||||
});
|
|
||||||
|
|
||||||
server.listen(port, () => console.log(`network-checker: probe listening on ${port}`));
|
|
||||||
|
|
||||||
for (const signal of ["SIGTERM", "SIGINT"] as const) {
|
|
||||||
process.on(signal, () => server.close(() => process.exit(0)));
|
|
||||||
}
|
|
||||||
@@ -1,155 +0,0 @@
|
|||||||
// What the mesh claims is reachable, and how to find out (novox/hq ADR 0145).
|
|
||||||
//
|
|
||||||
// The mesh asserts three things are callable (ADR 0144): what runs on the same machine, another
|
|
||||||
// machine's service exposed to the private network, and another machine's service exposed publicly.
|
|
||||||
// This decides what to dial for each and reads the answers. It opens connections and nothing more —
|
|
||||||
// the module that owns a service is the one that knows whether it is working.
|
|
||||||
//
|
|
||||||
// **The target is this module's own endpoint, and that is deliberate.** The obvious thing to dial is a
|
|
||||||
// service every machine has, and the services every machine has are the ones never closed — ssh above
|
|
||||||
// all. Dialling one of those would have passed throughout the outage this exists to catch, because what
|
|
||||||
// broke was a service exposed to the private network and ssh is admitted unconditionally. A probe on a
|
|
||||||
// port that cannot fail measures nothing.
|
|
||||||
|
|
||||||
import { connect } from "node:net";
|
|
||||||
import { lookup } from "node:dns";
|
|
||||||
|
|
||||||
/** One machine as the mesh's roster describes it. */
|
|
||||||
export interface Machine {
|
|
||||||
name: string;
|
|
||||||
fqdn: string;
|
|
||||||
address: string;
|
|
||||||
/** The name this machine is reached by from outside, where it has one. Absent for most machines, and
|
|
||||||
* a machine with no public face has no public claim to check. */
|
|
||||||
public?: string;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** The roster the mesh renders for this module: who this machine is, and who the others are. */
|
|
||||||
export interface Roster {
|
|
||||||
node: string;
|
|
||||||
machines: Machine[];
|
|
||||||
}
|
|
||||||
|
|
||||||
/** Which of the mesh's three claims a check is about, so a failure says which one broke. */
|
|
||||||
export type Claim = "this machine" | "the private network" | "the public network";
|
|
||||||
|
|
||||||
/** One thing to dial. */
|
|
||||||
export interface Target {
|
|
||||||
claim: Claim;
|
|
||||||
machine: string;
|
|
||||||
/** What to dial — a name where the point is that names resolve, an address where it is not. */
|
|
||||||
at: string;
|
|
||||||
port: number;
|
|
||||||
/** Whether `at` is a name that must resolve first, so a resolution failure is reported as one. */
|
|
||||||
byName: boolean;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** What one dial found. */
|
|
||||||
export interface Result extends Target {
|
|
||||||
ok: boolean;
|
|
||||||
/** Which step failed, so a reader is sent to the right place: the resolver, or the filter. */
|
|
||||||
failed?: "resolution" | "connection";
|
|
||||||
detail?: string;
|
|
||||||
ms: number;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* targetsFor is everything this machine should be able to reach, from the roster it was given.
|
|
||||||
*
|
|
||||||
* Its own machine first, because that is the case that distinguishes a caller on the machine from a
|
|
||||||
* caller in one of its containers — the one that broke. Then every other machine over the private
|
|
||||||
* network. The public claim is only checked where a public address is known for a machine, because a
|
|
||||||
* machine with no public face has nothing to fail.
|
|
||||||
*/
|
|
||||||
export function targetsFor(roster: Roster, probePort: number, publicPort?: number): Target[] {
|
|
||||||
const out: Target[] = [];
|
|
||||||
for (const m of roster.machines) {
|
|
||||||
const own = m.name === roster.node;
|
|
||||||
out.push({
|
|
||||||
claim: own ? "this machine" : "the private network",
|
|
||||||
machine: m.name,
|
|
||||||
at: m.address,
|
|
||||||
port: probePort,
|
|
||||||
byName: false,
|
|
||||||
});
|
|
||||||
// And by name, because a name that does not resolve and a port that does not answer are different
|
|
||||||
// faults with different owners.
|
|
||||||
out.push({
|
|
||||||
claim: own ? "this machine" : "the private network",
|
|
||||||
machine: m.name,
|
|
||||||
at: m.fqdn,
|
|
||||||
port: probePort,
|
|
||||||
byName: true,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
if (publicPort !== undefined) {
|
|
||||||
for (const m of roster.machines) {
|
|
||||||
if (!m.public) continue;
|
|
||||||
out.push({
|
|
||||||
claim: "the public network",
|
|
||||||
machine: m.name,
|
|
||||||
at: m.public,
|
|
||||||
port: publicPort,
|
|
||||||
byName: true,
|
|
||||||
});
|
|
||||||
}
|
|
||||||
}
|
|
||||||
return out;
|
|
||||||
}
|
|
||||||
|
|
||||||
/** dial opens a connection and closes it. Whether the port accepts is the whole of what is asked. */
|
|
||||||
export function dial(target: Target, timeoutMs: number): Promise<Result> {
|
|
||||||
const began = Date.now();
|
|
||||||
const done = (ok: boolean, failed?: Result["failed"], detail?: string): Result => ({
|
|
||||||
...target, ok, failed, detail, ms: Date.now() - began,
|
|
||||||
});
|
|
||||||
|
|
||||||
return new Promise<Result>((resolve) => {
|
|
||||||
const open = () => {
|
|
||||||
const socket = connect({ host: target.at, port: target.port });
|
|
||||||
const finish = (r: Result) => { socket.destroy(); resolve(r); };
|
|
||||||
socket.setTimeout(timeoutMs);
|
|
||||||
socket.once("connect", () => finish(done(true)));
|
|
||||||
socket.once("timeout", () => finish(done(false, "connection", "timed out")));
|
|
||||||
socket.once("error", (err: Error) => finish(done(false, "connection", err.message)));
|
|
||||||
};
|
|
||||||
|
|
||||||
if (!target.byName) { open(); return; }
|
|
||||||
// Resolved first and reported separately: a checker that says "unreachable" for a name the
|
|
||||||
// resolver never answered sends a reader to the filter, which is not where the fault is.
|
|
||||||
lookup(target.at, (err) => {
|
|
||||||
if (err) { resolve(done(false, "resolution", err.message)); return; }
|
|
||||||
open();
|
|
||||||
});
|
|
||||||
});
|
|
||||||
}
|
|
||||||
|
|
||||||
/** A path's running count of consecutive failures, keyed so it survives between runs. */
|
|
||||||
export type Counts = Record<string, number>;
|
|
||||||
|
|
||||||
/** keyOf names one path, stably, so a count follows it across runs. */
|
|
||||||
export function keyOf(t: Target): string {
|
|
||||||
return `${t.claim}|${t.machine}|${t.at}|${t.port}`;
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* tally folds this run's results into the counts carried from the last one.
|
|
||||||
*
|
|
||||||
* **One failure is not a fault.** A machine rebooting is ordinary, and a checker that cries at the
|
|
||||||
* first missed dial trains a reader to ignore it — which is worse than not checking (ADR 0145). A path
|
|
||||||
* is broken once it has failed on consecutive runs, and the count travels with the result so a reader
|
|
||||||
* can tell "briefly away" from "never worked".
|
|
||||||
*/
|
|
||||||
export function tally(results: Result[], before: Counts, threshold: number): {
|
|
||||||
counts: Counts; broken: Array<Result & { consecutive: number }>;
|
|
||||||
} {
|
|
||||||
const counts: Counts = {};
|
|
||||||
const broken: Array<Result & { consecutive: number }> = [];
|
|
||||||
for (const r of results) {
|
|
||||||
const key = keyOf(r);
|
|
||||||
const n = r.ok ? 0 : (before[key] ?? 0) + 1;
|
|
||||||
if (n > 0) counts[key] = n;
|
|
||||||
if (n >= threshold) broken.push({ ...r, consecutive: n });
|
|
||||||
}
|
|
||||||
return { counts, broken };
|
|
||||||
}
|
|
||||||
@@ -1,72 +0,0 @@
|
|||||||
import { strict as assert } from "node:assert";
|
|
||||||
import test from "node:test";
|
|
||||||
|
|
||||||
import { keyOf, tally, targetsFor, type Result, type Roster } from "../reach.js";
|
|
||||||
|
|
||||||
const roster: Roster = {
|
|
||||||
node: "here",
|
|
||||||
machines: [
|
|
||||||
{ name: "here", fqdn: "here.internal", address: "10.0.0.1" },
|
|
||||||
{ name: "there", fqdn: "there.internal", address: "10.0.0.2", public: "there.example.test" },
|
|
||||||
],
|
|
||||||
};
|
|
||||||
|
|
||||||
test("its own machine is checked, which is the case that distinguishes a caller on it from one in a container", () => {
|
|
||||||
const own = targetsFor(roster, 9876).filter((t) => t.claim === "this machine");
|
|
||||||
assert.equal(own.length, 2, "its own machine by address and by name");
|
|
||||||
assert.ok(own.some((t) => t.at === "10.0.0.1" && !t.byName));
|
|
||||||
assert.ok(own.some((t) => t.at === "here.internal" && t.byName));
|
|
||||||
});
|
|
||||||
|
|
||||||
test("every other machine is checked over the private network", () => {
|
|
||||||
const other = targetsFor(roster, 9876).filter((t) => t.claim === "the private network");
|
|
||||||
assert.deepEqual(other.map((t) => t.machine), ["there", "there"]);
|
|
||||||
});
|
|
||||||
|
|
||||||
test("the public claim is only checked where a machine has a public name", () => {
|
|
||||||
const pub = targetsFor(roster, 9876, 443).filter((t) => t.claim === "the public network");
|
|
||||||
assert.equal(pub.length, 1, "only the machine with a public name");
|
|
||||||
assert.equal(pub[0]!.at, "there.example.test");
|
|
||||||
assert.equal(pub[0]!.port, 443);
|
|
||||||
});
|
|
||||||
|
|
||||||
test("no public claim is made when no public port was given", () => {
|
|
||||||
assert.equal(targetsFor(roster, 9876).filter((t) => t.claim === "the public network").length, 0);
|
|
||||||
});
|
|
||||||
|
|
||||||
const failed = (at: string): Result => ({
|
|
||||||
claim: "this machine", machine: "here", at, port: 9876, byName: false,
|
|
||||||
ok: false, failed: "connection", ms: 1,
|
|
||||||
});
|
|
||||||
const passed = (at: string): Result => ({
|
|
||||||
claim: "this machine", machine: "here", at, port: 9876, byName: false, ok: true, ms: 1,
|
|
||||||
});
|
|
||||||
|
|
||||||
test("one failure is not a fault — a machine rebooting is ordinary", () => {
|
|
||||||
const { counts, broken } = tally([failed("10.0.0.1")], {}, 2);
|
|
||||||
assert.equal(broken.length, 0, "one missed dial says nothing");
|
|
||||||
assert.equal(counts[keyOf(failed("10.0.0.1"))], 1, "and is remembered");
|
|
||||||
});
|
|
||||||
|
|
||||||
test("a path that keeps failing is broken, and the count travels with it", () => {
|
|
||||||
const first = tally([failed("10.0.0.1")], {}, 2);
|
|
||||||
const second = tally([failed("10.0.0.1")], first.counts, 2);
|
|
||||||
assert.equal(second.broken.length, 1);
|
|
||||||
assert.equal(second.broken[0]!.consecutive, 2, "so a reader can tell briefly away from never worked");
|
|
||||||
});
|
|
||||||
|
|
||||||
test("a path that recovers stops being counted", () => {
|
|
||||||
const first = tally([failed("10.0.0.1")], {}, 2);
|
|
||||||
const second = tally([passed("10.0.0.1")], first.counts, 2);
|
|
||||||
assert.equal(second.broken.length, 0);
|
|
||||||
assert.deepEqual(second.counts, {}, "nothing carried forward for a path that works");
|
|
||||||
});
|
|
||||||
|
|
||||||
test("a count follows one path and not another", () => {
|
|
||||||
const a = failed("10.0.0.1");
|
|
||||||
const b = failed("10.0.0.2");
|
|
||||||
const first = tally([a, b], {}, 2);
|
|
||||||
const second = tally([a], first.counts, 2);
|
|
||||||
assert.equal(second.broken.length, 1, "only the path dialled this run is judged");
|
|
||||||
assert.equal(second.broken[0]!.at, "10.0.0.1");
|
|
||||||
});
|
|
||||||
@@ -1,12 +0,0 @@
|
|||||||
{
|
|
||||||
"compilerOptions": {
|
|
||||||
"target": "ES2022",
|
|
||||||
"module": "NodeNext",
|
|
||||||
"moduleResolution": "NodeNext",
|
|
||||||
"strict": true,
|
|
||||||
"esModuleInterop": true,
|
|
||||||
"skipLibCheck": true,
|
|
||||||
"noEmit": true
|
|
||||||
},
|
|
||||||
"include": ["reach.ts", "probe/index.ts", "check/index.ts", "test/*.ts"]
|
|
||||||
}
|
|
||||||
+21
-20
@@ -5,7 +5,7 @@
|
|||||||
"container-runtime"
|
"container-runtime"
|
||||||
],
|
],
|
||||||
"own-secrets": {
|
"own-secrets": {
|
||||||
"secret": "/var/lib/searxng-module/secret.secret",
|
"secret": "/var/lib/mesh/searxng/secret",
|
||||||
"broker": "/var/lib/mesh/searxng/broker"
|
"broker": "/var/lib/mesh/searxng/broker"
|
||||||
},
|
},
|
||||||
"listens": [
|
"listens": [
|
||||||
@@ -27,22 +27,14 @@
|
|||||||
{
|
{
|
||||||
"id": "state",
|
"id": "state",
|
||||||
"type": "directory",
|
"type": "directory",
|
||||||
"path": "/var/lib/searxng-module",
|
"mode": "0700",
|
||||||
"mode": "0700"
|
"place": "."
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"id": "valkey-data",
|
"id": "valkey-data",
|
||||||
"type": "directory",
|
"type": "directory",
|
||||||
"path": "/var/lib/searxng-module/valkey-data",
|
|
||||||
"mode": "0700"
|
"mode": "0700"
|
||||||
},
|
},
|
||||||
{
|
|
||||||
"id": "server-env",
|
|
||||||
"type": "file",
|
|
||||||
"path": "/var/lib/searxng-module/server.env",
|
|
||||||
"mode": "0600",
|
|
||||||
"content": "SEARXNG_SECRET=${secret:secret}\nSEARXNG_VALKEY_URL=valkey://valkey:6379/0\n"
|
|
||||||
},
|
|
||||||
{
|
{
|
||||||
"id": "net",
|
"id": "net",
|
||||||
"type": "network",
|
"type": "network",
|
||||||
@@ -63,30 +55,39 @@
|
|||||||
"warning"
|
"warning"
|
||||||
],
|
],
|
||||||
"volumes": [
|
"volumes": [
|
||||||
"/var/lib/searxng-module/valkey-data:/data"
|
"${dir:valkey-data}:/data"
|
||||||
]
|
]
|
||||||
},
|
},
|
||||||
|
{
|
||||||
|
"id": "settings",
|
||||||
|
"type": "file",
|
||||||
|
"path": "${dir:state}/settings.yml",
|
||||||
|
"mode": "0600",
|
||||||
|
"merge": "json",
|
||||||
|
"content": "{\n \"use_default_settings\": true,\n \"server\": {\n \"secret_key\": \"${secret:secret}\",\n \"base_url\": false,\n \"limiter\": false,\n \"image_proxy\": false,\n \"public_instance\": false\n },\n \"search\": {\n \"formats\": [\"html\", \"json\"]\n },\n \"valkey\": {\n \"url\": \"valkey://valkey:6379/0\"\n }\n}\n"
|
||||||
|
},
|
||||||
{
|
{
|
||||||
"id": "server",
|
"id": "server",
|
||||||
"type": "container",
|
"type": "container",
|
||||||
"name": "searxng",
|
"name": "searxng",
|
||||||
"image": "searxng/searxng@sha256:c7cc75852051bf6254afda6ed1b920dd1677d8efe4ab141bf558f02e582f4371",
|
"image": "searxng/searxng@sha256:cd8812607ab73730a0b1a0dc4990223fe1b9e383f6f35947114d0bef7f8bb441",
|
||||||
"network": "searxng",
|
"network": "searxng",
|
||||||
"env-file": [
|
|
||||||
"/var/lib/searxng-module/server.env"
|
|
||||||
],
|
|
||||||
"ports": [
|
"ports": [
|
||||||
"8080"
|
"8080"
|
||||||
],
|
],
|
||||||
"secrets-in-environment": "SEARXNG_SECRET is env-only, but settings.yml carries server.secret_key; convertible by mounting a generated settings.yml, not yet done"
|
"volumes": [
|
||||||
|
"${dir:state}/settings.yml:/etc/searxng/settings.yml:ro"
|
||||||
|
],
|
||||||
|
"restart-on": [
|
||||||
|
"settings"
|
||||||
|
]
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"id": "runtime-config",
|
"id": "runtime-config",
|
||||||
"type": "file",
|
"type": "file",
|
||||||
"path": "/var/lib/mesh/searxng/config.json",
|
"path": "/var/lib/mesh/searxng/config.json",
|
||||||
"mode": "0600",
|
"mode": "0600",
|
||||||
"content": "{}\n",
|
"content": "{}\n"
|
||||||
"merge": "json"
|
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"id": "runtime",
|
"id": "runtime",
|
||||||
@@ -118,7 +119,7 @@
|
|||||||
}
|
}
|
||||||
},
|
},
|
||||||
"binds": {
|
"binds": {
|
||||||
"route": "/var/lib/searxng-module/route.json"
|
"route": "${dir:state}/route.json"
|
||||||
},
|
},
|
||||||
"build": {
|
"build": {
|
||||||
"on": [
|
"on": [
|
||||||
|
|||||||
Reference in New Issue
Block a user