Compare commits

..
Author SHA1 Message Date
jschoubben 8f459c7023 letta: placed state, a route, accepted keys, and a runtime that can log in
The module named /var/lib/letta, a layout no definition may carry (ADR
0112); state is now a placed directory holding the bindings and secrets.

letta had no route, while the server it replaces is reached by its public
name (a workflow calls it there). It now contributes one for its `web`
endpoint; reach is the assignment's.

The server needs an OpenAI key for agents on OpenAI models, and nothing
gave it one: `openai-api-key` is an own-secret, accepted from the
operator (a key someone chose, not one the mesh can mint). The
server-password is accepted the same way where a server already has
clients. Both, and the database password inside LETTA_PG_URI, stay in the
server's environment: letta 0.6.x reads settings from the environment
only, and its startup.sh starts an embedded PostgreSQL unless
LETTA_PG_URI is set - the declared reason now says so.

The runtime's tools never authenticated: the client sent only a Bearer
token, and 0.6.x's --secure mode checks X-BARE-PASSWORD ("password <it>")
and answers 401 otherwise. The client now sends both. Its password comes
from the runtime config file (the key client.ts reads first) instead of an
env-file, so the runtime container no longer carries a secret in its
environment.

Image: the same 0.6.8 image, now pinned by the index digest ace runs
rather than its amd64 manifest.

Verified: catalogue tests with MESH_CATALOGUE set; tsc -p tsconfig.json in
the mesh-tools build image. Throwaway containers: a fresh 0.6.8 with its
embedded PG and two blocks made through the API; stopped, copied, dumped
from the copy; restored (schema letta + vector pre-made by the superuser,
--no-owner --role, search_path set on the database as the original had
it) into a grant-shaped database on the postgres module's pgvector image
(PG17); started in this shape: alembic finds nothing to do, both blocks
are there, a wrong password is refused with 401, and the patched client
lists agents. Test containers and data removed.
2026-09-30 12:22:12 +02:00
11 changed files with 42 additions and 744 deletions
+8 -4
View File
@@ -1,10 +1,10 @@
// The Letta API client — letta's own code, living in the module (novox/hq ADR 0039). Its tools
// import it; nothing outside letta does.
//
// Letta authenticates with a single server password, presented as a Bearer token. That password is
// a mesh own-secret, minted once and handed to both the server (LETTA_SERVER_PASSWORD) and this
// client (MESH_LETTA_PASSWORD) — so the module's tools are live without anything configured by hand.
// The runtime config file may still override the URL or password.
// Letta authenticates with a single server password. That password is a mesh own-secret handed to
// both the server (LETTA_SERVER_PASSWORD) and this client, through the runtime config file the mesh
// mounts (its `password` key) — so the module's tools are live without anything configured by hand.
// Where a server already has clients, the password is accepted rather than minted.
import { readFileSync } from "node:fs";
@@ -58,6 +58,10 @@ export class LettaClient {
...options,
headers: {
"Content-Type": "application/json",
// The server's --secure mode checks X-BARE-PASSWORD ("password <it>") and answers a Bearer
// token alone with 401 (letta/server/rest_api/app.py, 0.6.x). Both are sent: Bearer is what
// later servers read.
"X-BARE-PASSWORD": `password ${this.password}`,
Authorization: `Bearer ${this.password}`,
...(options.headers as Record<string, string> | undefined),
},
+21 -25
View File
@@ -5,21 +5,28 @@
"container-runtime"
],
"requires": [
"postgres-database"
"postgres-database",
"route"
],
"contributes": {
"postgres-database": {
"name": "letta"
},
"route": {
"label": "letta",
"endpoint": "web"
}
},
"binds": {
"postgres-database": "/var/lib/letta/database.json"
"postgres-database": "${dir:state}/database.json",
"route": "${dir:state}/route.json"
},
"secrets": {
"postgres-database": "/var/lib/letta/database.secret"
"postgres-database": "${dir:state}/database.secret"
},
"own-secrets": {
"server-password": "/var/lib/letta/server-password.secret",
"server-password": "${dir:state}/server-password.secret",
"openai-api-key": "${dir:state}/openai-api-key.secret",
"broker": "/var/lib/mesh/letta/broker"
},
"listens": [
@@ -28,7 +35,7 @@
"port": 8283,
"protocol": "tcp",
"from": "mesh",
"why": "the Letta agent server REST API and web UI; a public name is a route grant later"
"why": "the Letta agent server REST API and web UI, password-protected (--secure); a public name is the route's"
}
],
"resources": [
@@ -41,15 +48,15 @@
{
"id": "state",
"type": "directory",
"path": "/var/lib/letta",
"mode": "0700"
"mode": "0700",
"place": "."
},
{
"id": "server-env",
"type": "file",
"path": "/var/lib/letta/server.env",
"path": "${dir:state}/server.env",
"mode": "0600",
"content": "LETTA_PG_URI=postgresql://${bound:postgres-database:as}:${secret:postgres-database}@${bound:postgres-database:at}:${bound:postgres-database:port}/${bound:postgres-database:as}\nLETTA_SERVER_PASSWORD=${secret:server-password}\nSECURE=true\nTZ=Europe/Brussels\n"
"content": "LETTA_PG_URI=postgresql://${bound:postgres-database:as}:${secret:postgres-database}@${bound:postgres-database:at}:${bound:postgres-database:port}/${bound:postgres-database:as}\nLETTA_SERVER_PASSWORD=${secret:server-password}\nOPENAI_API_KEY=${secret:openai-api-key}\nSECURE=true\nTZ=Europe/Brussels\n"
},
{
"id": "net",
@@ -60,31 +67,24 @@
"id": "server",
"type": "container",
"name": "letta",
"image": "letta/letta@sha256:1d2e0692514287c5ed1a483e14e16ed945f8632d315539f5e66373bb7d7c471b",
"image": "letta/letta@sha256:bfd1e49ce45b9a208c941e832c1d1d194017ff210a3784b0ca6c323aed767a29",
"network": "letta",
"env-file": [
"/var/lib/letta/server.env"
"${dir:state}/server.env"
],
"ports": [
"8283"
],
"secrets-in-environment": "the letta image is env-driven and its file-source support could not be verified; the mesh runtime can take its password from config.json (client.ts) \u2014 not yet converted"
"secrets-in-environment": "letta 0.6.x reads its settings from the environment only (pydantic settings, no secrets_dir or _FILE twin), and its startup.sh starts an embedded PostgreSQL unless LETTA_PG_URI is set - so the database password travels inside that URI (startup.sh also echoes it to the log); LETTA_SERVER_PASSWORD and OPENAI_API_KEY have no file source either"
},
{
"id": "runtime-config",
"type": "file",
"path": "/var/lib/mesh/letta/config.json",
"mode": "0600",
"content": "{}\n",
"content": "{\n \"password\": \"${secret:server-password}\"\n}\n",
"merge": "json"
},
{
"id": "runtime-env",
"type": "file",
"path": "/var/lib/letta/runtime.env",
"mode": "0600",
"content": "MESH_LETTA_PASSWORD=${secret:server-password}\n"
},
{
"id": "runtime",
"type": "container",
@@ -99,14 +99,10 @@
"MESH_LETTA_URL": "http://letta:8283",
"MESH_LETTA_CONFIG_FILE": "/run/config/config.json"
},
"env-file": [
"/var/lib/letta/runtime.env"
],
"restart-on": [
"runtime-config"
],
"artifact": "runtime",
"secrets-in-environment": "the letta image is env-driven and its file-source support could not be verified; the mesh runtime can take its password from config.json (client.ts) \u2014 not yet converted"
"artifact": "runtime"
}
],
"build": {
+1 -4
View File
@@ -13,7 +13,7 @@ ARG RUNTIME_BASE
FROM ${BUILD_BASE} AS build
WORKDIR /app/modules/nodered
COPY . .
RUN node /app/node_modules/typescript/bin/tsc client.ts tools/index.ts mqtt/probe.ts mqtt/connection.ts mqtt/index.ts \
RUN node /app/node_modules/typescript/bin/tsc client.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE}
@@ -22,6 +22,3 @@ COPY --from=build /app/modules/nodered/dist /app/modules/nodered/dist
# provider's provisioner runs its reconcile loop in the same process, with the broker connected —
# the convention novox/hq issues 060/061 settled.
ENV MESH_TOOL_MODULES=/app/modules/nodered/dist/tools/index.js
# NOT dist/mqtt/index.js: that is a step the host runs to completion, named by the `mqtt`
# container's args as `mesh-tools run …` (novox/hq ADR 0052). Listed here it would run inside the
# serving sidecar too, and exit it.
+3 -38
View File
@@ -3,8 +3,7 @@
//
// Node-RED exposes a runtime admin API under its base URL: GET/POST /flows for the whole flow
// configuration, GET /nodes for installed node modules. A default install has no auth; when
// adminAuth is on, a bearer token is required — the module's settings accept the mesh-minted
// api-token, which the runtime config file carries as `token`.
// adminAuth is on, a bearer token (minted at /auth/token) is required.
import { readFileSync } from "node:fs";
@@ -82,35 +81,6 @@ export class NodeRedClient {
return modules.map((m) => ({ name: m.name, version: m.version, types: m.types ?? [] }));
}
/** The whole flow configuration with its revision (API v2), for a deploy that must not clobber
* a change made meanwhile. */
async flowsWithRev(): Promise<{ rev: string; flows: any[] }> {
const body = await this.req("/flows", { headers: this.headers({ "Node-RED-API-Version": "v2" }) });
return { rev: String(body?.rev ?? ""), flows: Array.isArray(body?.flows) ? body.flows : [] };
}
/** A node's stored credentials as Node-RED shows them: plain fields, and `has_<field>` for secret ones. */
async credentials(type: string, id: string): Promise<{ user?: string; has_password?: boolean }> {
return (await this.req(`/credentials/${encodeURIComponent(type)}/${encodeURIComponent(id)}`, { headers: this.headers() })) ?? {};
}
/**
* Deploy the flow configuration read at `rev`. Node-RED answers 409 when the flows changed since,
* rather than overwriting what someone deployed in between. A node carrying `credentials` has them
* stored (encrypted) and counts as changed, so a "nodes" deploy restarts it and nothing else.
*/
async deployFlowsAt(rev: string, config: any[], type = "nodes"): Promise<void> {
await this.req("/flows", {
method: "POST",
headers: this.headers({
"Content-Type": "application/json",
"Node-RED-API-Version": "v2",
"Node-RED-Deployment-Type": type,
}),
body: JSON.stringify({ rev, flows: config }),
});
}
/**
* Replace the whole flow configuration and deploy. Returns the new revision. `type` maps to
* Node-RED's deployment types — "full" (default), "nodes", or "flows".
@@ -118,13 +88,8 @@ export class NodeRedClient {
async deployFlows(config: any[], type = "full"): Promise<{ rev?: string; nodeCount: number }> {
const body = await this.req("/flows", {
method: "POST",
// v2 answers { rev }; v1 answers 204 with no body, which req() cannot parse.
headers: this.headers({
"Content-Type": "application/json",
"Node-RED-API-Version": "v2",
"Node-RED-Deployment-Type": type,
}),
body: JSON.stringify({ flows: config }),
headers: this.headers({ "Content-Type": "application/json", "Node-RED-Deployment-Type": type }),
body: JSON.stringify(config),
});
return { rev: body?.rev, nodeCount: config.length };
}
+7 -85
View File
@@ -5,8 +5,6 @@
"flows.deployed"
],
"own-secrets": {
"admin": "/var/lib/mesh/nodered/admin",
"api-token": "/var/lib/mesh/nodered/api-token",
"broker": "/var/lib/mesh/nodered/broker"
},
"capabilities": [
@@ -28,63 +26,26 @@
"path": "/var/lib/mesh/nodered",
"mode": "0700"
},
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
},
{
"id": "data",
"type": "directory",
"path": "/services/nodered/data",
"mode": "0700",
"owner": "1000:1000"
},
{
"id": "written",
"type": "directory",
"mode": "0700"
},
{
"id": "settings-code",
"type": "file",
"path": "${dir:state}/settings.js",
"mode": "0600",
"owner": "1000:1000",
"content": "// Node-RED's settings, written by the mesh from the nodered module. What an assignment may change\n// is settings.json beside this file (merged key by key); the credentials are the mesh's secrets and\n// reach Node-RED only through this file. The flows' own credentials stay encrypted in the user\n// directory under the key Node-RED keeps there (.config.runtime.json), which is data, not this.\nconst fs = require(\"fs\");\nconst path = require(\"path\");\nconst crypto = require(\"crypto\");\n\nconst ADMIN_PASSWORD = \"${secret:admin}\";\nconst API_TOKEN = \"${secret:api-token}\";\nconst ADMIN = { username: \"admin\", permissions: \"*\" };\n\nconst settings = JSON.parse(fs.readFileSync(path.join(__dirname, \"settings.json\"), \"utf8\"));\n// The mesh's keys, not Node-RED's: endpoints lands in every merged file; timeZone is the\n// assignment's way to set the zone flows schedule and format in; mqtt names the broker nodes the\n// module's MQTT step keeps pointed at the mesh's broker; topics is what nodered asks the broker for.\nif (settings.timeZone) process.env.TZ = settings.timeZone;\ndelete settings.timeZone;\ndelete settings.endpoints;\ndelete settings.mqtt;\ndelete settings.topics;\n\nfunction same(a, b) {\n const x = crypto.createHash(\"sha256\").update(String(a)).digest();\n const y = crypto.createHash(\"sha256\").update(String(b)).digest();\n return crypto.timingSafeEqual(x, y);\n}\n\n// The admin secret is a password, or, accepted from an existing install, the bcrypt hash its\n// settings held, so the password people already use keeps working.\nfunction passwordMatches(given) {\n if (/^\\$2[aby]\\$\\d\\d\\$/.test(ADMIN_PASSWORD)) return require(\"bcryptjs\").compare(String(given), ADMIN_PASSWORD);\n return Promise.resolve(same(given, ADMIN_PASSWORD));\n}\n\nmodule.exports = Object.assign(settings, {\n uiPort: 1880,\n adminAuth: {\n type: \"credentials\",\n users: (username) => Promise.resolve(username === ADMIN.username ? ADMIN : null),\n authenticate: (username, password) =>\n username === ADMIN.username\n ? passwordMatches(password).then((ok) => (ok ? ADMIN : null))\n : Promise.resolve(null),\n // The module's own tools call the admin API with this bearer token.\n tokens: (token) => Promise.resolve(same(token, API_TOKEN) ? { username: \"mesh\", permissions: \"*\" } : null),\n },\n});\n"
},
{
"id": "settings",
"type": "file",
"path": "${dir:state}/settings.json",
"mode": "0600",
"owner": "1000:1000",
"merge": "json",
"content": "{\n \"flowFile\": \"flows.json\",\n \"flowFilePretty\": true,\n \"diagnostics\": { \"enabled\": true, \"ui\": true },\n \"runtimeState\": { \"enabled\": false, \"ui\": false },\n \"logging\": { \"console\": { \"level\": \"info\", \"metrics\": false, \"audit\": false } },\n \"exportGlobalContextKeys\": false,\n \"externalModules\": {},\n \"editorTheme\": { \"projects\": { \"enabled\": false } },\n \"functionExternalModules\": true,\n \"debugMaxLength\": 1000,\n \"mqttReconnectTime\": 15000,\n \"serialReconnectTime\": 15000\n}\n"
},
{
"id": "server",
"type": "container",
"name": "nodered",
"image": "nodered/node-red@sha256:a649dd711d55490151a2c39a8e48ad0c44325488fbc0e66315f2d2e19e5e1ace",
"image": "nodered/node-red@sha256:02a2b92a41b73d2bc388238b86e4fcaab7fb5466373adb24e1df6aa5845265ff",
"env": {
"TZ": "Etc/UTC"
},
"ports": [
"1880"
],
"args": [
"--settings",
"/config/settings.js"
],
"volumes": [
"${dir:data}:/data",
"${dir:state}/settings.js:/config/settings.js:ro",
"${dir:state}/settings.json:/config/settings.json:ro"
],
"restart-on": [
"settings-code",
"settings"
"/services/nodered/data:/data"
]
},
{
@@ -92,7 +53,8 @@
"type": "file",
"path": "/var/lib/mesh/nodered/config.json",
"mode": "0600",
"content": "{\n \"token\": \"${secret:api-token}\"\n}\n"
"content": "{}\n",
"merge": "json"
},
{
"id": "runtime",
@@ -105,66 +67,26 @@
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_NODERED_URL": "http://127.0.0.1:${port:1880}",
"MESH_NODERED_URL": "http://127.0.0.1:1880",
"MESH_NODERED_CONFIG_FILE": "/run/config/config.json"
},
"restart-on": [
"runtime-config"
],
"artifact": "runtime"
},
{
"id": "mqtt",
"type": "container",
"name": "mesh-nodered-mqtt",
"network": "host",
"run-once": true,
"volumes": [
"/var/lib/mesh/nodered/config.json:/run/config/config.json:ro",
"${dir:written}:/var/lib/nodered-provisions",
"${dir:state}/mqtt-topic.json:/run/provisions/mqtt-topic.json:ro",
"${dir:state}/mqtt-topic.secret:/run/provisions/mqtt-topic.secret:ro",
"${dir:state}/settings.json:/run/provisions/settings.json:ro"
],
"env": {
"MESH_NODERED_URL": "http://127.0.0.1:${port:1880}",
"MESH_NODERED_CONFIG_FILE": "/run/config/config.json",
"MESH_PROVISIONS_DIR": "/run/provisions",
"MESH_WRITTEN_DIR": "/var/lib/nodered-provisions"
},
"args": [
"run",
"/app/modules/nodered/dist/mqtt/index.js"
],
"restart-on": [
"bound-mqtt-topic",
"secret-mqtt-topic",
"settings"
],
"artifact": "runtime"
}
],
"requires": [
"mqtt-topic",
"route"
],
"contributes": {
"mqtt-topic": {
"topics": [
"#"
]
},
"route": {
"label": "nodered",
"endpoint": "web"
}
},
"binds": {
"route": "${dir:state}/route.json",
"mqtt-topic": "${dir:state}/mqtt-topic.json"
},
"secrets": {
"mqtt-topic": "${dir:state}/mqtt-topic.secret"
"route": "/var/lib/mesh/nodered/route.json"
},
"build": {
"on": [
-232
View File
@@ -1,232 +0,0 @@
// Node-RED's MQTT broker config node, pointed at the broker the mesh bound — `mqtt-topic`.
//
// **Why a step.** Node-RED keeps a broker as a config node in its flows (`flows.json`) and the login
// and password in its encrypted credentials file, both of them Node-RED's to write. So this reads the
// binding and the pair credential and makes the broker node say the same thing through Node-RED's
// admin API — `GET /flows`, then `POST /flows` with the changed node and its `credentials`, deployed
// as "nodes" so only what changed restarts — with the module's own `api-token`.
//
// **Which broker nodes are the mesh's.** Never guessed: a flow may talk to a broker that has nothing
// to do with this mesh. The step owns the node it creates itself (id `mesh-mqtt-topic`, "mesh:
// mqtt-topic") and the ones an assignment names in settings (`mqtt.brokers`: node ids — how ace's
// existing broker node, which every one of its MQTT flows uses, is handed over). With none named and
// none made yet, it makes one, so a fresh Node-RED has a broker the flows can pick.
//
// **Only the connection, and only when it differs.** Host, port, TLS off (the broker serves plain
// MQTT), login, password. Every other field of the node — client id, keepalive, birth/close/will
// messages — is left as it is. Node-RED never hands a stored password back, so the step keeps a
// digest of what it last wrote: equal host/port/login and an equal digest is "already as the mesh
// says".
//
// **Nothing loses its connection without someone seeing it.** The broker is asked first whether it
// takes the delivered login; if not, nothing is written and the step fails saying why.
//
// Pure logic over two seams (Node-RED, the broker), tested against fakes (test/mqtt.test.ts).
import { createHash } from "node:crypto";
import type { Probe } from "./probe.js";
export const PROVISION = "mqtt-topic";
/** The id and name of the broker node the step makes when none is named. */
export const MESH_BROKER_ID = "mesh-mqtt-topic";
export const MESH_BROKER_NAME = "mesh: mqtt-topic";
/** What the mesh wrote at `binds.mqtt-topic`. */
export interface Binding {
provision?: string;
from?: string;
at?: string;
as?: string;
serves?: Record<string, unknown>;
}
export type Outcome =
| { what: string; result: "unchanged"; note?: string }
| { what: string; result: "written"; fields: string[]; note?: string }
| { what: string; result: "refused"; problem: string };
/** A flow node; a broker config node carries `broker`, `port`, `usetls`. */
export interface FlowNode {
id: string;
type: string;
[key: string]: unknown;
}
/** Node-RED's admin API, as the step uses it. */
export interface NodeRed {
/** The whole flow configuration and its revision (API v2). */
flows(): Promise<{ rev: string; flows: FlowNode[] }>;
/** A node's stored credentials as Node-RED shows them: the user, and only whether a password is set. */
credentials(type: string, id: string): Promise<{ user?: string; has_password?: boolean }>;
/** Deploy the configuration against the revision it was read at; "nodes" restarts only what changed. */
deploy(rev: string, flows: FlowNode[]): Promise<void>;
}
export interface Marks {
get(name: string): Promise<string | undefined>;
set(name: string, digest: string): Promise<void>;
}
export interface Deps {
nodered: NodeRed;
probe: Probe;
marks: Marks;
}
export interface Wanted {
host: string;
port: number;
user: string;
password: string;
}
export function digest(...parts: (string | number)[]): string {
return createHash("sha256").update(parts.map(String).join("\u0000")).digest("hex");
}
function isLoopback(host: string): boolean {
const h = host.toLowerCase();
return h === "localhost" || h === "::1" || h === "[::1]" || /^127\./.test(h);
}
/**
* The broker and login the mesh says Node-RED uses. A loopback `at` — what the mesh hands a machine
* that is not on the private network — is refused: from Node-RED's own container it is Node-RED.
*/
export function wanted(binding: Binding | undefined, credential: string | undefined): { ok: true; want: Wanted } | { ok: false; problem: string } {
if (!binding) return { ok: false, problem: `no binding for ${PROVISION} was delivered — the mesh writes it before this step runs` };
const host = typeof binding.at === "string" ? binding.at.trim() : "";
if (!host) return { ok: false, problem: `the ${PROVISION} binding names no host (at)` };
if (isLoopback(host)) {
return {
ok: false,
problem:
`the ${PROVISION} binding says the broker is at ${host}, which from Node-RED's own container is Node-RED ` +
`itself; put the machine on the private network so the broker has an address Node-RED can dial`,
};
}
const port = Number(binding.serves?.port);
if (!Number.isInteger(port) || port <= 0 || port > 65535) {
return { ok: false, problem: `the ${PROVISION} binding serves no usable port (${String(binding.serves?.port)})` };
}
const scheme = binding.serves?.scheme;
if (scheme !== undefined && scheme !== "mqtt") return { ok: false, problem: `the ${PROVISION} binding serves scheme ${String(scheme)}; this step writes plain MQTT` };
const user = typeof binding.as === "string" ? binding.as.trim() : "";
if (!user) return { ok: false, problem: `the ${PROVISION} binding names no login (as)` };
const password = (credential ?? "").replace(/\n$/, "");
if (!password) return { ok: false, problem: `the ${PROVISION} credential is empty or was not delivered` };
return { ok: true, want: { host, port, user, password } };
}
/** The broker node ids an assignment named in settings (`mqtt.brokers`), or none. */
export function namedBrokers(settings: unknown): string[] {
const brokers = (settings as { mqtt?: { brokers?: unknown } } | undefined)?.mqtt?.brokers;
return Array.isArray(brokers) ? brokers.filter((b): b is string => typeof b === "string" && b.length > 0) : [];
}
/** A new broker node, Node-RED 5's defaults, pointed at the broker. */
export function newBrokerNode(want: Wanted): FlowNode {
return {
id: MESH_BROKER_ID, type: "mqtt-broker", name: MESH_BROKER_NAME,
broker: want.host, port: String(want.port), clientid: "", autoConnect: true, usetls: false,
protocolVersion: "4", keepalive: "60", cleansession: true, autoUnsubscribe: true,
birthTopic: "", birthQos: "0", birthRetain: "false", birthPayload: "", birthMsg: {},
closeTopic: "", closeQos: "0", closeRetain: "false", closePayload: "", closeMsg: {},
willTopic: "", willQos: "0", willRetain: "false", willPayload: "", willMsg: {},
userProps: "", sessionExpiry: "",
};
}
const markFor = (id: string, w: Wanted): string => digest("nodered-mqtt", id, w.host, w.port, w.user, w.password);
function scrub(err: unknown, secret: string): string {
let text = err instanceof Error ? err.message : String(err);
for (const form of new Set([secret, encodeURIComponent(secret)])) text = text.split(form).join("***");
return text;
}
/**
* Bring the mesh's broker nodes in line with the binding: one outcome per node. Never throws. A node
* the settings name that is not in the flows is refused (the others are still put right).
*/
export async function reconcileBrokers(deps: Deps, binding: Binding | undefined, credential: string | undefined, named: readonly string[]): Promise<Outcome[]> {
const w = wanted(binding, credential);
if ("problem" in w) return [{ what: "mqtt", result: "refused", problem: w.problem }];
const want = w.want;
let note: string | undefined;
try {
const probe = await deps.probe(want.host, want.port, want.user, want.password, "#");
if (probe.connack === 4 || probe.connack === 5) {
return [{
what: "mqtt",
result: "refused",
problem:
`the broker at ${want.host}:${want.port} does not (yet) take the login ${want.user} with the delivered password ` +
`(CONNACK ${probe.connack}); mosquitto's provisioner creates it from the grant — nothing was written`,
}];
}
if (probe.connack !== 0) return [{ what: "mqtt", result: "refused", problem: `the broker at ${want.host}:${want.port} answered CONNACK ${probe.connack}; nothing was written` }];
if (probe.suback === 0x80) note = `warning: ${want.user} may not subscribe to every topic; flows subscribing outside its grant will get nothing`;
} catch (err) {
return [{ what: "mqtt", result: "refused", problem: `the broker at ${want.host}:${want.port} could not be asked: ${scrub(err, want.password)}; nothing was written` }];
}
const outcomes: Outcome[] = [];
for (let attempt = 0; attempt < 2; attempt++) {
outcomes.length = 0;
try {
const { rev, flows } = await deps.nodered.flows();
const targets = named.length > 0 ? [...named] : [MESH_BROKER_ID];
const changed: { id: string; fields: string[] }[] = [];
for (const id of targets) {
let node = flows.find((n) => n.id === id);
if (node && node.type !== "mqtt-broker") {
outcomes.push({ what: `broker ${id}`, result: "refused", problem: `node ${id} is a ${node.type}, not an mqtt-broker` });
continue;
}
if (!node) {
if (id !== MESH_BROKER_ID) {
outcomes.push({ what: `broker ${id}`, result: "refused", problem: `the settings name broker node ${id}, and Node-RED's flows have no such node` });
continue;
}
node = newBrokerNode(want);
flows.push(node);
node.credentials = { user: want.user, password: want.password };
changed.push({ id, fields: ["node"] });
continue;
}
const fields: string[] = [];
if (String(node.broker ?? "") !== want.host) fields.push("broker");
if (Number(node.port ?? 0) !== want.port) fields.push("port");
if (node.usetls === true) fields.push("usetls");
const creds = await deps.nodered.credentials("mqtt-broker", id);
if ((creds.user ?? "") !== want.user) fields.push("user");
if (!creds.has_password || (await deps.marks.get(`broker-${id}`)) !== markFor(id, want)) fields.push("password");
if (fields.length === 0) {
outcomes.push(note ? { what: `broker ${id}`, result: "unchanged", note } : { what: `broker ${id}`, result: "unchanged" });
continue;
}
node.broker = want.host;
node.port = String(want.port);
node.usetls = false;
node.credentials = { user: want.user, password: want.password };
changed.push({ id, fields });
}
if (changed.length > 0) {
await deps.nodered.deploy(rev, flows);
for (const c of changed) {
await deps.marks.set(`broker-${c.id}`, markFor(c.id, want));
outcomes.push({ what: `broker ${c.id}`, result: "written", fields: c.fields, ...(note ? { note } : {}) });
}
}
return outcomes;
} catch (err) {
// A deploy against a revision someone else changed meanwhile (409) is read again once.
if (attempt === 0 && /\b409\b/.test(String(err))) continue;
return [...outcomes, { what: "mqtt", result: "refused", problem: scrub(err, want.password) }];
}
}
return outcomes;
}
-102
View File
@@ -1,102 +0,0 @@
// nodered's MQTT step — run once by the host after Node-RED starts, and again whenever the
// `mqtt-topic` binding, its pair credential or the settings change (the container's `restart-on`,
// novox/hq ADR 0099). It points the mesh's broker config nodes at the broker the mesh bound, through
// Node-RED's admin API (connection.ts). It connects to no mesh broker.
//
// Exits non-zero when anything could not be put right, so the node reports the step failed and the
// host runs it again on the next apply. Declared last in the manifest, so its failing gates nothing
// else of nodered's (novox/hq ADR 0136). Never prints a password.
import { mkdir, readFile, rename, writeFile } from "node:fs/promises";
import { join } from "node:path";
import { NodeRedClient } from "../client.js";
import { namedBrokers, reconcileBrokers, type Binding, type Marks } from "./connection.js";
import { probeBroker } from "./probe.js";
const dir = process.env.MESH_PROVISIONS_DIR ?? "/run/provisions";
const writtenDir = process.env.MESH_WRITTEN_DIR ?? "/var/lib/nodered-provisions";
const waitSeconds = Number(process.env.MESH_NODERED_WAIT_SECONDS ?? "180");
const readIfThere = (path: string): Promise<string | undefined> => readFile(path, "utf8").catch(() => undefined);
const parse = <T>(raw: string | undefined): T | undefined => {
if (raw === undefined) return undefined;
try {
return JSON.parse(raw) as T;
} catch {
return undefined;
}
};
const marks: Marks = {
async get(name) {
return (await readIfThere(join(writtenDir, `${name}.digest`)))?.trim() || undefined;
},
async set(name, value) {
await mkdir(writtenDir, { recursive: true, mode: 0o700 });
const path = join(writtenDir, `${name}.digest`);
await writeFile(`${path}.tmp`, `${value}\n`, { mode: 0o600 });
await rename(`${path}.tmp`, path);
},
};
let client: NodeRedClient;
try {
client = NodeRedClient.fromEnv();
} catch (err) {
console.error(`[nodered-mqtt] ${err instanceof Error ? err.message : String(err)}`);
process.exit(1);
}
/** Node-RED answers the admin API once its flows are loaded and the token is good. */
async function ready(): Promise<boolean> {
const until = Date.now() + waitSeconds * 1000;
for (;;) {
try {
await client.flowsWithRev();
return true;
} catch (err) {
if (/\b(401|403)\b/.test(String(err))) {
console.error("[nodered-mqtt] Node-RED refuses the api-token — settings.js and this step disagree");
return false;
}
}
if (Date.now() >= until) return false;
await new Promise((r) => setTimeout(r, 2000));
}
}
if (!(await ready())) {
console.error(`[nodered-mqtt] Node-RED's admin API did not answer at ${client.baseUrl} within ${waitSeconds}s`);
process.exit(1);
}
const binding = parse<Binding>(await readIfThere(join(dir, "mqtt-topic.json")));
const secret = await readIfThere(join(dir, "mqtt-topic.secret"));
const settings = parse<unknown>(await readIfThere(join(dir, "settings.json")));
const outcomes = await reconcileBrokers(
{
nodered: {
flows: () => client.flowsWithRev(),
credentials: (type, id) => client.credentials(type, id),
deploy: (rev, flows) => client.deployFlowsAt(rev, flows, "nodes"),
},
probe: probeBroker,
marks,
},
binding,
secret,
namedBrokers(settings),
);
let failed = 0;
for (const o of outcomes) {
if (o.result === "unchanged") console.log(`[nodered-mqtt] ${o.what}: already as the mesh says${o.note ? ` — ${o.note}` : ""}`);
else if (o.result === "written") console.log(`[nodered-mqtt] ${o.what}: wrote ${o.fields.join(", ")}${o.note ? ` — ${o.note}` : ""}`);
else {
failed++;
console.error(`[nodered-mqtt] ${o.what}: ${o.problem}`);
}
}
process.exitCode = failed > 0 ? 1 : 0;
-117
View File
@@ -1,117 +0,0 @@
// Ask the broker, before Node-RED is told anything, whether it takes the login and password the
// mesh delivered — and whether that login may subscribe to every topic, as flows expect.
//
// One MQTT 3.1.1 session: CONNECT (clean, a throwaway client id, so no flow's session is taken
// over), read the CONNACK, optionally SUBSCRIBE once and read the SUBACK, DISCONNECT. No dependency:
// the handful of bytes MQTT needs for this are written here.
import { randomBytes } from "node:crypto";
import { connect } from "node:net";
export interface ProbeResult {
/** 0 accepted; 4 bad username or password; 5 not authorised. */
connack: number;
/** The SUBACK return code for the filter asked about: 0–2 granted, 0x80 refused. */
suback?: number;
}
export type Probe = (host: string, port: number, username: string, password: string, subscribe?: string) => Promise<ProbeResult>;
function str(v: string): Buffer {
const b = Buffer.from(v, "utf8");
const len = Buffer.alloc(2);
len.writeUInt16BE(b.length);
return Buffer.concat([len, b]);
}
function packet(type: number, body: Buffer): Buffer {
let remaining = body.length;
const lenBytes: number[] = [];
do {
let byte = remaining % 128;
remaining = Math.floor(remaining / 128);
if (remaining > 0) byte |= 0x80;
lenBytes.push(byte);
} while (remaining > 0);
return Buffer.concat([Buffer.from([type, ...lenBytes]), body]);
}
/** The first complete packet in `buf`: its type byte, its body, and how many bytes it took. */
export function firstPacket(buf: Buffer): { type: number; body: Buffer; used: number } | undefined {
if (buf.length < 2) return undefined;
let length = 0;
let multiplier = 1;
let i = 1;
for (;;) {
if (i >= buf.length) return undefined;
const byte = buf[i++];
length += (byte & 0x7f) * multiplier;
if ((byte & 0x80) === 0) break;
multiplier *= 128;
if (i > 4) throw new Error("malformed MQTT remaining length");
}
if (buf.length < i + length) return undefined;
return { type: buf[0], body: buf.subarray(i, i + length), used: i + length };
}
export const probeBroker: Probe = (host, port, username, password, subscribe) => {
const connectBody = Buffer.concat([
str("MQTT"),
Buffer.from([4, 0xc2, 0, 10]), // level 4 (3.1.1); username + password + clean session; keepalive 10s
str(`mesh-probe-${randomBytes(6).toString("hex")}`),
str(username),
str(password),
]);
return new Promise((resolve, reject) => {
const socket = connect({ host, port });
let buf = Buffer.alloc(0);
const result: ProbeResult = { connack: -1 };
const timer = setTimeout(() => {
socket.destroy();
reject(new Error(`no answer from the broker at ${host}:${port} within 10s`));
}, 10_000);
const finish = (): void => {
clearTimeout(timer);
if (result.connack === 0) socket.end(Buffer.from([0xe0, 0]));
else socket.destroy();
resolve(result);
};
socket.on("connect", () => socket.write(packet(0x10, connectBody)));
socket.on("data", (chunk) => {
buf = Buffer.concat([buf, chunk]);
for (;;) {
let p;
try {
p = firstPacket(buf);
} catch (err) {
clearTimeout(timer);
socket.destroy();
reject(err);
return;
}
if (!p) return;
buf = buf.subarray(p.used);
const kind = p.type >> 4;
if (kind === 2) {
result.connack = p.body[1] ?? -1;
if (result.connack !== 0 || !subscribe) return finish();
// SUBSCRIBE, packet id 1, one filter at QoS 0.
socket.write(packet(0x82, Buffer.concat([Buffer.from([0, 1]), str(subscribe), Buffer.from([0])])));
} else if (kind === 9) {
result.suback = p.body[2];
return finish();
}
}
});
socket.on("error", (err) => {
clearTimeout(timer);
reject(err);
});
socket.on("close", () => {
if (result.connack === -1) {
clearTimeout(timer);
reject(new Error(`the broker at ${host}:${port} closed the connection without answering`));
}
});
});
};
+1 -6
View File
@@ -1,14 +1,9 @@
{
"name": "@novox/module-nodered",
"version": "0.1.0",
"description": "nodered \u2014 flow-based automation. Its client and tools live here (novox/hq ADR 0039).",
"description": "nodered — flow-based automation. Its client and tools live here (novox/hq ADR 0039).",
"type": "module",
"private": true,
"scripts": {
"build": "tsc client.ts tools/index.ts mqtt/probe.ts mqtt/connection.ts mqtt/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.0"
},
-124
View File
@@ -1,124 +0,0 @@
// What holds nodered's MQTT step (mqtt/connection.ts): the broker nodes the mesh owns — the one it
// makes, or the ones settings name — are made to use the broker and login the mesh bound, only after
// the broker takes that login; every other field of a node is kept; nothing is deployed when nothing
// differs; a node that is not named is never touched; a loopback broker address is refused.
//
// Node-RED and the broker are fakes answering as the real ones do (admin API v2 of nodered/node-red
// 5.0.7, the build ace runs).
import { test } from "node:test";
import assert from "node:assert/strict";
import { MESH_BROKER_ID, namedBrokers, reconcileBrokers, type Binding, type FlowNode, type Marks, type NodeRed } from "../mqtt/connection.ts";
import type { Probe } from "../mqtt/probe.ts";
const MINTED = "mesh-minted-password";
const binding = (at = "ace.internal"): Binding => ({ provision: "mqtt-topic", from: "ace", at, as: "mesh_ace_nodered", serves: { scheme: "mqtt", port: 1883 } });
/** ace's flows, reduced: its one broker node (dead, zurag.be:1884) and a node that uses it. */
function aceFlows(): FlowNode[] {
return [
{ id: "2b0aece9c5f3b307", type: "mqtt-broker", name: "MQTT Broker", broker: "zurag.be", port: "1884", clientid: "", usetls: false, protocolVersion: "4", keepalive: "60" },
{ id: "fe0cae96f1e3ae4d", type: "mqtt in", topic: "stat/sonoff_office_light_switch/RESULT", broker: "2b0aece9c5f3b307", z: "t" },
{ id: "other-broker", type: "mqtt-broker", name: "someone else's", broker: "test.mosquitto.org", port: "1883" },
];
}
function fakeNodeRed(flows: FlowNode[], creds: Record<string, { user?: string; password?: string }> = {}) {
let rev = "r1";
const deploys: FlowNode[][] = [];
const nodered: NodeRed = {
async flows() {
return { rev, flows: structuredClone(flows) };
},
async credentials(_type, id) {
const c = creds[id] ?? {};
return { user: c.user, has_password: Boolean(c.password) };
},
async deploy(at, next) {
if (at !== rev) throw new Error("Node-RED /flows: 409 version_mismatch");
for (const n of next) {
if (n.credentials) creds[n.id] = { ...(creds[n.id] ?? {}), ...(n.credentials as object) };
}
flows.splice(0, flows.length, ...next.map(({ credentials: _c, ...n }) => n as FlowNode));
deploys.push(next);
rev = `r${deploys.length + 1}`;
},
};
return { nodered, deploys, flows, creds };
}
const marks = (): Marks & { store: Map<string, string> } => {
const store = new Map<string, string>();
return { store, get: async (k) => store.get(k), set: async (k, v) => void store.set(k, v) };
};
const takes: Probe = async (_h, _p, user, pass) => ({ connack: user === "mesh_ace_nodered" && pass === MINTED ? 0 : 5, suback: 0 });
test("ace: the named broker node is moved to the bound broker and login; the other broker is not touched", async () => {
const f = fakeNodeRed(aceFlows(), { "2b0aece9c5f3b307": { user: "luffy", password: "old" }, "other-broker": { user: "x", password: "y" } });
const m = marks();
const out = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: m }, binding(), `${MINTED}\n`, ["2b0aece9c5f3b307"]);
assert.deepEqual(out, [{ what: "broker 2b0aece9c5f3b307", result: "written", fields: ["broker", "port", "user", "password"] }]);
const node = f.flows.find((n) => n.id === "2b0aece9c5f3b307");
assert.deepEqual(node, { ...aceFlows()[0], broker: "ace.internal", port: "1883", usetls: false });
assert.deepEqual(f.creds["2b0aece9c5f3b307"], { user: "mesh_ace_nodered", password: MINTED });
assert.deepEqual(f.flows.find((n) => n.id === "other-broker"), aceFlows()[2]);
assert.deepEqual(f.creds["other-broker"], { user: "x", password: "y" });
// Only the changed node carried credentials in the deploy.
assert.deepEqual(f.deploys[0].filter((n) => n.credentials).map((n) => n.id), ["2b0aece9c5f3b307"]);
// Again: nothing differs, nothing is deployed.
const again = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: m }, binding(), MINTED, ["2b0aece9c5f3b307"]);
assert.deepEqual(again, [{ what: "broker 2b0aece9c5f3b307", result: "unchanged" }]);
assert.equal(f.deploys.length, 1);
});
test("fresh: with nothing named, the step makes its own broker node", async () => {
const f = fakeNodeRed([]);
const out = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: marks() }, binding(), MINTED, []);
assert.deepEqual(out, [{ what: `broker ${MESH_BROKER_ID}`, result: "written", fields: ["node"] }]);
assert.equal(f.flows[0].type, "mqtt-broker");
assert.equal(f.flows[0].broker, "ace.internal");
assert.deepEqual(f.creds[MESH_BROKER_ID], { user: "mesh_ace_nodered", password: MINTED });
});
test("a login the broker does not take is never written", async () => {
const f = fakeNodeRed(aceFlows());
const out = await reconcileBrokers({ nodered: f.nodered, probe: async () => ({ connack: 5 }), marks: marks() }, binding(), MINTED, ["2b0aece9c5f3b307"]);
assert.equal(out[0].result, "refused");
assert.equal(f.deploys.length, 0);
});
test("a named node that is not there, or a loopback broker, is refused", async () => {
const f = fakeNodeRed(aceFlows());
const out = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: marks() }, binding(), MINTED, ["gone"]);
assert.match((out[0] as { problem: string }).problem, /no such node/);
const lo = await reconcileBrokers({ nodered: f.nodered, probe: takes, marks: marks() }, binding("127.0.0.1"), MINTED, []);
assert.match((lo[0] as { problem: string }).problem, /Node-RED itself/);
assert.equal(f.deploys.length, 0);
});
test("a deploy that raced another is read again once", async () => {
const f = fakeNodeRed(aceFlows());
let first = true;
const racing: NodeRed = {
...f.nodered,
async flows() {
const got = await f.nodered.flows();
if (first) {
first = false;
return { ...got, rev: "stale" };
}
return got;
},
};
const out = await reconcileBrokers({ nodered: racing, probe: takes, marks: marks() }, binding(), MINTED, ["2b0aece9c5f3b307"]);
assert.equal(out[0].result, "written");
assert.equal(f.deploys.length, 1);
});
test("settings name broker nodes under mqtt.brokers", () => {
assert.deepEqual(namedBrokers({ mqtt: { brokers: ["a", "", 3, "b"] } }), ["a", "b"]);
assert.deepEqual(namedBrokers({ endpoints: {} }), []);
assert.deepEqual(namedBrokers(undefined), []);
});
+1 -7
View File
@@ -8,11 +8,5 @@
"skipLibCheck": true,
"noEmit": true
},
"include": [
"client.ts",
"tools/index.ts",
"mqtt/probe.ts",
"mqtt/connection.ts",
"mqtt/index.ts"
]
"include": ["client.ts", "tools/index.ts"]
}