Compare commits

..
Author SHA1 Message Date
jschoubben 4991c09c1f Convert the resolver from the module it replaces: forward, answer at 127.0.0.1, point the runtime at it
Read against hal/modules/dnsmasq-app (hal dnsmasq-app conversion, hq 08-connectivity).
On the machines it runs, the predecessor's dnsmasq answers every name: the mesh's own
itself, the rest forwarded to 1.1.1.1 and 8.8.8.8, its module's defaults; resolv.conf
names it alone at 127.0.0.1, and the container runtime's dns is the machine's tunnel
address, so the host and every container resolve the world through it. The nox module
forwarded nothing, listened on 127.0.0.55 — a convention of its own beside the one every
machine already followed — and read a machines file that the module's `facts` already
asks the mesh for, so taking it would have left an adopted machine with a resolv.conf
pointing at an address nothing answered on, and no upstream for anything else.

Now the resolver keeps `no-resolv` (the documented loop — finding its own address in
resolv.conf and becoming its own upstream — stays impossible) and forwards to the same
two explicit upstreams; listens on mesh0 and 127.0.0.1, which systemd-resolved does not
hold; requires `mesh-addressing`, since its data is the mesh's addresses; and writes the
runtime's `dns` into daemon.json beside whatever the machine had (ADR 0102), at this
machine's own address — `${machine:address}`, new in the controller. The runtime is not
restarted for it: it reads the key at start, not on reload, and a restart stops every
container; on the machine this replaces the value is already there.

resolv-conf names the resolver alone, as the predecessor's file did; its placeholder second
line was a fallback nothing ever reached. resolved-split-dns follows the address. The
mesh's suffix as a local domain comes with the machines file, so a mesh name the resolver
does not know is refused here rather than asked upstream. mDNS is not carried: no module
does it and the design says mesh names are not multicast names.
2026-09-23 23:53:51 +02:00
175 changed files with 1970 additions and 3372 deletions
+56
View File
@@ -0,0 +1,56 @@
{
"module": "amqp-email-forwarder",
"version": "1",
"slug": "emailfwd",
"capabilities": [
"container-runtime"
],
"requires": [
"amqp"
],
"contributes": {},
"binds": {
"amqp": "/var/lib/amqp-email-forwarder/amqp.json"
},
"secrets": {
"amqp": "/var/lib/amqp-email-forwarder/amqp.secret"
},
"own-secrets": {
"smtp-user": "/var/lib/amqp-email-forwarder/smtp-user.secret",
"smtp-password": "/var/lib/amqp-email-forwarder/smtp-password.secret"
},
"resources": [
{
"id": "state",
"type": "directory",
"path": "/var/lib/amqp-email-forwarder",
"mode": "0700"
},
{
"id": "app-env",
"type": "file",
"path": "/var/lib/amqp-email-forwarder/app.env",
"mode": "0600",
"content": "AMQP_HOST=${bound:amqp:at}\nAMQP_PORT=${bound:amqp:port}\nAMQP_USER=${bound:amqp:as}\nAMQP_VHOST=EMAILDELIVERY_T\nAMQP_EXCHANGE=News.TransactionalEmailing.Command\nAMQP_QUEUE=email-forwarder\nAMQP_URL=amqp://${bound:amqp:as}:${secret:amqp}@${bound:amqp:at}:${bound:amqp:port}/EMAILDELIVERY_T\nSMTP_HOST=mail.novox.be\nSMTP_PORT=587\nSMTP_USER=${secret:smtp-user}\nSMTP_PASSWORD=${secret:smtp-password}\n"
},
{
"id": "net",
"type": "network",
"name": "amqp-email-forwarder"
},
{
"id": "app",
"type": "container",
"name": "amqp-email-forwarder",
"image": "registry-api.novox.be/novox/amqp-email-forwarder@sha256:f76d34646d9d3b2098c72688a63f6ae656f1888ffcb2f90a9c8dd2a44ad7f8af",
"network": "amqp-email-forwarder",
"env-file": [
"/var/lib/amqp-email-forwarder/app.env"
],
"restart-on": [
"app-env"
],
"secrets-in-environment": "the application's own code reads AMQP_URL, SMTP_USER and SMTP_PASSWORD from the environment (amqp-email-forwarder app.js); converting is that repository's change"
}
]
}
+31
View File
@@ -0,0 +1,31 @@
# amqp-ping's runtime: the tool runtime, carrying this module's compiled code.
#
# **Built from this module's own directory and nothing else.** The sdk and the tool runtime are not
# copied out of neighbouring checkouts — they are in the base image, which is published like any
# other artifact. That is what makes this buildable by the mesh from a repository and a path
# (novox/hq ADR 0069) rather than only on a workstation that happens to have the siblings.
#
# Two bases, named rather than pinned: the image this is COMPILED in, and the image it RUNS in.
# They are different images on purpose — the first carries a compiler and the second must not, or
# every running container would carry one it never invokes. The mesh answers both with the copies it
# holds, because a fingerprint written here would name one particular copy and no other mesh has it
# (novox/hq issue 044). Declared in module.json's `build.on`; deliberately no defaults, so a build
# nobody told stops here and says which module to build first.
ARG BUILD_BASE
ARG RUNTIME_BASE
FROM ${BUILD_BASE} AS build
# Compiled under /app/modules, so resolving `@novox/mesh-sdk` walks up to the base's own
# node_modules — the module is compiled against exactly the sdk it will run against.
WORKDIR /app/modules/amqp-ping
COPY . .
# The compiler is invoked by its real path, not through node_modules/.bin. Those are symlinks to
# a launcher that requires its library relatively, and the base image resolves them when copying —
# leaving a launcher whose relative require no longer points at anything.
RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE}
COPY --from=build /app/modules/amqp-ping/dist /app/modules/amqp-ping/dist
# Declared rather than derived from which files happen to exist: the module knows what it serves.
ENV MESH_TOOL_MODULES=/app/modules/amqp-ping/dist/index.js
+191
View File
@@ -0,0 +1,191 @@
// amqp-ping's AMQP client — the demo consumer's own code (novox/hq ADR 0039). It speaks AMQP 0-9-1
// directly over a raw TCP socket (node:net), the way redis's client speaks RESP: the module carries
// NO npm dependency beyond @novox/mesh-sdk — no amqplib, no CLI in the image. It does exactly one
// thing, the round-trip that proves the grant works: connect, authenticate with PLAIN to the vhost
// the mesh named, declare a queue, publish one message and get it back.
//
// This is the consumer half of the `amqp` interface. It connects as the login the mesh derived
// (`${bound:amqp:as}`) with the password the mesh minted (`${secret:amqp}`) to a vhost of that SAME
// name — the provider named the vhost after the login, so the consumer must too. Nothing here is
// hardcoded: user AND vhost are both the bound login, and a wrong vhost is refused by the broker.
import { createConnection, type Socket } from "node:net";
import { readFileSync } from "node:fs";
const FRAME_END = 0xce;
const PROTOCOL_HEADER = Buffer.from([0x41, 0x4d, 0x51, 0x50, 0x00, 0x00, 0x09, 0x01]); // "AMQP" 0-9-1
export interface AmqpConn {
readonly host: string;
readonly port: number;
readonly user: string;
readonly password: string;
readonly vhost: string;
}
/** Build the connection facts from the environment the mesh's env-file set (see the module manifest). */
export function connFromEnv(env: NodeJS.ProcessEnv = process.env): AmqpConn {
const host = env.MESH_AMQP_HOST ?? "";
const port = Number(env.MESH_AMQP_PORT ?? "5672") || 5672;
const user = env.MESH_AMQP_USER ?? "";
const vhost = env.MESH_AMQP_VHOST ?? user; // the provider names the vhost after the login
const password = env.MESH_AMQP_PASSWORD ?? readMaybe(env.MESH_AMQP_PASSWORD_FILE);
if (!host || !user || !password) {
throw new Error(`amqp-ping: connection is not fully set yet (host=${host} user=${user} password=${password ? "set" : "unset"})`);
}
return { host, port, user, password, vhost };
}
// --- wire helpers ---------------------------------------------------------------------------------
function shortstr(s: string): Buffer {
const b = Buffer.from(s, "utf8");
const o = Buffer.alloc(1 + b.length);
o.writeUInt8(b.length, 0);
b.copy(o, 1);
return o;
}
function longstr(s: Buffer | string): Buffer {
const b = Buffer.isBuffer(s) ? s : Buffer.from(s, "utf8");
const o = Buffer.alloc(4 + b.length);
o.writeUInt32BE(b.length, 0);
b.copy(o, 4);
return o;
}
function u16(n: number): Buffer {
const o = Buffer.alloc(2);
o.writeUInt16BE(n, 0);
return o;
}
function u32(n: number): Buffer {
const o = Buffer.alloc(4);
o.writeUInt32BE(n, 0);
return o;
}
function frame(type: number, channel: number, payload: Buffer): Buffer {
const o = Buffer.alloc(7 + payload.length + 1);
o.writeUInt8(type, 0);
o.writeUInt16BE(channel, 1);
o.writeUInt32BE(payload.length, 3);
payload.copy(o, 7);
o.writeUInt8(FRAME_END, 7 + payload.length);
return o;
}
function method(channel: number, classId: number, methodId: number, ...parts: Buffer[]): Buffer {
return frame(1, channel, Buffer.concat([u16(classId), u16(methodId), ...parts]));
}
interface MethodWaiter {
classId: number;
methodId: number;
resolve: (args: Buffer) => void;
reject: (e: Error) => void;
}
/**
* Connect, authenticate to the vhost, declare a queue, publish one message and get it back. Returns
* the body that came back — the caller checks it equals what went out. Throws on any protocol error,
* including the broker's `NOT_ALLOWED` refusal of a vhost the login has no permission on (the
* isolation the provider builds, seen from the consumer's side).
*/
export function roundTrip(conn: AmqpConn, queue = "amqp-ping", payload?: string): Promise<string> {
const body = Buffer.from(payload ?? `ping-${Date.now()}`);
return new Promise<string>((resolve, reject) => {
const sock: Socket = createConnection({ host: conn.host, port: conn.port });
let buf = Buffer.alloc(0);
const waiters: MethodWaiter[] = [];
let lastBody: Buffer | null = null;
let done = false;
const fail = (e: Error): void => {
if (done) return;
done = true;
sock.destroy();
reject(e);
};
const expect = (classId: number, methodId: number): Promise<Buffer> =>
new Promise((res, rej) => waiters.push({ classId, methodId, resolve: res, reject: rej }));
sock.on("error", (e) => fail(e));
sock.on("close", () => fail(new Error("amqp connection closed before the round-trip completed")));
sock.on("data", (chunk: Buffer) => {
buf = Buffer.concat([buf, chunk]);
for (;;) {
if (buf.length < 7) return;
const type = buf.readUInt8(0);
const size = buf.readUInt32BE(3);
if (buf.length < 7 + size + 1) return;
const framePayload = buf.subarray(7, 7 + size);
buf = buf.subarray(7 + size + 1);
if (type === 1) {
const classId = framePayload.readUInt16BE(0);
const methodId = framePayload.readUInt16BE(2);
const args = framePayload.subarray(4);
const w = waiters.shift();
if (!w) continue;
if (w.classId === classId && w.methodId === methodId) w.resolve(args);
else w.reject(new Error(`expected method ${w.classId}/${w.methodId}, got ${classId}/${methodId}: ${args.toString("utf8")}`));
} else if (type === 3) {
lastBody = framePayload; // a content body frame
}
// type 2 (content header) and type 8 (heartbeat) need no handling for this round-trip.
}
});
sock.on("connect", () => {
void (async () => {
try {
sock.write(PROTOCOL_HEADER);
await expect(10, 10); // Connection.Start
const response = Buffer.concat([
Buffer.from([0]), Buffer.from(conn.user, "utf8"), Buffer.from([0]), Buffer.from(conn.password, "utf8"),
]);
// Connection.Start-Ok: empty client-properties table, PLAIN, the SASL response, locale.
sock.write(method(0, 10, 11, u32(0), shortstr("PLAIN"), longstr(response), shortstr("en_US")));
const tune = await expect(10, 30); // Connection.Tune
const frameMax = tune.readUInt32BE(2) || 131072;
sock.write(method(0, 10, 31, u16(tune.readUInt16BE(0)), u32(frameMax), u16(0))); // Tune-Ok, no heartbeat
sock.write(method(0, 10, 40, shortstr(conn.vhost), shortstr(""), Buffer.from([0]))); // Connection.Open
await expect(10, 41); // Open-Ok — authenticated and into the vhost
sock.write(method(1, 20, 10, shortstr(""))); // Channel.Open
await expect(20, 11);
// Queue.Declare: reserved, queue, bits(auto-delete=1), empty arguments table.
sock.write(method(1, 50, 10, u16(0), shortstr(queue), Buffer.from([0b00001000]), u32(0)));
await expect(50, 11);
// Basic.Publish to the default exchange, routing-key = queue; then content header + body.
sock.write(method(1, 60, 40, u16(0), shortstr(""), shortstr(queue), Buffer.from([0])));
const bodySize = Buffer.alloc(8);
bodySize.writeBigUInt64BE(BigInt(body.length), 0);
sock.write(frame(2, 1, Buffer.concat([u16(60), u16(0), bodySize, u16(0)]))); // content header, no properties
sock.write(frame(3, 1, body)); // content body
await new Promise((r) => setTimeout(r, 200));
lastBody = null;
sock.write(method(1, 60, 70, u16(0), shortstr(queue), Buffer.from([1]))); // Basic.Get, no-ack
await expect(60, 71); // Get-Ok (a Get-Empty would arrive as 60/72 and reject the expect)
await new Promise((r) => setTimeout(r, 200));
const received = lastBody ? (lastBody as Buffer).toString("utf8") : "";
sock.write(method(0, 10, 50, u16(200), shortstr("bye"), u16(0), u16(0))); // Connection.Close
await expect(10, 51).catch(() => undefined);
done = true;
sock.end();
resolve(received);
} catch (e) {
fail(e instanceof Error ? e : new Error(String(e)));
}
})();
});
});
}
function readMaybe(path: string | undefined): string {
if (!path) return "";
try {
return readFileSync(path, "utf8").replace(/\n$/, "");
} catch {
return "";
}
}
+53
View File
@@ -0,0 +1,53 @@
// amqp-ping — a tiny demo consumer of the mesh `amqp` interface, run as a long-lived container by
// `mesh-tools run` (it never returns, so the container stays up). It exists to PROVE the grant end to
// end: the mesh gave it a scoped login and a vhost of that name on the lavinmq provider, and this
// connects with exactly those and round-trips a message.
//
// The connection facts arrive the way every consumer's do — the mesh writes them into an env-file the
// container reads (novox/hq ADR 0048): MESH_AMQP_HOST/PORT from the binding, MESH_AMQP_USER and
// MESH_AMQP_VHOST both from `${bound:amqp:as}` (the provider named the vhost after the login, so the
// consumer uses the login for both — the db-name lesson applied to AMQP), and MESH_AMQP_PASSWORD from
// `${secret:amqp}`.
//
// It retries: on first boot the provider may not have provisioned this consumer yet (the reconcile is
// asynchronous and cross-container), so a refused or unreachable connection is a "not yet", not a
// failure — it waits and tries again until the round-trip succeeds, then holds the connection idle
// and re-pings on a slow cadence so the container is a stable, running proof.
import { connFromEnv, roundTrip } from "./client.js";
async function sleep(ms: number): Promise<void> {
await new Promise((r) => setTimeout(r, ms));
}
async function pingOnce(): Promise<boolean> {
try {
const conn = connFromEnv();
const sent = `ping-${Date.now()}`;
const got = await roundTrip(conn, "amqp-ping", sent);
if (got === sent) {
console.log(`[amqp-ping] round-trip ok as ${conn.user} on vhost ${conn.vhost} (${conn.host}:${conn.port})`);
return true;
}
console.error(`[amqp-ping] round-trip mismatch: sent ${sent}, got ${got}`);
return false;
} catch (err) {
console.error(`[amqp-ping] not ready yet: ${err instanceof Error ? err.message : err}`);
return false;
}
}
// Wait for the first successful round-trip — the proof this consumer's grant works — then stay up.
let first = false;
for (let i = 0; !first; i++) {
first = await pingOnce();
if (!first) await sleep(3000);
}
console.log("[amqp-ping] connected and round-tripped; holding steady");
for (;;) {
await sleep(30000);
await pingOnce();
}
// changed by the one-node test at build 66e54af151df
// changed by the one-node test at build 4fb41636cffd
// changed by the one-node test at build a69f083bf6a6
+89
View File
@@ -0,0 +1,89 @@
{
"module": "amqp-ping",
"slug": "ping",
"version": "1",
"capabilities": [
"container-runtime"
],
"requires": [
"amqp"
],
"contributes": {},
"binds": {
"amqp": "/var/lib/amqp-ping/amqp.json"
},
"secrets": {
"amqp": "/var/lib/amqp-ping/amqp.secret"
},
"own-secrets": {
"broker": "/var/lib/mesh/amqp-ping/broker"
},
"resources": [
{
"id": "mesh-state",
"type": "directory",
"path": "/var/lib/mesh/amqp-ping",
"mode": "0700"
},
{
"id": "state",
"type": "directory",
"path": "/var/lib/amqp-ping",
"mode": "0700"
},
{
"id": "amqp-env",
"type": "file",
"path": "/var/lib/amqp-ping/amqp.env",
"mode": "0600",
"content": "MESH_AMQP_HOST=${bound:amqp:at}\nMESH_AMQP_PORT=${bound:amqp:port}\nMESH_AMQP_USER=${bound:amqp:as}\nMESH_AMQP_VHOST=${bound:amqp:as}\n"
},
{
"id": "net",
"type": "network",
"name": "amqp-ping"
},
{
"id": "runtime",
"type": "container",
"name": "amqp-ping",
"network": "amqp-ping",
"volumes": [
"/var/lib/mesh/amqp-ping/broker:/run/secrets/broker:ro",
"/var/lib/amqp-ping/amqp.secret:/run/secrets/amqp:ro"
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_AMQP_PASSWORD_FILE": "/run/secrets/amqp"
},
"env-file": [
"/var/lib/amqp-ping/amqp.env"
],
"restart-on": [
"amqp-env"
],
"artifact": "runtime"
}
],
"build": {
"on": [
{
"arg": "BUILD_BASE",
"module": "mesh-tools",
"artifact": "build"
},
{
"arg": "RUNTIME_BASE",
"module": "mesh-tools",
"artifact": "runtime"
}
],
"artifacts": [
{
"name": "runtime",
"kind": "image",
"from": "Dockerfile"
}
]
}
}
+14
View File
@@ -0,0 +1,14 @@
{
"name": "@novox/module-amqp-ping",
"version": "0.1.0",
"description": "amqp-ping — a demo consumer of the mesh amqp interface. Connects with its scoped grant and round-trips one message to prove the broker the mesh gave it (novox/hq ADR 0039).",
"type": "module",
"private": true,
"dependencies": {
"@novox/mesh-sdk": "^0.1.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
"typescript": "^5.6.0"
}
}
+12
View File
@@ -0,0 +1,12 @@
{
"compilerOptions": {
"target": "ES2022",
"module": "NodeNext",
"moduleResolution": "NodeNext",
"strict": true,
"esModuleInterop": true,
"skipLibCheck": true,
"noEmit": true
},
"include": ["client.ts", "index.ts"]
}
+1 -1
View File
@@ -18,7 +18,7 @@
"broker": "/var/lib/mesh/anthropic-consumer/broker"
},
"emits": [
"usage.session"
"module.anthropic-consumer.usage.session"
],
"resources": [
{
+1 -1
View File
@@ -123,7 +123,7 @@ async function emitUsage(body: Record<string, unknown>): Promise<void> {
await new Promise<void>((resolve) => {
const child = spawn(
process.execPath,
[main, "emit", "usage.session", JSON.stringify(body)],
[main, "emit", "module.anthropic-consumer.usage.session", JSON.stringify(body)],
{ stdio: "inherit" },
);
child.on("exit", () => resolve());
+1 -1
View File
@@ -18,7 +18,7 @@
"broker": "/var/lib/mesh/anthropic-manager/broker"
},
"emits": [
"usage.read"
"module.anthropic-manager.usage.read"
],
"resources": [
{
+1 -1
View File
@@ -143,7 +143,7 @@ async function emitUsage(body: Record<string, unknown>): Promise<void> {
const main = process.env.MESH_TOOLS_MAIN ?? "/app/dist/main.js";
const { spawn } = await import("node:child_process");
await new Promise<void>((resolve) => {
const child = spawn(process.execPath, [main, "emit", "usage.read", JSON.stringify(body)], {
const child = spawn(process.execPath, [main, "emit", "module.anthropic-manager.usage.read", JSON.stringify(body)], {
stdio: "inherit",
});
child.on("exit", () => resolve());
+1 -1
View File
@@ -3,7 +3,7 @@
"version": "1",
"slug": "audit",
"consumes": [
"**"
"#"
],
"own-secrets": {
"broker": "/var/lib/audit-logger/broker"
+3 -3
View File
@@ -15,16 +15,16 @@ test("audit-logger records every event to the trail as one line each", async ()
const path = join(dir, "audit.log");
// The audit-logger's whole behaviour: consume everything, record it.
await on("**", async (event) => record(event, path));
await on("#", async (event) => record(event, path));
process.env.MESH_MODULE = "umami";
process.env.MESH_NODE = "anchor";
await emit("site.created", { domain: "my-app" });
await emit("module.umami.site.created", { domain: "my-app" });
await emit("node.anchor.joined", { role: "worker" }); // a node event, not a module one
const lines = (await readFile(path, "utf8")).trim().split("\n").map((l) => JSON.parse(l));
assert.equal(lines.length, 2);
assert.deepEqual(lines.map((l) => l.type), ["umami.site.created", "node.anchor.joined"]);
assert.deepEqual(lines.map((l) => l.type), ["module.umami.site.created", "node.anchor.joined"]);
assert.equal(lines[0].source, "umami");
assert.equal(lines[0].node, "anchor");
assert.equal(lines[0].body.domain, "my-app");
+1 -2
View File
@@ -14,7 +14,7 @@
},
"route": {
"label": "baserow",
"endpoint": "web"
"port": 80
}
},
"binds": {
@@ -30,7 +30,6 @@
},
"listens": [
{
"name": "web",
"port": 80,
"protocol": "tcp",
"from": "mesh",
+1 -1
View File
@@ -24,7 +24,7 @@ async function pollHistory(): Promise<void> {
for (const entry of entries) {
if (seen.has(entry.id)) continue;
if (primed) {
await emit("subtitle.downloaded", {
await emit("module.bazarr.subtitle.downloaded", {
kind: entry.kind,
title: entry.title,
language: entry.language,
+3 -4
View File
@@ -5,7 +5,7 @@
"container-runtime"
],
"emits": [
"subtitle.downloaded"
"module.bazarr.subtitle.downloaded"
],
"own-secrets": {
"broker": "/var/lib/mesh/bazarr/broker",
@@ -13,7 +13,6 @@
},
"listens": [
{
"name": "web",
"port": 6767,
"protocol": "tcp",
"from": "mesh",
@@ -94,7 +93,7 @@
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_BAZARR_URL": "http://127.0.0.1:${port:6767}",
"MESH_BAZARR_URL": "http://127.0.0.1:6767",
"MESH_BAZARR_API_KEY_FILE": "/run/secrets/api-key",
"MESH_BAZARR_CONFIG_FILE": "/run/config/config.json",
"MESH_BAZARR_CONFIG_DIR": "/var/lib/bazarr/config"
@@ -111,7 +110,7 @@
"contributes": {
"route": {
"label": "subs",
"endpoint": "web"
"port": 6767
}
},
"binds": {
+2 -2
View File
@@ -45,12 +45,12 @@ async function pollQueue(bookshelf: BookshelfClient): Promise<void> {
if (primed) {
// Entered the queue since last look — Bookshelf grabbed a release.
for (const [id, item] of now) {
if (!inQueue.has(id)) await emit("book.grabbed", { title: item.title, status: item.status });
if (!inQueue.has(id)) await emit("module.bookshelf.book.grabbed", { title: item.title, status: item.status });
}
// Left the queue — imported and done, unless it was last seen failing.
for (const [id, item] of inQueue) {
if (!now.has(id) && !FAILED_STATUSES.has(item.status)) {
await emit("download.completed", { title: item.title });
await emit("module.bookshelf.download.completed", { title: item.title });
}
}
}
+4 -5
View File
@@ -6,8 +6,8 @@
"container-runtime"
],
"emits": [
"book.grabbed",
"download.completed"
"module.bookshelf.book.grabbed",
"module.bookshelf.download.completed"
],
"consumes": [],
"own-secrets": {
@@ -15,7 +15,6 @@
},
"listens": [
{
"name": "web",
"port": 8787,
"protocol": "tcp",
"from": "mesh",
@@ -76,7 +75,7 @@
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_BOOKSHELF_URL": "http://127.0.0.1:${port:8787}",
"MESH_BOOKSHELF_URL": "http://127.0.0.1:8787",
"MESH_BOOKSHELF_CONFIG_DIR": "/var/lib/bookshelf/config"
},
"artifact": "runtime"
@@ -88,7 +87,7 @@
"contributes": {
"route": {
"label": "books",
"endpoint": "web"
"port": 8787
}
},
"binds": {
-25
View File
@@ -1,25 +0,0 @@
ARG GO_BASE
ARG ALPINE_BASE
# builder's own image: the build machine itself, compiled into a container.
#
# **The source is not vendored here.** builder's actual code — cmd/mesh-builder, internal/builder,
# internal/catalogue — lives in the mesh-controller repository, the same control plane it is one
# half of. This module ships the packaging, not a second copy of the source, so the build context
# is the mesh-controller repository root (declared under build.artifacts[].context), and this
# Dockerfile compiles ./cmd/mesh-builder from it — the same shape route-proxy already uses for the
# same reason.
FROM ${GO_BASE} AS build
WORKDIR /src
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux go build -trimpath -o /mesh-builder ./cmd/mesh-builder
# Unlike mesh-controller's own FROM scratch (ADR 0006: nothing to audit but one binary), the build
# machine's whole job is shelling out to git and docker — it needs a real userland to do that in,
# not a second copy of either tool vendored into this image. apk installs both from the base's own
# packages, not fetched on its own at build time.
FROM ${ALPINE_BASE}
RUN apk add --no-cache docker-cli git
COPY --from=build /mesh-builder /usr/local/bin/mesh-builder
ENTRYPOINT ["/usr/local/bin/mesh-builder"]
+28 -38
View File
@@ -6,22 +6,19 @@
],
"claims": [
{
"name": "mesh-build-machine",
"scope": "mesh"
"name": "the-build-machine",
"scope": "node"
}
],
"requires": [
"artifact-store",
"npm-package-registry"
"artifact-store"
],
"emits": [
"module.builder.built"
],
"binds": {
"npm-package-registry": "/var/lib/mesh/builder/package-registry.json"
},
"secrets": {
"npm-package-registry": "/var/lib/mesh/builder/package-registry.secret"
},
"own-secrets": {
"broker": "/var/lib/mesh/builder/broker"
"broker": "/var/lib/mesh/builder/broker",
"npm-password": "/var/lib/mesh/builder/npm-password"
},
"resources": [
{
@@ -41,13 +38,27 @@
"type": "file",
"path": "/var/lib/mesh/builder/builder.env",
"mode": "0600",
"content": "MESH_BROKER_FILE=/run/mesh/broker\nMESH_NODE=${machine:name}\nMESH_REGISTRY=${bound:artifact-store:at}:${bound:artifact-store:port}\nMESH_PACKAGE_BINDING=/run/mesh/package-registry.json\nMESH_NPM_TOKEN_FILE=/run/mesh/package-registry.secret\nMESH_WORKSPACE=/var/lib/builder/workspace\n"
"content": "MESH_BROKER_FILE=/run/mesh/broker\nMESH_NODE=${machine:name}\nMESH_REGISTRY=${bound:artifact-store:at}:${bound:artifact-store:port}\nMESH_PACKAGE_BINDING=/run/mesh/package-registry.json\nMESH_NPM_TOKEN_FILE=/run/mesh/npm-password\nMESH_WORKSPACE=/var/lib/builder/workspace\n"
},
{
"id": "package-binding",
"type": "file",
"path": "/var/lib/mesh/builder/package-registry.json",
"mode": "0600",
"merge": "json",
"protected": [
"provision",
"from",
"at",
"as"
],
"content": "{\"provision\": \"package-registry\", \"from\": \"gitea\", \"at\": \"127.0.0.1\", \"as\": \"mesh-builder\", \"serves\": {\"scheme\": \"http\", \"port\": 3000, \"npm-path\": \"/api/packages/novox/npm/\"}}\n"
},
{
"id": "server",
"type": "container",
"name": "mesh-builder",
"artifact": "server",
"image": "mesh-builder@sha256:0000000000000000000000000000000000000000000000000000000000000000",
"env-file": [
"/var/lib/mesh/builder/builder.env"
],
@@ -57,32 +68,11 @@
"/var/run/docker.sock:/var/run/docker.sock"
],
"restart-on": [
"builder-env"
"builder-env",
"package-binding",
"needs-npm-password"
],
"network": "host"
}
],
"build": {
"artifacts": [
{
"name": "server",
"kind": "image",
"from": "Dockerfile",
"context": {
"repository": "https://git.novox.be/novox/mesh-controller.git",
"ref": "main"
}
}
],
"on": [
{
"arg": "GO_BASE",
"image": "golang@sha256:8ac98ca534ac3f51e1f420a1dd2c15e74c75cfa0f23f3ad27eb5d7236c349a0c"
},
{
"arg": "ALPINE_BASE",
"image": "alpine@sha256:d9e853e87e55526f6b2917df91a2115c36dd7c696a35be12163d44e6e2a4b6bc"
}
]
}
]
}
-56
View File
@@ -1,56 +0,0 @@
{
"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"
]
}
]
}
+2 -2
View File
@@ -22,8 +22,8 @@
"broker": "/var/lib/mesh/cloudflare-dns/broker"
},
"emits": [
"record.created",
"record.removed"
"module.cloudflare-dns.record.created",
"module.cloudflare-dns.record.removed"
],
"resources": [
{
+2 -2
View File
@@ -20,7 +20,7 @@ runProvisioner("public-dns", {
async create(p: Provision): Promise<void> {
const fqdn = cloudflare.nameFor(p.as);
await cloudflare.upsert(fqdn);
await announce("record.created", {
await announce("module.cloudflare-dns.record.created", {
name: fqdn,
target: cloudflare.ingress,
consumer: p.consumer ?? "",
@@ -30,7 +30,7 @@ runProvisioner("public-dns", {
async remove(p: { as: string }): Promise<void> {
const fqdn = cloudflare.nameFor(p.as);
await cloudflare.remove(fqdn);
await announce("record.removed", { name: fqdn, consumer: p.as });
await announce("module.cloudflare-dns.record.removed", { name: fqdn, consumer: p.as });
},
});
+2 -3
View File
@@ -11,7 +11,7 @@
"contributes": {
"route": {
"label": "de-spiegel",
"endpoint": "web"
"port": 35621
}
},
"binds": {
@@ -23,7 +23,6 @@
},
"listens": [
{
"name": "web",
"port": 35621,
"protocol": "tcp",
"from": "mesh",
@@ -59,7 +58,7 @@
"/var/lib/de-spiegel/server.env"
],
"ports": [
"35621"
"35621:35621"
],
"secrets-in-environment": "the application's own code reads SMTP_AUTH_USER/PASS from the environment (de-spiegel server/index.js); converting is that repository's change"
}
-43
View File
@@ -1,43 +0,0 @@
# dhcpcd
The uplink seat's module for a machine whose own network is dhcpcd's (novox/hq ADR 0117). It
asks two things of dhcpcd, and nothing else: leave the resolver file to the mesh, and leave the
private network's interface alone. It never declares an interface, an address, a route, a
wireless network or its credentials — the link dhcpcd keeps is the only channel the mesh reaches
the machine over.
## What it writes
Two lines into `/etc/dhcpcd.conf`, as the mesh's marked region (`into: block`) — dhcpcd reads no
drop-in directory, so the mesh writes into its one file rather than over it (ADR 0102):
- `nohook resolv.conf` — dhcpcd's resolv.conf hook rewrites `/etc/resolv.conf` on every lease it
takes or renews, which would silently replace the resolver `resolv-conf` names.
- `denyinterfaces mesh0` — dhcpcd never asks for a lease on the private network's interface, and
never takes it down. dhcpcd leaves a point-to-point interface alone by default; this says so
rather than relying on it.
**At the start of the file** (`at: start`). Both are global options, and dhcpcd reads every line
after an `interface` or `ssid` line as that interface's own. A configured machine's file ends in
exactly such a block (the interface, its static address), so appended at the end these two would
quietly apply to one interface only.
## Why it declares no service
dhcpcd is the machine's, not the mesh's. The mesh never starts, stops or enables it: stopping it
drops the address the machine is reached at, and a module unassigned by mistake must not be able
to do that. And there is nothing to reload it with — `dhcpcd.service` reports `CanReload=no`, and
a restart drops the lease. So the two lines take effect at **dhcpcd's next start**.
On an adopted machine that is normally no gap: the predecessor wrote the same `nohook` line, and
it is already in force. **On a machine that was not adopted, it is one:** until dhcpcd next
starts (a reboot, or the operator restarting it in a window of their choosing), a lease renewal
still rewrites `/etc/resolv.conf`, and `resolv-conf` puts it back at the next push. Assign this
module before `resolv-conf` on such a machine, and restart dhcpcd once, by hand, when losing the
link for a moment is acceptable.
## One manager per machine
It claims `the-uplink`: a machine runs one network manager, and assigning a second module that
claims the seat is refused. Assigning this one to a machine whose network is NetworkManager's
installs the package and writes the two lines, and starts nothing.
-30
View File
@@ -1,30 +0,0 @@
{
"module": "dhcpcd",
"version": "1",
"capabilities": [
"package-manager",
"service-manager"
],
"claims": [
{
"name": "node-uplink",
"scope": "node"
}
],
"resources": [
{
"id": "package",
"type": "package",
"package": "dhcpcd"
},
{
"id": "config",
"type": "file",
"path": "/etc/dhcpcd.conf",
"mode": "0644",
"into": "block",
"at": "start",
"content": "# The mesh's two lines (module dhcpcd, novox/hq ADR 0117). Global options, so\n# kept above any interface line; read at dhcpcd's next start.\nnohook resolv.conf\ndenyinterfaces mesh0\n"
}
]
}
+1 -1
View File
@@ -33,7 +33,7 @@ async function pollCatalog(): Promise<void> {
for (const tag of tags) {
const id = `${repo}:${tag}`;
if (!seen.has(id)) {
if (primed) await emit("image.pushed", { repo, tag });
if (primed) await emit("module.registry.image.pushed", { repo, tag });
seen.add(id);
}
}
+3 -10
View File
@@ -10,14 +10,14 @@
"claims": [
{
"name": "the-artifact-store",
"scope": "mesh"
"scope": "node"
}
],
"capabilities": [
"container-runtime"
],
"emits": [
"image.pushed"
"module.registry.image.pushed"
],
"own-secrets": {
"broker": "/var/lib/mesh/registry/broker"
@@ -29,7 +29,6 @@
},
"listens": [
{
"name": "registry",
"port": 5000,
"protocol": "tcp",
"from": "mesh",
@@ -43,12 +42,6 @@
"path": "/var/lib/mesh/registry",
"mode": "0700"
},
{
"id": "registry-data",
"type": "directory",
"path": "/var/lib/mesh-registry",
"mode": "0700"
},
{
"id": "store",
"type": "container",
@@ -58,7 +51,7 @@
"5000:5000"
],
"volumes": [
"/var/lib/mesh-registry:/var/lib/registry"
"mesh-registry-data:/var/lib/registry"
]
}
]
+2 -2
View File
@@ -20,10 +20,10 @@ async function poll(): Promise<void> {
const now = new Map((await dnsmasq.answeredNames()).map((a) => [a.name, a.address]));
if (primed) {
for (const [name, address] of now) {
if (!known.has(name)) await emit("name.added", { name, address });
if (!known.has(name)) await emit("module.dnsmasq.name.added", { name, address });
}
for (const [name] of known) {
if (!now.has(name)) await emit("name.removed", { name });
if (!now.has(name)) await emit("module.dnsmasq.name.removed", { name });
}
}
known.clear();
File diff suppressed because one or more lines are too long
+2 -10
View File
@@ -6,7 +6,7 @@
],
"claims": [
{
"name": "node-intrusion-prevention",
"name": "the-intrusion-prevention",
"scope": "node"
}
],
@@ -33,7 +33,7 @@
"type": "file",
"path": "/etc/fail2ban/jail.local",
"mode": "0644",
"content": "[INCLUDES]\n\nbefore = paths-arch.conf\n\n[DEFAULT]\n\n# Never act on the machine itself or on a tunnel peer: the mesh's private range is\n# ${machine:mesh-range}, named here rather than written as a value the module cannot\n# know (novox/hq ADR 0112). Without this, fail2ban could ban the mesh's own nodes.\nignoreip = 127.0.0.1/8 ::1 ${machine:mesh-range}\n\nbantime = 10m\nfindtime = 10m\nmaxretry = 5\n\n# Ban through iptables, not through a firewall front-end the machine may not have. ufw is\n# installed on two of this mesh's machines and absent on the other two, and fail2ban finds out\n# only at ban time: the service reports healthy, the jail counts the attempt, the ban command\n# exits 127, and nothing is blocked. Proven on 2026-09-28 -- 'ufw: command not found' on a\n# machine the mesh reported as protected.\n#\n# The action below is this module's own, already used by the recidive jail on every machine\n# here, and it bans in DOCKER-USER as well as INPUT, so a container's published port is\n# covered too.\nbanaction = iptables-allports-dualchain\nbanaction_allports = iptables-allports-dualchain\n\n[sshd]\nenabled = true\nport = ssh\nlogpath = %(sshd_log)s\nbackend = %(sshd_backend)s\n"
"content": "[INCLUDES]\n\nbefore = paths-arch.conf\n\n[DEFAULT]\n\nbantime = 10m\nfindtime = 10m\nmaxretry = 5\n\nbanaction = ufw\nbanaction_allports = iptables-allports\n\n[sshd]\nenabled = true\nport = ssh\nlogpath = %(sshd_log)s\nbackend = %(sshd_backend)s\n"
},
{
"id": "jail-sshd",
@@ -42,14 +42,6 @@
"mode": "0644",
"content": "[sshd]\nenabled = true\nport = ssh\nlogpath = %(sshd_log)s\nbackend = %(sshd_backend)s\nmaxretry = 5\n"
},
{
"id": "log",
"type": "file",
"path": "/var/log/fail2ban.log",
"mode": "0640",
"create-once": true,
"content": ""
},
{
"id": "jail-recidive",
"type": "file",
+1 -1
View File
@@ -23,7 +23,7 @@ COPY . .
# The compiler is invoked by its real path rather than through node_modules/.bin, whose entries are
# symlinks to a launcher that requires its library relatively — resolved away when the base image
# was assembled.
RUN node /app/node_modules/typescript/bin/tsc client.ts token.ts index.ts provisioner/index.ts tools/index.ts \
RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts provisioner/index.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE}
+45 -137
View File
@@ -4,13 +4,10 @@
// does.
import { readFileSync } from "node:fs";
import { ConfiguredToken, MintedToken, type TokenSource } from "./token.js";
/** A repository, trimmed to what the mesh cares about. */
export interface GiteaRepo {
full_name: string;
/** The URL a build clones — what a module records as its source. */
clone_url?: string;
name: string;
owner: string;
private: boolean;
@@ -36,9 +33,6 @@ export interface GiteaPull {
title: string;
state: string;
merged: boolean;
/** The commit the merge produced — what a build of the base branch is made from. */
merge_commit_sha?: string;
merged_at?: string;
user?: string;
head?: string;
base?: string;
@@ -59,60 +53,44 @@ function meshConfig(file?: string): Record<string, string> {
export class GiteaClient {
readonly baseUrl: string;
private readonly tokens: TokenSource;
private cachedUsername: string | null = null;
/** A token given as a string is one somebody configured; a source decides for itself (token.ts). */
constructor(url: string, token: string | TokenSource) {
constructor(
url: string,
private readonly token: string,
) {
this.baseUrl = url.replace(/\/+$/, "");
this.tokens = typeof token === "string" ? new ConfiguredToken(token) : token;
}
/**
* Build from the module's resolved environment. The URL comes from MESH_GITEA_URL (the mesh's own
* name), falling back to the bare GITEA_URL and to the forge's loopback port. The token, in order:
* one configured in settings or the environment (MESH_GITEA_TOKEN / GITEA_TOKEN), which wins; else
* one the module mints for itself with the admin account the vault delivered and keeps in its own
* state (token.ts; hq issue 100). Throws only when neither is possible, naming what is missing,
* rather than hand back a client that fails on first use.
* Build from the module's resolved environment. URL and token come from MESH_GITEA_URL /
* MESH_GITEA_TOKEN (the mesh's own names), falling back to the bare GITEA_* names and, for the
* URL, to the forge's loopback port. A token is required — without one there is no authenticated
* call to make, so this throws rather than hand back a client that fails on first use.
*/
static fromEnv(env: NodeJS.ProcessEnv = process.env): GiteaClient {
const cfg = meshConfig(env.MESH_GITEA_CONFIG_FILE);
const url = cfg.url ?? env.MESH_GITEA_URL ?? env.GITEA_URL ?? `http://127.0.0.1:${env.GITEA_PORT ?? "3000"}`;
const configured = cfg.token ?? env.MESH_GITEA_TOKEN ?? env.GITEA_TOKEN;
if (configured) return new GiteaClient(url, new ConfiguredToken(configured));
return new GiteaClient(url, MintedToken.fromEnv(url, env));
const token = cfg.token ?? env.MESH_GITEA_TOKEN ?? env.GITEA_TOKEN;
if (!token) throw new Error("no Gitea token — set MESH_GITEA_TOKEN");
return new GiteaClient(url, token);
}
/**
* One authenticated call. A 401 is the forge saying the token is not one it knows — the case
* after the forge's data was restored, or after somebody revoked it — so the source is asked to
* renew once and the call is repeated with the new token. A configured token has nothing to renew
* with, and its source says so.
*/
private async request<T = unknown>(path: string, options: RequestInit = {}): Promise<T> {
let token = await this.tokens.current();
let res = await this.send(path, options, token);
if (res.status === 401) {
token = await this.tokens.renew(token);
res = await this.send(path, options, token);
}
const res = await fetch(`${this.baseUrl}/api/v1${path}`, {
...options,
headers: {
"Content-Type": "application/json",
Authorization: `token ${this.token}`,
...(options.headers as Record<string, string> | undefined),
},
});
if (!res.ok) throw new Error(`Gitea API ${path}: ${res.status} ${await res.text()}`);
if (res.status === 204) return null as T;
const text = await res.text();
return (text ? JSON.parse(text) : null) as T;
}
private send(path: string, options: RequestInit, token: string): Promise<Response> {
return fetch(`${this.baseUrl}/api/v1${path}`, {
...options,
headers: {
"Content-Type": "application/json",
Authorization: `token ${token}`,
...(options.headers as Record<string, string> | undefined),
},
});
}
/** Generic authenticated API call — the escape hatch for endpoints without a dedicated method.
* Path is relative to /api/v1. */
async api<T = unknown>(path: string, options: RequestInit = {}): Promise<T> {
@@ -121,23 +99,9 @@ export class GiteaClient {
// ---- Repositories ----
/** Every repository this token can see, one page. `/user/repos` is only what the token's own
* user owns — for the mesh's administrator that is nothing, which is how the forge watched an
* empty list and announced no merge (2026-09-28). The search endpoint is the forge's whole view. */
async listRepos(page = 1, limit = 20): Promise<GiteaRepo[]> {
const found = await this.request<{ data?: any[] }>(`/repos/search?page=${page}&limit=${limit}`);
return (found?.data ?? []).map(GiteaClient.mapRepo);
}
/** Every repository, all pages. */
async listAllRepos(): Promise<GiteaRepo[]> {
const all: GiteaRepo[] = [];
for (let page = 1; page < 100; page++) {
const batch = await this.listRepos(page, 50);
all.push(...batch);
if (batch.length < 50) break;
}
return all;
const repos = await this.request<any[]>(`/user/repos?page=${page}&limit=${limit}`);
return (repos ?? []).map(GiteaClient.mapRepo);
}
async createRepo(data: {
@@ -229,17 +193,6 @@ export class GiteaClient {
return GiteaClient.mapPull(await this.request<any>(`/repos/${owner}/${repo}/pulls/${index}`));
}
/** The files a merged pull request changed, as paths from the repository's root.
*
* `limit` is what is asked for, and a merge that changed more says so rather than being read
* page by page: what the mesh does with a partial list is treat the whole repository as changed,
* so more pages would buy nothing. */
async listPullFiles(owner: string, repo: string, index: number, limit = 100): Promise<{ paths: string[]; truncated: boolean }> {
const files = await this.request<any[]>(`/repos/${owner}/${repo}/pulls/${index}/files?limit=${limit}`);
const paths = (files ?? []).map((f) => String(f?.filename ?? "")).filter((p) => p !== "");
return { paths, truncated: paths.length >= limit };
}
async createPullRequest(
owner: string,
repo: string,
@@ -262,7 +215,6 @@ export class GiteaClient {
private static mapRepo(r: any): GiteaRepo {
return {
full_name: r.full_name,
clone_url: r.clone_url ?? undefined,
name: r.name,
owner: r.owner?.login ?? r.full_name?.split("/")[0] ?? "unknown",
private: Boolean(r.private),
@@ -290,8 +242,6 @@ export class GiteaClient {
title: p.title,
state: p.state,
merged: Boolean(p.merged),
merge_commit_sha: p.merge_commit_sha ?? undefined,
merged_at: p.merged_at ?? undefined,
user: p.user?.login,
head: p.head?.ref,
base: p.base?.ref,
@@ -385,34 +335,24 @@ export class GiteaAdmin {
GiteaAdmin.fail("/orgs", res);
}
/** Ensure the org's package team exists with exactly these units, and return its id. Found or
* created, the units are applied either way — a team is configuration the reconcile loop owns,
* the same as a user's password, so a unit this code gains reaches a team that already exists
* rather than only the next mesh raised from scratch. A lost create race is resolved by
* re-listing. */
/** Ensure the org's package team exists, granting read+write on packages, and return its id. The
* team is found by name if it is already there, created otherwise; a lost create race is resolved
* by re-listing. */
async ensureTeam(org: string, team: string, packageWrite: boolean): Promise<number> {
// The units a consumer needs, and no more. `units_map` is exhaustive — a unit not named is a
// unit the team does not have — so code read must be said here: without it gitea answers a
// member's clone of a private repository with "not found", which is how the builder's first
// credentialed clone failed against a team that named only packages.
const units = {
permission: "read",
units_map: { "repo.code": "read", "repo.packages": packageWrite ? "write" : "read" },
includes_all_repositories: true,
can_create_org_repo: false,
};
const found = await this.findTeam(org, team);
if (found !== null) {
const patch = await this.request(`/teams/${found}`, {
method: "PATCH",
body: JSON.stringify({ name: team, ...units }),
});
if (patch.status === 200) return found;
GiteaAdmin.fail(`/teams/${found}`, patch);
}
if (found !== null) return found;
const res = await this.request(`/orgs/${encodeURIComponent(org)}/teams`, {
method: "POST",
body: JSON.stringify({ name: team, ...units }),
body: JSON.stringify({
name: team,
permission: "read",
// Package access is a per-unit grant; the team needs write on the packages unit and nothing
// else. includes_all_repositories keeps the team's repo view whole without widening its
// repo permission beyond read.
units_map: { "repo.packages": packageWrite ? "write" : "read" },
includes_all_repositories: true,
can_create_org_repo: false,
}),
});
if (res.status === 201) return Number(res.body?.id);
if (res.status === 422 || res.status === 409) {
@@ -423,18 +363,14 @@ export class GiteaAdmin {
}
private async findTeam(org: string, team: string): Promise<number | null> {
const res = await this.request(`/orgs/${encodeURIComponent(org)}/teams?limit=50`);
const res = await this.request(`/orgs/${encodeURIComponent(org)}/teams`);
if (res.status !== 200) return null;
const match = (res.body as any[] | null)?.find((t) => t?.name === team);
return match ? Number(match.id) : null;
}
/** Ensure a user exists with exactly this password. Created if absent; if already there, its
* password is patched — so the mesh minting a new secret takes on the next reconcile.
*
* The edit path is taken only when the user actually exists. A 422 from the create is also what
* a plain validation failure returns, and reading it as "already there" made the follow-up edit
* 404 — burying the create's own message, which is the one that says what is actually wrong. */
* password is patched — so the mesh minting a new secret takes on the next reconcile. */
async ensureUser(username: string, password: string, email: string): Promise<void> {
const res = await this.request("/admin/users", {
method: "POST",
@@ -442,18 +378,13 @@ export class GiteaAdmin {
});
if (res.status === 201) return;
if (res.status === 422 || res.status === 409) {
const seen = await this.request(`/users/${encodeURIComponent(username)}`);
if (seen.status === 200) {
const patch = await this.request(`/admin/users/${encodeURIComponent(username)}`, {
method: "PATCH",
// login_name is required by the admin edit endpoint; for a local user it is the username.
// active and prohibit_login: a deactivated or login-prohibited user is refused like a wrong
// password, so the provisioner's check reports it lost; applying again must undo both.
body: JSON.stringify({ login_name: username, password, must_change_password: false, active: true, prohibit_login: false }),
});
if (patch.status === 200) return;
GiteaAdmin.fail(`/admin/users/${username}`, patch);
}
const patch = await this.request(`/admin/users/${encodeURIComponent(username)}`, {
method: "PATCH",
// login_name is required by the admin edit endpoint; for a local user it is the username.
body: JSON.stringify({ login_name: username, password, must_change_password: false }),
});
if (patch.status === 200) return;
GiteaAdmin.fail(`/admin/users/${username}`, patch);
}
GiteaAdmin.fail("/admin/users", res);
}
@@ -468,29 +399,6 @@ export class GiteaAdmin {
GiteaAdmin.fail(`/teams/${teamId}/members/${username}`, res);
}
/**
* Whether a consumer's user logs in with exactly this password and is still a member of the
* package team. Read-only: the password is checked as the consumer presents it, basic auth on the
* API, and membership through the admin API. `false` for a refused login or a missing member; any
* other answer rejects (novox/hq issue 120).
*/
async holdsTeamMember(org: string, team: string, username: string, password: string): Promise<boolean> {
const me = await fetch(`${this.baseUrl}/api/v1/user`, {
headers: { Authorization: "Basic " + Buffer.from(`${username}:${password}`).toString("base64") },
});
if (me.status === 401 || me.status === 403) return false;
if (me.status !== 200) throw new Error(`Gitea GET /user as ${username}: ${me.status}`);
const teams = await this.request(`/orgs/${encodeURIComponent(org)}/teams?limit=50`);
if (teams.status === 404) return false;
if (teams.status !== 200) GiteaAdmin.fail(`/orgs/${org}/teams`, teams);
const found = (teams.body as { id: number; name: string }[]).find((t) => t.name === team);
if (!found) return false;
const member = await this.request(`/teams/${found.id}/members/${encodeURIComponent(username)}`);
if (member.status === 200 || member.status === 204) return true;
if (member.status === 404) return false;
GiteaAdmin.fail(`/teams/${found.id}/members/${username}`, member);
}
/** Delete a user, purging what they own. A 404 means the mesh already withdrew them — success, not
* an error, so a re-run of remove is safe. */
async deleteUser(username: string): Promise<void> {
+5 -98
View File
@@ -16,9 +16,7 @@
import { emit } from "@novox/mesh-sdk/events";
import { GiteaClient } from "./client.js";
// Without a way to a token — configured, or mintable with the admin account (token.ts) — there is
// nothing to watch; log and stay quiet rather than crash the runtime. With one, the first poll mints
// or reuses the token, so the runtime's start also shows what it did about it.
// Without a token there is nothing to watch; log and stay quiet rather than crash the runtime.
let gitea: GiteaClient | null = null;
try {
gitea = GiteaClient.fromEnv();
@@ -31,11 +29,11 @@ try {
const seen = new Set<string>();
let primed = false;
async function pollRepos(client: GiteaClient): Promise<void> {
const repos = await client.listAllRepos();
const repos = await client.listRepos(1, 50);
for (const repo of repos) {
if (!seen.has(repo.full_name)) {
if (primed) {
await emit("repo.created", {
await emit("module.gitea.repo.created", {
full_name: repo.full_name,
owner: repo.owner,
name: repo.name,
@@ -49,104 +47,13 @@ async function pollRepos(client: GiteaClient): Promise<void> {
primed = true;
}
// **A merge is announced whoever made it.** The merge tool below emits at the instant it acts; a
// merge made in the forge's own pages or over its API would emit nothing, and the mesh would go on
// believing every module current with its source (novox/hq 04-ISSUES/131). So merged pull requests
// are watched the way repositories are: what the forge holds, asked for on a tick, announced once.
// What has been announced is kept beside the module's state, so a restart does not announce the
// whole history again — and the first tick on a machine with no record announces nothing, because
// everything it sees then predates the watching.
import { existsSync, mkdirSync, readFileSync, renameSync, writeFileSync } from "node:fs";
import { join } from "node:path";
const mergedRecord = process.env.MESH_GITEA_STATE_DIR ? join(process.env.MESH_GITEA_STATE_DIR, "merged-announced.json") : null;
const announced = new Set<string>();
let primedMerges = false;
// since is the moment the watching began: a merge made before it is history, whatever page of the
// forge's listing it surfaces on. Without it, an old merge past the first page — pushed into view
// as newer pull requests were updated — was announced as if it had just happened, and the mesh
// rebuilt everything built from that repository, once per old merge (2026-09-28).
let since = "";
if (mergedRecord && existsSync(mergedRecord)) {
try {
const kept = JSON.parse(readFileSync(mergedRecord, "utf8")) as string[] | { announced: string[]; since: string };
const list = Array.isArray(kept) ? kept : kept.announced;
for (const sha of list) announced.add(sha);
since = Array.isArray(kept) ? new Date().toISOString() : kept.since;
primedMerges = true;
} catch {
// An unreadable record is treated as no record: prime again rather than re-announce history.
}
}
function keepAnnounced(): void {
if (!mergedRecord) return;
mkdirSync(join(mergedRecord, ".."), { recursive: true });
const tmp = mergedRecord + ".tmp";
writeFileSync(tmp, JSON.stringify({ announced: [...announced].slice(-2000), since }));
renameSync(tmp, mergedRecord);
}
async function pollMerged(client: GiteaClient): Promise<void> {
const repos = await client.listAllRepos();
let changed = false;
for (const repo of repos) {
const pulls = await client.listPullRequests(repo.owner, repo.name, { state: "closed", sort: "recentupdate", limit: "20" });
for (const pull of pulls) {
if (!pull.merged || !pull.merge_commit_sha || announced.has(pull.merge_commit_sha)) continue;
// Announced only if merged since the watching began; recorded either way, so it is looked
// at once.
const fresh = !!pull.merged_at && !!since && pull.merged_at > since;
if (primedMerges && fresh) {
// What it changed, asked for only now: a module is rebuilt because a file inside its own
// directory moved, and without this every module built from a repository is rebuilt for a
// change to any of them (novox/hq 04-ISSUES/131).
const changed = await client.listPullFiles(repo.owner, repo.name, pull.number);
await emit("pull.merged", {
owner: repo.owner,
repo: repo.name,
number: pull.number,
title: pull.title,
head: pull.head,
base: pull.base,
merge_commit_sha: pull.merge_commit_sha,
merged_at: pull.merged_at,
clone_url: repo.clone_url,
html_url: pull.html_url,
paths: changed.paths,
paths_truncated: changed.truncated,
});
// Said, because a trigger that fires silently is indistinguishable from one that did not
// fire (novox/hq 04-ISSUES/131) — this line is how an operator knows the mesh was told.
console.log(`[gitea] announced merge ${repo.full_name}#${pull.number} (${pull.merge_commit_sha.slice(0, 8)}) into ${pull.base}`);
}
announced.add(pull.merge_commit_sha);
changed = true;
}
}
if (!primedMerges) since = new Date().toISOString();
if (!primedMerges || changed) keepAnnounced();
primedMerges = true;
}
if (gitea) {
const client = gitea;
// A poll that fails says so once, not once a minute: the same reason repeating (the forge not up
// yet, the admin account refused on a restored forge) is one fact, and a recovery is worth a line.
let failing: string | null = null;
const tick = (fn: () => Promise<void>, everyMs: number): void => {
const run = (): void =>
void fn()
.then(() => {
if (failing !== null) console.log("[gitea] watching again");
failing = null;
})
.catch((err) => {
const why = err instanceof Error ? err.message : String(err);
if (why !== failing) console.error(`[gitea] not watching until this clears — ${why}`);
failing = why;
});
const run = (): void => void fn().catch((err) => console.error(`[gitea] ${err}`));
setInterval(run, everyMs);
run();
};
tick(() => pollRepos(client), 60_000);
tick(() => pollMerged(client), 30_000);
console.log("[gitea] watching for new repositories and merged pull requests");
console.log("[gitea] watching for new repositories");
}
+30 -64
View File
@@ -11,80 +11,56 @@
"name": "gitea"
},
"route": {
"web": {
"label": "git",
"endpoint": "web"
},
"internal-api-refused": {
"label": "git",
"path": "/api/internal",
"deny": true,
"priority": 100000
}
"label": "git",
"port": 3000
}
},
"binds": {
"postgres-database": "${dir:state}/database.json",
"route": "${dir:state}/route.json"
"postgres-database": "/var/lib/gitea/database.json",
"route": "/var/lib/gitea/route.json"
},
"secrets": {
"postgres-database": "${dir:state}/database.secret",
"postgres-database": "/var/lib/gitea/database.secret",
"secret": {
"internal-token": "${dir:state}/internal-token.secret",
"admin": "${dir:state}/admin.secret"
"internal-token": "/var/lib/gitea/internal-token.secret",
"admin": "/var/lib/gitea/admin.secret"
}
},
"capabilities": [
"container-runtime"
],
"emits": [
"repo.created",
"issue.opened",
"pull.merged"
"module.gitea.repo.created",
"module.gitea.issue.opened",
"module.gitea.pull.merged"
],
"listens": [
{
"name": "web",
"port": 3000,
"protocol": "tcp",
"from": "mesh",
"why": "the forge, over http"
},
{
"name": "ssh",
"port": 22,
"port": 2222,
"protocol": "tcp",
"from": "mesh",
"why": "git over ssh, gitea's own unmodified sshd. Published on the machine's own side at 222, the mesh's fixed public convention \u2014 not 22, which the machine's own daemon holds and a module does not take"
"why": "git over ssh. Not 22: the machine's own daemon holds that, and a module does not take it"
}
],
"serves": {
"npm-package-registry": {
"package-registry": {
"scheme": "http",
"port": 3000,
"npm-path": "/api/packages/novox/npm/"
},
"git": {
"scheme": "http",
"port": 3000
}
},
"receives": {
"npm-package-registry": "${dir:grants}/npm.json"
"package-registry": "/var/lib/gitea/grants/mesh.json"
},
"grants": {
"npm-package-registry": "${dir:grants}"
"package-registry": "/var/lib/gitea/grants"
},
"claims": [
{
"name": "npm-package-registry",
"scope": "mesh"
},
{
"name": "git",
"scope": "mesh"
}
],
"own-secrets": {
"broker": "/var/lib/mesh/gitea/broker"
},
@@ -95,33 +71,29 @@
"path": "/var/lib/mesh/gitea",
"mode": "0700"
},
{
"id": "runtime-state",
"type": "directory",
"path": "/var/lib/mesh/gitea/state",
"mode": "0700"
},
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
"path": "/var/lib/gitea",
"mode": "0700"
},
{
"id": "grants",
"type": "directory",
"path": "/var/lib/gitea/grants",
"mode": "0700"
},
{
"id": "server-env",
"type": "file",
"path": "${dir:state}/server.env",
"path": "/var/lib/gitea/server.env",
"mode": "0600",
"content": "GITEA__security__INTERNAL_TOKEN=${secret:internal-token}\nGITEA__database__DB_TYPE=postgres\nGITEA__database__HOST=${bound:postgres-database:at}:${bound:postgres-database:port}\nGITEA__database__NAME=${bound:postgres-database:as}\nGITEA__database__USER=${bound:postgres-database:as}\nGITEA__database__PASSWD=${secret:postgres-database}\n"
},
{
"id": "data",
"type": "directory",
"path": "/services/gitea/gitea",
"mode": "0700",
"owner": "1000:1000"
},
@@ -136,14 +108,14 @@
"USER_GID": "1000"
},
"env-file": [
"${dir:state}/server.env"
"/var/lib/gitea/server.env"
],
"ports": [
"3000",
"222:22"
"2222:22"
],
"volumes": [
"${dir:data}:/data"
"/services/gitea/gitea:/data"
],
"secrets-in-environment": "gitea honours GITEA__database__PASSWD__FILE and GITEA__security__INTERNAL_TOKEN__FILE; convertible, awaiting a bed that proves it"
},
@@ -159,11 +131,11 @@
"MESH_GITEA_ADMIN_USER": "mesh-admin"
},
"env-file": [
"${dir:state}/server.env"
"/var/lib/gitea/server.env"
],
"volumes": [
"${dir:data}:/data",
"${dir:state}/admin.secret:/run/secrets/admin:ro"
"/services/gitea/gitea:/data",
"/var/lib/gitea/admin.secret:/run/secrets/admin:ro"
],
"args": [
"/bin/sh",
@@ -188,9 +160,8 @@
"volumes": [
"/var/lib/mesh/gitea/broker:/run/secrets/broker:ro",
"/var/lib/mesh/gitea/config.json:/run/config/config.json:ro",
"${dir:grants}:${dir:grants}:ro",
"${dir:state}/admin.secret:/run/secrets/admin:ro",
"/var/lib/mesh/gitea/state:/run/state"
"/var/lib/gitea/grants:/var/lib/gitea/grants:ro",
"/var/lib/gitea/admin.secret:/run/secrets/admin:ro"
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
@@ -198,8 +169,7 @@
"MESH_GITEA_CONFIG_FILE": "/run/config/config.json",
"MESH_GITEA_ADMIN_USER": "mesh-admin",
"MESH_GITEA_ADMIN_PASSWORD_FILE": "/run/secrets/admin",
"MESH_GITEA_STATE_DIR": "/run/state",
"MESH_RECEIVES": "${dir:grants}/npm.json"
"MESH_RECEIVES": "/var/lib/gitea/grants/mesh.json"
},
"artifact": "runtime",
"restart-on": [
@@ -209,11 +179,7 @@
],
"provides": [
{
"name": "npm-package-registry",
"scope": "mesh"
},
{
"name": "git",
"name": "package-registry",
"scope": "mesh"
}
],
+1 -5
View File
@@ -4,12 +4,8 @@
"description": "gitea — git hosting. Its API client, tools and events live here (novox/hq ADR 0039).",
"type": "module",
"private": true,
"scripts": {
"build": "tsc client.ts token.ts index.ts provisioner/index.ts tools/index.ts --module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist",
"test": "npm run build && node --test --experimental-strip-types 'test/*.test.ts'"
},
"dependencies": {
"@novox/mesh-sdk": "^0.1.1"
"@novox/mesh-sdk": "^0.1.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
+7 -27
View File
@@ -1,15 +1,9 @@
// gitea's provisioner — the adapter that makes gitea a provider of the mesh
// `npm-package-registry` interface. The reconcile loop, the contributions file, and reading the
// mesh's minted password are the sdk harness's; this writes only the per-service half: how gitea
// creates and removes a consumer's npm credential (novox/hq ADR 0048/0076).
// gitea's provisioner — the adapter that makes gitea a provider of the mesh `package-registry`
// interface. The reconcile loop, the contributions file, and reading the mesh's minted password are
// the sdk harness's; this writes only the per-service half: how gitea creates and removes a
// consumer's npm credential (novox/hq ADR 0048/0076).
//
// **A package registry seat is one per ecosystem (novox/hq ADR 0109).** gitea holds the npm seat
// (ADR 0110). Adding cargo or PyPI is adding a provision — another `provides` entry, another
// `receives` path and another registration below — not widening this one. `git`, which gitea also
// provides, mints nothing and so registers nothing here: the mesh's own repositories are public,
// and a clone credential is not yet decided (ADR 0111).
//
// The `npm-package-registry` interface: a consumer authenticates to the npm registry at
// The `package-registry` interface: a consumer authenticates to the npm registry at
// `/api/packages/novox/npm/` with basic auth, as `as` with the password the mesh minted, and can
// read and write packages under the `@novox` scope. The registry's npm owner is the gitea org
// `novox`; a consumer is a gitea *user* placed on that org's package team.
@@ -32,11 +26,7 @@ const PACKAGE_TEAM = "packages";
const gitea = GiteaAdmin.fromEnv();
// Where this registration's contributions land comes from $MESH_RECEIVES, never a path written
// here: the mesh writes the file where the manifest's `receives` says, and a second copy of that
// path in code would drift from it. One variable carries one path, so a second registration in this
// module needs the mesh to say where each provision's file is — not yet possible, and not faked.
runProvisioner("npm-package-registry", {
runProvisioner("package-registry", {
async create(p: Provision): Promise<void> {
// The org and its package team are the same for every consumer; ensuring them per-create is
// idempotent and needs no separate bootstrap step.
@@ -44,21 +34,11 @@ runProvisioner("npm-package-registry", {
const teamId = await gitea.ensureTeam(ORG, PACKAGE_TEAM, true);
// The user carries the consumer's login and the mesh's minted password, set every run so a
// rotation takes. Membership of the package team is what grants read+write on packages.
//
// The address is gitea's own convention for one that is not real: its email validation
// requires a dotted domain, so `@localhost` was refused at create — the fault that had this
// grant retrying for a day — while `@noreply.localhost` is the shape gitea itself gives
// hidden addresses.
await gitea.ensureUser(p.as, p.password, `${p.as}@noreply.localhost`);
await gitea.ensureUser(p.as, p.password, `${p.as}@localhost`);
await gitea.addUserToTeam(teamId, p.as);
},
async remove(p: { as: string }): Promise<void> {
await gitea.deleteUser(p.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).
async holds(p: Provision): Promise<boolean> {
return gitea.holdsTeamMember(ORG, PACKAGE_TEAM, p.as, p.password);
},
});
-313
View File
@@ -1,313 +0,0 @@
// What holds the module to its own token (token.ts; hq issue 100, the forge's tools): minted with
// the delivered admin account on the first call and kept at 0600, reused on the next start, minted
// afresh when the forge rejects it or the kept file is gone, and a refused admin account reported in
// plain words rather than crash-looped. A configured token still wins. And the tools register once
// there is a way to a token at all — before one exists.
//
// The forge is a fake: the four routes the module touches, with the same status codes gitea gives.
// Run against the compiled module (npm test builds first), the way the runtime loads it.
import { test, after } from "node:test";
import assert from "node:assert/strict";
import { createServer, type IncomingMessage, type ServerResponse } from "node:http";
import { mkdtemp, readFile, rm, stat, writeFile } from "node:fs/promises";
import { tmpdir } from "node:os";
import { join } from "node:path";
import { collectTools } from "@novox/mesh-sdk/tools";
import { AdminRefused, MintedToken, TOKEN_SCOPES } from "../dist/token.js";
import { GiteaClient } from "../dist/client.js";
import "../dist/tools/index.js";
const ADMIN = "mesh-admin";
const PASSWORD = "the-vault-minted-this";
// ---- A fake forge: what the module sends, and what gitea would answer. ----
interface Forge {
url: string;
mints: number;
lastScopes: string[] | null;
tokens: Map<string, string>;
admins: Map<string, string>;
close(): Promise<void>;
}
function fakeForge(): Promise<Forge> {
const forge = {
mints: 0,
lastScopes: null as string[] | null,
tokens: new Map<string, string>(), // name -> value
scopesOf: new Map<string, string[]>(), // value -> scopes, so a route can enforce them like gitea does
admins: new Map([[ADMIN, PASSWORD]]),
};
// write:X implies read:X — gitea's own rule (models/auth/access_token_scope.go).
const covers = (scopes: string[], required: string): boolean =>
scopes.includes(required) || scopes.includes(`write:${required.split(":")[1]}`);
const json = (res: ServerResponse, status: number, body: unknown): void => {
res.writeHead(status, { "Content-Type": "application/json" });
res.end(body === null ? "" : JSON.stringify(body));
};
const body = (req: IncomingMessage): Promise<any> =>
new Promise((resolve) => {
let text = "";
req.on("data", (c) => (text += c));
req.on("end", () => resolve(text ? JSON.parse(text) : null));
});
const basic = (req: IncomingMessage): string | null => {
const h = req.headers.authorization ?? "";
if (!h.startsWith("Basic ")) return null;
const [user, pass] = Buffer.from(h.slice(6), "base64").toString().split(":");
return forge.admins.get(user) === pass ? user : null;
};
const server = createServer(async (req, res) => {
const url = new URL(req.url ?? "/", "http://fake");
const tokens = url.pathname.match(/^\/api\/v1\/users\/([^/]+)\/tokens(?:\/([^/]+))?$/);
if (tokens) {
const user = basic(req);
if (user === null || user !== decodeURIComponent(tokens[1])) return json(res, 401, { message: "auth required" });
if (req.method === "POST") {
const { name, scopes } = await body(req);
if (forge.tokens.has(name)) return json(res, 400, { message: "token name has already been used" });
forge.mints++;
forge.lastScopes = scopes;
const sha1 = `minted-${forge.mints}-${Math.random().toString(36).slice(2)}`;
forge.tokens.set(name, sha1);
forge.scopesOf.set(sha1, scopes);
return json(res, 201, { id: forge.mints, name, sha1, scopes, token_last_eight: sha1.slice(-8) });
}
if (req.method === "DELETE" && tokens[2]) {
const name = decodeURIComponent(tokens[2]);
if (!forge.tokens.has(name)) return json(res, 404, { message: "token not found" });
forge.tokens.delete(name);
return json(res, 204, null);
}
return json(res, 405, { message: "method not allowed" });
}
if (url.pathname === "/api/v1/user/repos") {
const h = req.headers.authorization ?? "";
const value = h.startsWith("token ") ? h.slice(6) : "";
if (![...forge.tokens.values()].includes(value)) return json(res, 401, { message: "token is required" });
// gitea 1.27.3: GET /user/repos sits under the `user` scope category, not `repository` —
// confirmed against the live forge. A token without read:user (or write:user) is refused here.
const scopes = forge.scopesOf.get(value) ?? [];
if (!covers(scopes, "read:user")) {
return json(res, 403, {
message: `token does not have at least one of required scope(s), required=[read:user]`,
});
}
return json(res, 200, [
{ full_name: "novox/hq", name: "hq", owner: { login: "novox" }, private: true, html_url: "http://fake/novox/hq" },
]);
}
return json(res, 404, { message: "no such route in the fake" });
});
return new Promise((resolve) => {
server.listen(0, "127.0.0.1", () => {
const { port } = server.address() as { port: number };
resolve({
url: `http://127.0.0.1:${port}`,
get mints() { return forge.mints; },
get lastScopes() { return forge.lastScopes; },
tokens: forge.tokens,
admins: forge.admins,
close: () => new Promise((r) => server.close(() => r())),
});
});
});
}
// ---- What the runtime's environment gives the module. ----
async function delivered(forge: Forge): Promise<{ env: NodeJS.ProcessEnv; file: string; logs: string[] }> {
const dir = await mkdtemp(join(tmpdir(), "gitea-"));
const passwordFile = join(dir, "admin.secret");
await writeFile(passwordFile, PASSWORD + "\n", { mode: 0o600 });
const state = join(dir, "state");
return {
env: {
MESH_GITEA_URL: forge.url,
MESH_GITEA_ADMIN_USER: ADMIN,
MESH_GITEA_ADMIN_PASSWORD_FILE: passwordFile,
MESH_GITEA_STATE_DIR: state,
},
file: join(state, "token"),
logs: [],
};
}
/** A client as a fresh process would build it: a new source over the kept file, its log captured. */
function minted(env: NodeJS.ProcessEnv, logs: string[]): GiteaClient {
const source = new MintedToken({
url: env.MESH_GITEA_URL!,
admin: env.MESH_GITEA_ADMIN_USER!,
passwordFile: env.MESH_GITEA_ADMIN_PASSWORD_FILE!,
file: join(env.MESH_GITEA_STATE_DIR!, "token"),
log: (l) => logs.push(l),
});
return new GiteaClient(env.MESH_GITEA_URL!, source);
}
const forge = await fakeForge();
after(() => forge.close());
test("first start: mints with the admin account, keeps the token at 0600, asks for two scopes only", async () => {
const { env, file, logs } = await delivered(forge);
const repos = await minted(env, logs).listRepos();
assert.equal(repos[0]?.full_name, "novox/hq");
assert.equal(forge.mints, 1);
assert.deepEqual(forge.lastScopes, ["write:repository", "write:issue", "read:user"]);
assert.deepEqual(forge.lastScopes, [...TOKEN_SCOPES]);
const token = forge.tokens.get("mesh-tools")!;
assert.equal(await readFile(file, "utf8"), token + "\n");
assert.equal((await stat(file)).mode & 0o777, 0o600);
// Said that it minted, and where it keeps it — never what it is.
assert.ok(logs.some((l) => l.startsWith("minted a token")), logs.join("\n"));
assert.ok(logs.every((l) => !l.includes(token) && !l.includes(PASSWORD)), logs.join("\n"));
});
test("second start: reuses the kept token, mints nothing", async () => {
const { env, logs } = await delivered(forge);
await minted(env, logs).listRepos();
const before = forge.mints;
const again: string[] = [];
await minted(env, again).listRepos();
assert.equal(forge.mints, before);
assert.ok(again.some((l) => l.startsWith("reusing the token kept at")), again.join("\n"));
assert.ok(again.every((l) => !l.includes(forge.tokens.get("mesh-tools")!)), again.join("\n"));
});
test("the forge rejects the kept token (its data was restored): minted afresh, once, and the call goes through", async () => {
const { env, file, logs } = await delivered(forge);
const client = minted(env, logs);
await client.listRepos();
const before = forge.mints;
forge.tokens.clear(); // the forge no longer knows any token — a restore from the predecessor
const repos = await client.listRepos();
assert.equal(repos.length, 1);
assert.equal(forge.mints, before + 1);
assert.equal(await readFile(file, "utf8"), forge.tokens.get("mesh-tools") + "\n");
assert.ok(logs.some((l) => l.startsWith("the forge rejected the kept token")), logs.join("\n"));
});
test("the kept file is gone but the forge still holds a token by that name: replaced, not refused", async () => {
const { env, file, logs } = await delivered(forge);
await minted(env, logs).listRepos();
const before = forge.mints;
await rm(file);
const repos = await minted(env, logs).listRepos();
assert.equal(repos.length, 1);
assert.equal(forge.mints, before + 1);
assert.equal([...forge.tokens.keys()].filter((n) => n === "mesh-tools").length, 1);
assert.ok(logs.some((l) => l.includes('already holds a token named "mesh-tools"')), logs.join("\n"));
});
test("concurrent first calls share one mint", async () => {
const { env, logs } = await delivered(forge);
const client = minted(env, logs);
const before = forge.mints;
await Promise.all([client.listRepos(), client.listRepos(), client.listRepos()]);
assert.equal(forge.mints, before + 1);
});
test("the admin account is refused: said plainly, nothing kept, and the next call fails the same way rather than crashing", async () => {
const { env, file, logs } = await delivered(forge);
forge.admins.delete(ADMIN); // the forge's data came from a predecessor; mesh-admin was never created there
try {
const client = minted(env, logs);
const before = forge.mints;
await assert.rejects(client.listRepos(), (err: unknown) => {
assert.ok(err instanceof AdminRefused, String(err));
assert.match(err.message, /refused the admin account "mesh-admin" \(401\)/);
assert.match(err.message, /admin-bootstrap step creates it/);
assert.match(err.message, /came from a predecessor/);
assert.ok(!err.message.includes(PASSWORD));
return true;
});
await assert.rejects(client.listRepos(), AdminRefused);
assert.equal(forge.mints, before);
await assert.rejects(stat(file), /ENOENT/);
// The account appears (the operator created it): the very next call mints and works.
forge.admins.set(ADMIN, PASSWORD);
assert.equal((await client.listRepos()).length, 1);
assert.equal(forge.mints, before + 1);
} finally {
forge.admins.set(ADMIN, PASSWORD);
}
});
test("one process shares one source per kept file — the watcher and the tools never renew against each other", async () => {
const { env } = await delivered(forge);
assert.equal(MintedToken.fromEnv(env.MESH_GITEA_URL!, env), MintedToken.fromEnv(env.MESH_GITEA_URL!, env));
});
test("a second process finds the token the first renewed, and reuses it instead of minting over it", async () => {
const { env, logs } = await delivered(forge);
const first = minted(env, logs);
const second = minted(env, logs);
await first.listRepos();
await second.listRepos(); // both hold the same kept token
const before = forge.mints;
forge.tokens.clear();
await first.listRepos(); // renews: one mint
await second.listRepos(); // rejected too — but the kept file already carries the renewed one
assert.equal(forge.mints, before + 1);
assert.ok(logs.some((l) => l.includes("is newer — reusing it")), logs.join("\n"));
});
test("a configured token wins, and is reported rather than minted over when the forge rejects it", async () => {
const { env } = await delivered(forge);
const before = forge.mints;
const client = GiteaClient.fromEnv({ ...env, MESH_GITEA_TOKEN: "one-somebody-pasted-in" });
await assert.rejects(client.listRepos(), /rejected the configured Gitea token \(401\)/);
assert.equal(forge.mints, before);
});
test("nothing to mint with and no token: the client says what is missing", async () => {
assert.throws(
() => GiteaClient.fromEnv({ MESH_GITEA_URL: forge.url, MESH_GITEA_ADMIN_USER: ADMIN }),
/set MESH_GITEA_TOKEN, or MESH_GITEA_ADMIN_PASSWORD_FILE, MESH_GITEA_STATE_DIR/,
);
});
test("the tools register once there is a way to a token, and the first call mints it", async () => {
const { env } = await delivered(forge);
const before = forge.mints;
const withAdmin = collectTools(env).find((c) => c.module === "gitea")!;
const withNothing = collectTools({}).find((c) => c.module === "gitea")!;
assert.equal(withNothing.tools.length, 0);
assert.deepEqual(
withAdmin.tools.map((t) => t.name),
[
"gitea_list_repos", "gitea_create_repo", "gitea_delete_repo",
"gitea_list_issues", "gitea_get_issue", "gitea_create_issue", "gitea_close_issue", "gitea_add_comment",
"gitea_list_pull_requests", "gitea_get_pull_request", "gitea_create_pull_request", "gitea_merge_pull_request",
"gitea_list_labels", "gitea_create_label",
"gitea_api",
],
);
assert.equal(forge.mints, before, "registering must not mint — the forge may not be up yet");
const result = (await withAdmin.tools.find((t) => t.name === "gitea_list_repos")!.run({})) as { repos: unknown[] };
assert.equal(result.repos.length, 1);
assert.equal(forge.mints, before + 1);
});
-253
View File
@@ -1,253 +0,0 @@
// The token the forge's tools and watcher authenticate with — and where it comes from.
//
// Nobody configures it. The forge is raised by the mesh, so there is no operator holding a token to
// paste in, and pasting one into settings would put a secret in the inventory in plaintext. What
// the mesh does deliver is the admin account: a login the manifest names and a password the vault
// minted and the host unsealed into a file (novox/hq ADR 0086). That account is enough to mint a
// token, so the module mints its own (hq issue 100, the forge's tools):
//
// - at first use, when none is kept: POST /users/{admin}/tokens over basic auth, with the two
// scopes the tools and the watcher need, and no more;
// - kept in the module's own state, a 0600 file, and read back on the next start — the forge
// hands a token's value out exactly once, so a token not kept is a token lost;
// - re-minted when the forge rejects it (401) or the kept file is gone. The one case that is not
// a fault: the forge's data was restored from a predecessor and the token the file names never
// existed there.
//
// An explicitly configured token still wins, and is never minted over: if it is rejected, that is
// reported, not repaired — somebody chose it.
//
// The token is never logged. Lines say that one was minted, reused or renewed, and where it is
// kept; never what it is.
import { chmodSync, mkdirSync, readFileSync, renameSync, writeFileSync } from "node:fs";
import { dirname, join } from "node:path";
/** The name the token carries in the forge's own list — one per mesh runtime, found by name. */
export const TOKEN_NAME = "mesh-tools";
/**
* The least the fifteen tools and the watcher need (gitea's route groups, 1.20+ scoped tokens):
* write:repository — create/delete repositories, pull requests (list/get/open/merge);
* write:issue — issues, comments, labels;
* read:user — GET /user/repos, which the watcher's poll and gitea_list_repos both call.
* It sits under the `user` category despite listing repositories, not `repository`
* — confirmed against the running forge (1.27.3), which answered
* `required=[read:user]` to a token carrying only the other two.
* Nothing under /admin, /orgs or write:user — the escape-hatch tool reaches only what these three cover.
*/
export const TOKEN_SCOPES: readonly string[] = ["write:repository", "write:issue", "read:user"];
/** Where a client's token comes from, and what to do when the forge says it is wrong. */
export interface TokenSource {
/** The token to authenticate with now; minted, read or configured. */
current(): Promise<string>;
/** The forge answered 401 to `rejected`. A fresh token, or a plain error when there is nothing to renew with. */
renew(rejected: string): Promise<string>;
}
/** A token somebody set — in settings or the environment. Never minted over. */
export class ConfiguredToken implements TokenSource {
constructor(private readonly token: string) {}
async current(): Promise<string> {
return this.token;
}
async renew(): Promise<string> {
throw new Error(
"the forge rejected the configured Gitea token (401). It was set explicitly (settings or MESH_GITEA_TOKEN), " +
"so the module does not mint over it — fix it, or unset it and the module mints its own",
);
}
}
/** The forge would not take the admin account: it is missing, or its password is not the one the mesh holds. */
export class AdminRefused extends Error {
constructor(admin: string, status: number) {
super(
`the forge refused the admin account "${admin}" (${status}) — it does not exist there, or its password is not ` +
`the one the vault delivered. The admin-bootstrap step creates it on a forge the mesh raised; a forge whose data ` +
`came from a predecessor does not have it. Create "${admin}" on the forge with the delivered password and the ` +
`token is minted on the next call — the tools stay registered and the watcher keeps trying`,
);
this.name = "AdminRefused";
}
}
export interface MintedTokenOptions {
/** The forge, e.g. http://127.0.0.1:3000. */
readonly url: string;
/** The admin login the manifest names. */
readonly admin: string;
/** The file the host unsealed the admin password into (ADR 0086). Read at mint time, so a rotation takes. */
readonly passwordFile: string;
/** Where the token is kept: a 0600 file in the module's own state. */
readonly file: string;
readonly name?: string;
readonly scopes?: readonly string[];
readonly log?: (line: string) => void;
readonly fetch?: typeof fetch;
}
/** The token the module mints for itself, kept in its state and renewed when the forge rejects it. */
export class MintedToken implements TokenSource {
private held: string | null = null;
private readFile = false;
private inflight: Promise<string> | null = null;
private readonly name: string;
private readonly scopes: readonly string[];
private readonly log: (line: string) => void;
private readonly fetchImpl: typeof fetch;
constructor(private readonly opts: MintedTokenOptions) {
this.name = opts.name ?? TOKEN_NAME;
this.scopes = opts.scopes ?? TOKEN_SCOPES;
this.log = opts.log ?? ((line) => console.log(`[gitea] ${line}`));
this.fetchImpl = opts.fetch ?? fetch;
}
/**
* Build from the runtime's environment: the forge's URL, the admin login and password file the
* manifest hands the runtime, and the module's state directory (MESH_GITEA_STATE_DIR, a directory
* the runtime mounts writable). Throws, naming what is missing, rather than hand back a source
* that cannot mint.
*/
static fromEnv(url: string, env: NodeJS.ProcessEnv = process.env): MintedToken {
const opts = MintedToken.optionsFromEnv(url, env);
// One source per kept file in a process. The watcher and the tools entrypoint both build a
// client in the same runtime; two sources over one file would each renew on a 401 and drop the
// other's token by name, forever. Shared, a renewal is one renewal.
const shared = MintedToken.shared.get(opts.file);
if (shared) return shared;
const source = new MintedToken(opts);
MintedToken.shared.set(opts.file, source);
return source;
}
private static readonly shared = new Map<string, MintedToken>();
private static optionsFromEnv(url: string, env: NodeJS.ProcessEnv): MintedTokenOptions {
const admin = env.MESH_GITEA_ADMIN_USER;
const passwordFile = env.MESH_GITEA_ADMIN_PASSWORD_FILE;
const stateDir = env.MESH_GITEA_STATE_DIR;
const missing = [
admin ? null : "MESH_GITEA_ADMIN_USER",
passwordFile ? null : "MESH_GITEA_ADMIN_PASSWORD_FILE",
stateDir ? null : "MESH_GITEA_STATE_DIR",
].filter((v): v is string => v !== null);
if (missing.length) {
throw new Error(`no Gitea token, and nothing to mint one with — set MESH_GITEA_TOKEN, or ${missing.join(", ")}`);
}
return { url, admin: admin!, passwordFile: passwordFile!, file: join(stateDir!, "token") };
}
async current(): Promise<string> {
if (this.held !== null) return this.held;
if (!this.readFile) {
this.readFile = true;
const kept = this.read();
if (kept !== null) {
this.held = kept;
this.log(`reusing the token kept at ${this.opts.file}`);
return kept;
}
}
return this.mint("no token kept — minting one");
}
async renew(rejected: string): Promise<string> {
// Another caller already renewed while this one was in flight with the old token.
if (this.held !== null && this.held !== rejected) return this.held;
// Or another process did, and kept it: use what is kept before minting over it.
const kept = this.read();
if (kept !== null && kept !== rejected) {
this.held = kept;
this.log(`the forge rejected the token held; the kept one at ${this.opts.file} is newer — reusing it`);
return kept;
}
this.held = null;
return this.mint("the forge rejected the kept token — minting a fresh one");
}
/** One mint at a time: concurrent first calls share it, rather than each minting its own. */
private mint(why: string): Promise<string> {
if (this.inflight === null) {
this.log(why);
this.inflight = this.doMint().finally(() => {
this.inflight = null;
});
}
return this.inflight;
}
private read(): string | null {
try {
const token = readFileSync(this.opts.file, "utf8").replace(/\n$/, "");
return token.length ? token : null;
} catch (err) {
if ((err as NodeJS.ErrnoException).code === "ENOENT") return null;
throw new Error(`cannot read the kept Gitea token at ${this.opts.file}: ${(err as Error).message}`);
}
}
/** Write the token at 0600, whole or not at all: a temp file beside it, then a rename. */
private keep(token: string): void {
mkdirSync(dirname(this.opts.file), { recursive: true, mode: 0o700 });
const tmp = `${this.opts.file}.tmp`;
writeFileSync(tmp, token + "\n", { mode: 0o600 });
chmodSync(tmp, 0o600);
renameSync(tmp, this.opts.file);
}
private async doMint(): Promise<string> {
let password: string;
try {
password = readFileSync(this.opts.passwordFile, "utf8").replace(/\n$/, "");
} catch (err) {
throw new Error(`cannot read the admin password at ${this.opts.passwordFile}: ${(err as Error).message}`);
}
const authorization = "Basic " + Buffer.from(`${this.opts.admin}:${password}`).toString("base64");
const tokens = `${this.opts.url.replace(/\/+$/, "")}/api/v1/users/${encodeURIComponent(this.opts.admin)}/tokens`;
const call = async (method: string, path = "", body?: unknown): Promise<{ status: number; body: any }> => {
const res = await this.fetchImpl(tokens + path, {
method,
headers: { "Content-Type": "application/json", Authorization: authorization },
...(body === undefined ? {} : { body: JSON.stringify(body) }),
});
const text = await res.text();
let parsed: any = null;
if (text) {
try { parsed = JSON.parse(text); } catch { parsed = text; }
}
return { status: res.status, body: parsed };
};
let res = await call("POST", "", { name: this.name, scopes: this.scopes });
if (res.status === 401 || res.status === 403) throw new AdminRefused(this.opts.admin, res.status);
if (res.status === 400 || res.status === 422) {
// The forge still holds a token by this name whose value we no longer have — the kept file
// went while the forge's data stayed. It is ours to replace: drop it by name and mint again.
this.log(`the forge already holds a token named "${this.name}" — replacing it`);
const dropped = await call("DELETE", `/${encodeURIComponent(this.name)}`);
if (dropped.status !== 204 && dropped.status !== 404) {
throw new Error(`Gitea DELETE /users/${this.opts.admin}/tokens/${this.name}: ${dropped.status} ${detail(dropped.body)}`);
}
res = await call("POST", "", { name: this.name, scopes: this.scopes });
}
if (res.status !== 201 && res.status !== 200) {
throw new Error(`Gitea POST /users/${this.opts.admin}/tokens: ${res.status} ${detail(res.body)}`);
}
const token = typeof res.body?.sha1 === "string" ? res.body.sha1 : null;
if (!token) throw new Error(`Gitea POST /users/${this.opts.admin}/tokens: ${res.status} but no token in the reply`);
this.keep(token);
this.held = token;
this.log(`minted a token for "${this.opts.admin}" (${this.scopes.join(", ")}), kept at ${this.opts.file}`);
return token;
}
}
function detail(body: unknown): string {
return typeof body === "string" ? body : JSON.stringify(body);
}
+5 -17
View File
@@ -124,7 +124,7 @@ export function getGiteaTools(gitea: GiteaClient): ToolDefinition[] {
labels: labelIds,
});
// The mesh just opened an issue — announce it the moment it exists.
await emit("issue.opened", {
await emit("module.gitea.issue.opened", {
owner,
repo,
number: issue.number,
@@ -231,24 +231,15 @@ export function getGiteaTools(gitea: GiteaClient): ToolDefinition[] {
// Read the PR first, so the merged event carries a title and branches, not just a number.
const pull = await gitea.getPullRequest(owner, repo, number);
await gitea.mergePullRequest(owner, repo, number, method, deleteBranch);
// Read it again: the merge commit only exists now, and it is what a build is made from.
const merged = await gitea.getPullRequest(owner, repo, number);
// And what it changed, so the mesh rebuilds the modules whose own files moved rather than
// every module built from the repository (novox/hq 04-ISSUES/131).
const changed = await gitea.listPullFiles(owner, repo, number);
await emit("pull.merged", {
await emit("module.gitea.pull.merged", {
owner,
repo,
number,
title: pull.title,
head: pull.head,
base: pull.base,
merge_commit_sha: merged.merge_commit_sha,
merged_at: merged.merged_at,
method,
html_url: pull.html_url,
paths: changed.paths,
paths_truncated: changed.truncated,
});
return { merged: true, number, method, deleted_branch: deleteBranch };
},
@@ -304,15 +295,12 @@ export function getGiteaTools(gitea: GiteaClient): ToolDefinition[] {
];
}
// The tools exist when the client has a way to a token: one configured, or the admin account to mint
// one with (token.ts). The mint itself happens on the first call, not here — a contributor is
// synchronous, and a forge not yet answering must not keep the runtime from serving. Without either
// way, gitea contributes none rather than failing the whole runtime, and says why.
// The tools exist only when a token can be found; without one, gitea contributes none rather than
// failing the whole runtime.
registerModuleTools("gitea", (env) => {
try {
return getGiteaTools(GiteaClient.fromEnv(env));
} catch (err) {
console.log(`[gitea] no tools — ${err instanceof Error ? err.message : String(err)}`);
} catch {
return [];
}
});
+1 -1
View File
@@ -8,5 +8,5 @@
"skipLibCheck": true,
"noEmit": true
},
"include": ["client.ts", "token.ts", "index.ts", "provisioner/index.ts", "tools/index.ts"]
"include": ["client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts"]
}
+1 -1
View File
@@ -40,7 +40,7 @@ async function pollAlerts(client: GrafanaClient): Promise<void> {
for (const key of now) {
if (!firing.has(key)) {
const a = byKey.get(key)!;
await emit("alert.firing", { name: a.name, labels: a.labels, activeAt: a.activeAt });
await emit("module.grafana.alert.firing", { name: a.name, labels: a.labels, activeAt: a.activeAt });
}
}
}
+2 -3
View File
@@ -2,7 +2,7 @@
"module": "grafana",
"version": "1",
"emits": [
"alert.firing"
"module.grafana.alert.firing"
],
"own-secrets": {
"admin": "/var/lib/grafana-module/admin.secret",
@@ -13,7 +13,6 @@
],
"listens": [
{
"name": "web",
"port": 3000,
"protocol": "tcp",
"from": "mesh",
@@ -97,7 +96,7 @@
"contributes": {
"route": {
"label": "grafana",
"endpoint": "web"
"port": 3000
}
},
"binds": {
+2 -3
View File
@@ -11,7 +11,7 @@
"contributes": {
"route": {
"label": "hello",
"endpoint": "web"
"port": 8080
}
},
"binds": {
@@ -19,7 +19,6 @@
},
"listens": [
{
"name": "web",
"port": 8080,
"protocol": "tcp",
"from": "mesh",
@@ -51,7 +50,7 @@
"name": "hello-web",
"network": "hello-web",
"ports": [
"8080"
"8080:8080"
],
"volumes": [
"/var/lib/hello-web/index.html:/www/index.html:ro"
+1 -1
View File
@@ -44,7 +44,7 @@ async function pollStates(): Promise<void> {
for (const s of states) {
const prev = lastState.get(s.entity_id);
if (primed && prev !== undefined && prev !== s.state) {
await emit("state.changed", {
await emit("module.home-assistant.state.changed", {
entity: s.entity_id,
name: nameOf(s),
from: prev,
+2 -3
View File
@@ -6,7 +6,7 @@
"container-runtime"
],
"emits": [
"state.changed"
"module.home-assistant.state.changed"
],
"own-secrets": {
"broker": "/var/lib/mesh/home-assistant/broker",
@@ -14,7 +14,6 @@
},
"listens": [
{
"name": "web",
"port": 8123,
"protocol": "tcp",
"from": "mesh",
@@ -86,7 +85,7 @@
"contributes": {
"route": {
"label": "home-assistant",
"endpoint": "web"
"port": 8123
}
},
"binds": {
+2 -2
View File
@@ -23,7 +23,7 @@ async function pollMounts(): Promise<void> {
if (primed) {
for (const [mount, m] of now) {
if (!live.has(mount)) {
await emit("stream.started", {
await emit("module.icecast.stream.started", {
mount,
name: m.name,
description: m.description,
@@ -33,7 +33,7 @@ async function pollMounts(): Promise<void> {
}
for (const [mount, m] of live) {
if (!now.has(mount)) {
await emit("stream.stopped", { mount, name: m.name });
await emit("module.icecast.stream.stopped", { mount, name: m.name });
}
}
}
+2 -3
View File
@@ -5,15 +5,14 @@
"container-runtime"
],
"emits": [
"stream.started",
"stream.stopped"
"module.icecast.stream.started",
"module.icecast.stream.stopped"
],
"own-secrets": {
"broker": "/var/lib/mesh/icecast/broker"
},
"listens": [
{
"name": "stream",
"port": 8000,
"protocol": "tcp",
"from": "mesh",
-1
View File
@@ -9,7 +9,6 @@
},
"listens": [
{
"name": "api",
"port": 8086,
"protocol": "tcp",
"from": "mesh",
+6 -11
View File
@@ -14,15 +14,12 @@
"mongodb-database": {
"name": "invoicing"
},
"s3-bucket": {
"bucket": "invoicing"
},
"route": {
"site": {
"label": "invoicing",
"endpoint": "web"
},
"api": {
"label": "invoicing-api",
"endpoint": "api"
}
"label": "invoicing",
"port": 80
}
},
"binds": {
@@ -36,14 +33,12 @@
},
"listens": [
{
"name": "web",
"port": 80,
"protocol": "tcp",
"from": "mesh",
"why": "the invoicing web frontend; a public name is a route grant later"
},
{
"name": "api",
"port": 9000,
"protocol": "tcp",
"from": "mesh",
@@ -68,7 +63,7 @@
"type": "file",
"path": "/var/lib/invoicing/api.env",
"mode": "0600",
"content": "NODE_ENV=production\nPORT=9000\nMONGO_URL=mongodb://${bound:mongodb-database:as}:${secret:mongodb-database}@${bound:mongodb-database:at}:${bound:mongodb-database:port}/${bound:mongodb-database:as}?authSource=${bound:mongodb-database:as}\nMONGO_DB=${bound:mongodb-database:as}\nMINIO_BUCKET=mesh-novox-invoice\nMINIO_ENDPOINT=${bound:s3-bucket:at}\nMINIO_PORT=${bound:s3-bucket:port}\nMINIO_ACCESSKEY=${bound:s3-bucket:as}\nMINIO_SECRET=${secret:s3-bucket}\n"
"content": "NODE_ENV=production\nPORT=9000\nMONGO_URL=mongodb://${bound:mongodb-database:as}:${secret:mongodb-database}@${bound:mongodb-database:at}:${bound:mongodb-database:port}/invoicing?authSource=admin\nMINIO_BUCKET=invoicing\nMINIO_ENDPOINT=${bound:s3-bucket:at}\nMINIO_PORT=${bound:s3-bucket:port}\nMINIO_ACCESSKEY=${bound:s3-bucket:as}\nMINIO_SECRET=${secret:s3-bucket}\n"
},
{
"id": "net",
+1 -2
View File
@@ -6,7 +6,6 @@
],
"listens": [
{
"name": "web",
"port": 9117,
"protocol": "tcp",
"from": "mesh",
@@ -83,7 +82,7 @@
"contributes": {
"route": {
"label": "indexers",
"endpoint": "web"
"port": 9117
}
},
"binds": {
+6 -6
View File
@@ -28,17 +28,17 @@ async function announce(type: string, body: Record<string, unknown>): Promise<vo
export const events = {
userCreated: (realm: string, username: string, email?: string) =>
announce("user.created", { realm, username, ...(email ? { email } : {}) }),
announce("module.keycloak.user.created", { realm, username, ...(email ? { email } : {}) }),
userDeleted: (realm: string, userId: string) =>
announce("user.deleted", { realm, userId }),
announce("module.keycloak.user.deleted", { realm, userId }),
passwordReset: (realm: string, userId: string) =>
announce("password.reset", { realm, userId }),
announce("module.keycloak.password.reset", { realm, userId }),
clientCreated: (realm: string, clientId: string, name?: string) =>
announce("client.created", { realm, clientId, ...(name ? { name } : {}) }),
announce("module.keycloak.client.created", { realm, clientId, ...(name ? { name } : {}) }),
groupCreated: (realm: string, name: string) =>
announce("group.created", { realm, name }),
announce("module.keycloak.group.created", { realm, name }),
roleCreated: (realm: string, name: string) =>
announce("role.created", { realm, name }),
announce("module.keycloak.role.created", { realm, name }),
};
console.log("[keycloak] event surface ready — identity, client, group and role changes are announced");
+9 -12
View File
@@ -11,7 +11,7 @@
},
"route": {
"label": "keycloak",
"endpoint": "web"
"port": 8080
}
},
"binds": {
@@ -25,16 +25,15 @@
"container-runtime"
],
"emits": [
"user.created",
"user.deleted",
"password.reset",
"client.created",
"group.created",
"role.created"
"module.keycloak.user.created",
"module.keycloak.user.deleted",
"module.keycloak.password.reset",
"module.keycloak.client.created",
"module.keycloak.group.created",
"module.keycloak.role.created"
],
"listens": [
{
"name": "web",
"port": 8080,
"protocol": "tcp",
"from": "mesh",
@@ -89,9 +88,7 @@
"env": {
"KC_DB": "postgres",
"KC_HTTP_ENABLED": "true",
"KC_HEALTH_ENABLED": "true",
"KC_HOSTNAME": "https://keycloak.novox.be",
"KC_PROXY_HEADERS": "xforwarded"
"KC_HEALTH_ENABLED": "true"
},
"env-file": [
"/var/lib/keycloak/admin.env",
@@ -121,7 +118,7 @@
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_KEYCLOAK_URL": "http://127.0.0.1:${port:8080}",
"MESH_KEYCLOAK_URL": "http://127.0.0.1:8080",
"MESH_KEYCLOAK_CONFIG_FILE": "/run/config/config.json"
},
"restart-on": [
+42
View File
@@ -0,0 +1,42 @@
# lavinmq's runtime: the tool runtime, carrying this module's compiled bootstrap, provisioner,
# tools and event consumer.
#
# **Built from this module's own directory and nothing else.** The sdk is in the base image, so
# nothing is copied out of a neighbouring checkout — which is what lets the mesh build this from a
# repository and a path (novox/hq ADR 0069) rather than only on a workstation that happens to have
# the siblings.
#
# Two bases, named rather than pinned: the image this is COMPILED in, and the image it RUNS in.
# They are different images on purpose — the first carries a compiler and the second must not, or
# every running container would carry one it never invokes. The mesh answers both with the copies it
# holds, because a fingerprint written here would name one particular copy and no other mesh has it
# (novox/hq issue 044). Declared in module.json's `build.on`; deliberately no defaults, so a build
# nobody told stops here and says which module to build first.
ARG BUILD_BASE
ARG RUNTIME_BASE
FROM ${BUILD_BASE} AS build
# Compiled under /app/modules so `@novox/mesh-sdk` resolves upward into the base's own
# node_modules — the module is compiled against exactly the sdk it will run against.
WORKDIR /app/modules/lavinmq
COPY . .
# The compiler is invoked by its real path rather than through node_modules/.bin, whose entries are
# symlinks to a launcher that requires its library relatively — resolved away when the base image
# was assembled.
#
# Four entrypoints and a client, because this module is four things: a run-once bootstrap that
# writes the broker's configuration before it first starts, a provisioner that grants consumers
# their own vhost and user, a set of tools, and an event consumer.
RUN node /app/node_modules/typescript/bin/tsc \
client.ts index.ts bootstrap/index.ts provisioner/index.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE}
COPY --from=build /app/modules/lavinmq/dist /app/modules/lavinmq/dist
# What the runtime loads from this module in serve mode: its event consumer, its tools, and its
# provisioner — all three in one process, so the provisioner's reconcile loop runs with the broker
# connected (novox/hq issues 060/061; the provisioner used to be run as a separate container `args`
# command, which meant it served no tools and, once, ran nowhere at all). The run-once bootstrap is
# NOT listed here — it is named in its own container's `args`, because it runs to completion before
# the broker starts rather than serving. One image, because they are one module and share a client.
ENV MESH_TOOL_MODULES=/app/modules/lavinmq/dist/index.js,/app/modules/lavinmq/dist/tools/index.js,/app/modules/lavinmq/dist/provisioner/index.js
+45
View File
@@ -0,0 +1,45 @@
// lavinmq's run-once bootstrap — lavinmq's own code (novox/hq ADR 0039), run once before the broker
// first starts (ADR 0052). lavinmq's default admin is set at first boot from a config file's
// `default_password_hash`, and that value is a HASH of the mesh-minted admin password, not the
// password itself — a form the mesh's plain-secret delivery cannot produce and no `${secret:...}`
// placeholder can compute. So this step computes it: it reads the admin password the mesh minted and
// the host unsealed, hashes it the way lavinmq expects (see client.rabbitHash), and writes the config
// file the broker container reads with `--config`. The host runs it to completion and only then
// starts the broker the manifest places after it — so the broker's first boot finds a config with an
// admin it can authenticate, and the provisioner (which reaches the management API as that admin)
// can do its work.
//
// It runs in the module's own runtime image, under the module's own account, as `mesh-tools run`
// imports it — no broker connection, because writing a config file is an offline operation and there
// is no broker to reach yet.
//
// lavinmq consults `default_user`/`default_password_hash` only on a first boot with an empty data
// dir; a later boot uses the persisted user database and ignores them. So this seeds the admin once,
// and a rotation of the admin secret does not re-key an already-initialised broker — the same
// first-boot-only shape the RabbitMQ-compatible default user has always had.
import { writeFileSync } from "node:fs";
import { readFileSync } from "node:fs";
import { rabbitHash } from "../client.js";
const adminUser = process.env.MESH_PROVISION_ADMIN_USER ?? process.env.MESH_LAVINMQ_ADMIN_USER ?? "mesh-admin";
const passwordFile = process.env.MESH_PROVISION_PASSWORD_FILE ?? process.env.MESH_LAVINMQ_ADMIN_PASSWORD_FILE ?? "/run/secrets/default";
const configOut = process.env.MESH_LAVINMQ_CONFIG_OUT ?? "/var/lib/lavinmq-module/lavinmq.ini";
const dataDir = process.env.MESH_LAVINMQ_DATA_DIR ?? "/var/lib/lavinmq";
const password = readFileSync(passwordFile, "utf8").replace(/\n$/, "");
if (!password) {
throw new Error(`[lavinmq:bootstrap] admin password file ${passwordFile} is empty — cannot seed the admin`);
}
// The broker reads only what it needs to authenticate its admin on first boot; bind/ports come from
// the container's entrypoint (`-b 0.0.0.0`), so this file names the admin and nothing else about the
// network.
const ini =
"[main]\n" +
`data_dir = ${dataDir}\n` +
`default_user = ${adminUser}\n` +
`default_password_hash = ${rabbitHash(password)}\n`;
writeFileSync(configOut, ini, { mode: 0o600 });
console.log(`[lavinmq:bootstrap] wrote ${configOut} with admin '${adminUser}' (password hashed for lavinmq)`);
+150
View File
@@ -0,0 +1,150 @@
// lavinmq's admin client — lavinmq's own code, living in the module (novox/hq ADR 0039). Both this
// module's tools and its provisioner import it, and nothing outside lavinmq does.
//
// It drives lavinmq through its HTTP management API (the RabbitMQ-compatible surface lavinmq serves
// on 15672), not a hand-rolled AMQP admin stack: the module may take NO npm dependency beyond
// @novox/mesh-sdk, and the management API is exactly the admin surface — create/remove a vhost, a
// user, and its permissions — reached with `fetch` (global on node 22) and HTTP Basic auth. One
// boundary, `api()`, and every method is built on it. This is the module's one impure seam, the way
// postgres's is `psql` and redis's is a RESP socket.
//
// **The login and password are the mesh's, not the provisioner's (novox/hq 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. lavinmq creates exactly that user with exactly that password on a vhost of the same
// name — a name or password the provisioner invented is one the consumer could never present.
import { createHash, randomBytes } from "node:crypto";
import { readFileSync } from "node:fs";
export interface LavinmqConn {
/** Base URL of the management API, e.g. http://lavinmq:15672 (no trailing /api). */
readonly base: string;
readonly adminUser: string;
readonly adminPassword: string;
}
export class LavinmqClient {
constructor(private readonly conn: LavinmqConn) {}
/**
* Build from the module's resolved environment. Reads MESH_LAVINMQ_* first (the documented
* names), falling back to the MESH_PROVISION_* keys the manifest already sets on the provisioner
* container so the module runs unchanged there. Throws if it cannot find a management endpoint and
* an admin password — the right failure, because without them nothing it does can work.
*/
static fromEnv(env: NodeJS.ProcessEnv = process.env): LavinmqClient {
const base = (env.MESH_LAVINMQ_MANAGEMENT ?? env.MESH_PROVISION_LAVINMQ ?? "").replace(/\/+$/, "");
const adminUser = env.MESH_LAVINMQ_ADMIN_USER ?? env.MESH_PROVISION_ADMIN_USER ?? "mesh-admin";
const adminPassword = env.MESH_LAVINMQ_ADMIN_PASSWORD ?? readSecretFile(env.MESH_PROVISION_PASSWORD_FILE);
if (!base || !adminPassword) {
throw new Error("lavinmq management endpoint or admin password is not set — lavinmq's own code cannot reach the server");
}
return new LavinmqClient({ base, adminUser, adminPassword });
}
/** One request against the management API. A non-2xx reply rejects, carrying the body for the log. */
async api(method: string, path: string, body?: unknown): Promise<unknown> {
const headers: Record<string, string> = {
Authorization: "Basic " + Buffer.from(`${this.conn.adminUser}:${this.conn.adminPassword}`).toString("base64"),
};
if (body !== undefined) headers["Content-Type"] = "application/json";
const resp = await fetch(`${this.conn.base}/api${path}`, {
method,
headers,
body: body !== undefined ? JSON.stringify(body) : undefined,
});
if (!resp.ok) {
throw new Error(`lavinmq management API ${method} ${path} -> ${resp.status}: ${await resp.text()}`);
}
const text = await resp.text();
return text ? JSON.parse(text) : null;
}
/** True once the management API answers — the server has finished starting. */
async ready(): Promise<boolean> {
try {
await this.api("GET", "/overview");
return true;
} catch {
return false;
}
}
/** Block until the management API answers, or throw once the budget is spent. */
async waitReady(retries = 30, delayMs = 1000): Promise<void> {
for (let i = 0; i < retries; i++) {
if (await this.ready()) return;
await new Promise((r) => setTimeout(r, delayMs));
}
throw new Error("lavinmq management API did not become ready");
}
/**
* Create (or reset to a known state) one consumer's broker: a vhost and a user both named for the
* consumer's login, with the login owning full permissions on exactly that vhost. Idempotent — a
* PUT of a vhost or user that exists is a no-op or a password reset, so a reconcile can call it
* again without harm. The consumer connects as `<login>` to vhost `<login>` and can reach nothing
* else (novox/hq ADR 0048).
*/
async createConsumer(login: string, password: string): Promise<void> {
const v = encodeURIComponent(login);
const u = encodeURIComponent(login);
await this.api("PUT", `/vhosts/${v}`);
await this.api("PUT", `/users/${u}`, { password, tags: "" });
await this.api("PUT", `/permissions/${v}/${u}`, { configure: ".*", write: ".*", read: ".*" });
}
/** Remove a consumer's vhost and user, idempotently. A DELETE of what is already gone is tolerated. */
async removeConsumer(login: string): Promise<void> {
const v = encodeURIComponent(login);
const u = encodeURIComponent(login);
try {
await this.api("DELETE", `/vhosts/${v}`);
} catch (err) {
console.error(`[lavinmq] delete vhost ${login} failed (continuing): ${err}`);
}
try {
await this.api("DELETE", `/users/${u}`);
} catch (err) {
console.error(`[lavinmq] delete user ${login} failed (continuing): ${err}`);
}
}
/** The vhosts, for the amqp_list_vhosts tool. */
async listVhosts(): Promise<{ name: string; messages: number }[]> {
const vhosts = (await this.api("GET", "/vhosts")) as { name: string; messages?: number }[];
return vhosts.map((v) => ({ name: v.name, messages: v.messages ?? 0 }));
}
/** The queues on one vhost (default the root vhost), for the amqp_list_queues tool. */
async listQueues(vhost = "/"): Promise<{ name: string; messages: number; consumers: number }[]> {
const queues = (await this.api("GET", `/queues/${encodeURIComponent(vhost)}`)) as
{ name: string; messages?: number; consumers?: number }[];
return queues.map((q) => ({ name: q.name, messages: q.messages ?? 0, consumers: q.consumers ?? 0 }));
}
}
/**
* The RabbitMQ-compatible SHA-256 password hash lavinmq's `default_password_hash` expects:
* base64( salt[4] || sha256( salt || utf8(password) ) ). The salt is any four bytes — random here,
* because a fixed salt buys nothing and a fresh one is free. Verified against `lavinmqctl
* hash_password`: a hash produced here is accepted by the server unchanged.
*/
export function rabbitHash(password: string, salt: Buffer = randomBytes(4)): string {
const digest = createHash("sha256").update(Buffer.concat([salt, Buffer.from(password, "utf8")])).digest();
return Buffer.concat([salt, digest]).toString("base64");
}
/** Generate a URL-safe password. */
export function generatePassword(): string {
return randomBytes(24).toString("base64url");
}
function readSecretFile(path: string | undefined): string | undefined {
if (!path) return undefined;
try {
return readFileSync(path, "utf8").trim();
} catch {
return undefined;
}
}
+25
View File
@@ -0,0 +1,25 @@
// lavinmq's events entrypoint, loaded by the per-node tool host (the provisioner container runs
// ./provisioner separately). The broker lifecycle events are EMITTED from the provisioner, where the
// lifecycle actually happens (novox/hq ADR 0041/0042):
// module.lavinmq.amqp.provisioned — a consumer's vhost + user was created
// module.lavinmq.amqp.deprovisioned — that vhost + user was removed
// Here in the tool host we react to them, keeping a lightweight audit trail of who was granted a
// broker and who lost one — observability the provider itself is best placed to log.
import { on } from "@novox/mesh-sdk/events";
interface AmqpEvent {
consumer?: string;
user: string;
vhost?: string;
}
await on<AmqpEvent>("module.lavinmq.amqp.provisioned", async (e) => {
console.log(`[lavinmq] broker provisioned for ${e.body.consumer ?? "?"} (user ${e.body.user}, vhost ${e.body.vhost})`);
});
await on<AmqpEvent>("module.lavinmq.amqp.deprovisioned", async (e) => {
console.log(`[lavinmq] broker deprovisioned (user ${e.body.user})`);
});
console.log("[lavinmq] auditing broker lifecycle events");
+139
View File
@@ -0,0 +1,139 @@
{
"module": "lavinmq",
"version": "1",
"provides": [
{
"name": "amqp",
"scope": "mesh"
}
],
"claims": [
{
"name": "mesh-broker",
"scope": "mesh"
}
],
"capabilities": [
"container-runtime"
],
"emits": [
"module.lavinmq.amqp.provisioned",
"module.lavinmq.amqp.deprovisioned"
],
"consumes": [
"module.lavinmq.amqp.provisioned",
"module.lavinmq.amqp.deprovisioned"
],
"serves": {
"amqp": {
"port": 5672
}
},
"receives": {
"amqp": "/var/lib/lavinmq-module/grants/mesh.json"
},
"grants": {
"amqp": "/var/lib/lavinmq-module/grants"
},
"own-secrets": {
"admin": "/var/lib/lavinmq-module/admin.secret",
"broker": "/var/lib/mesh/lavinmq/broker"
},
"listens": [
{
"port": 5671,
"protocol": "tcp",
"from": "mesh",
"why": "the mesh bus \u2014 amqps, every module's events and the control plane, reached over the overlay"
},
{
"port": 5672,
"protocol": "tcp",
"from": "mesh",
"why": "modules on any machine that were granted a queue"
}
],
"guards": [
15672
],
"resources": [
{
"id": "mesh-state",
"type": "directory",
"path": "/var/lib/mesh/lavinmq",
"mode": "0700"
},
{
"id": "state",
"type": "directory",
"path": "/var/lib/lavinmq-module",
"mode": "0700"
},
{
"id": "grants-dir",
"type": "directory",
"path": "/var/lib/lavinmq-module/grants",
"mode": "0700"
},
{
"id": "server",
"type": "container",
"name": "mesh-broker",
"image": "cloudamqp/lavinmq@sha256:3eb54c12916d700a978c2ea86e6362cd4974b0e3189508718006d4e6d341246b",
"ports": [
"5671:5671",
"5672:5672",
"127.0.0.1:15672:15672"
],
"volumes": [
"mesh-broker-data:/var/lib/lavinmq",
"mesh-broker-tls:/tls:ro"
],
"args": [
"--amqps-port=5671",
"--cert=/tls/tls.crt",
"--key=/tls/tls.key"
]
},
{
"id": "runtime",
"type": "container",
"name": "mesh-lavinmq",
"artifact": "runtime",
"network": "host",
"volumes": [
"/var/lib/mesh/lavinmq/broker:/run/secrets/broker:ro",
"/var/lib/lavinmq-module/grants:/var/lib/lavinmq-module/grants:ro",
"/var/lib/lavinmq-module/admin.secret:/run/secrets/admin:ro"
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_RECEIVES": "/var/lib/lavinmq-module/grants/mesh.json",
"MESH_PROVISION_LAVINMQ": "http://127.0.0.1:15672",
"MESH_PROVISION_ADMIN_USER": "guest",
"MESH_PROVISION_PASSWORD_FILE": "/run/secrets/admin"
}
}
],
"build": {
"on": [
{
"arg": "BUILD_BASE",
"module": "mesh-tools",
"artifact": "build"
},
{
"arg": "RUNTIME_BASE",
"module": "mesh-tools",
"artifact": "runtime"
}
],
"artifacts": [
{
"name": "runtime",
"kind": "image",
"from": "Dockerfile"
}
]
}
}
+14
View File
@@ -0,0 +1,14 @@
{
"name": "@novox/module-lavinmq",
"version": "0.1.0",
"description": "lavinmq — provides the mesh amqp interface (a per-consumer AMQP message broker). Its management client, provisioner, run-once bootstrap, tools and events live here (novox/hq ADR 0039).",
"type": "module",
"private": true,
"dependencies": {
"@novox/mesh-sdk": "^0.1.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
"typescript": "^5.6.0"
}
}
+51
View File
@@ -0,0 +1,51 @@
// lavinmq's provisioner — the adapter that makes lavinmq a provider of the mesh `amqp` interface.
// The reconcile loop, the contributions file, and reading the mesh's minted password are the sdk
// harness's; this writes only the per-service half: how lavinmq creates and removes a consumer's own
// broker (novox/hq ADR 0039/0040/0048).
//
// The `amqp` interface: a consumer connects as `as` with the password the mesh minted, to a vhost
// named for that same login — its own message broker, isolated from every other consumer's by the
// vhost boundary. It is a broker of its own, not a shared account on the mesh's control-plane broker.
//
// **The login and password are the mesh's, not the provisioner's (novox/hq 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. lavinmq creates exactly that user with exactly that password — a name or password the
// provisioner invented is one the consumer could never present.
//
// Vhost-per-login is the isolation model, the exact analog of postgres's database-per-login: the
// consumer owns one vhost, named for its login, and a user with full rights on that vhost and no
// rights anywhere else. lavinmq enforces it — a user with no permission on `/` is refused the moment
// it opens that vhost (`NOT_ALLOWED`), so a login is a broker the consumer alone can reach.
import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner";
import { emit } from "@novox/mesh-sdk/events";
import { LavinmqClient } from "../client.js";
const lavinmq = LavinmqClient.fromEnv();
/** Emit a lifecycle event without letting a broker hiccup fail the provisioning itself. */
async function announce(type: string, body: Record<string, string>): Promise<void> {
try {
await emit(type, body);
} catch (err) {
console.error(`[provisioner:amqp] emit ${type} failed: ${err}`);
}
}
runProvisioner("amqp", {
async create(p: Provision): Promise<void> {
// The vhost and the user share the consumer's login, so one cannot reach another's broker.
await lavinmq.waitReady();
await lavinmq.createConsumer(p.as, p.password);
await announce("module.lavinmq.amqp.provisioned", {
consumer: p.consumer ?? "",
user: p.as,
vhost: p.as,
});
},
async remove(p: { as: string }): Promise<void> {
await lavinmq.removeConsumer(p.as);
await announce("module.lavinmq.amqp.deprovisioned", { user: p.as, vhost: p.as });
},
});
+32
View File
@@ -0,0 +1,32 @@
// lavinmq's tools — lavinmq's own code (novox/hq ADR 0039), importing lavinmq's own management
// client. They return structured data; the mesh serves them through the sdk's tool harness.
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
import { LavinmqClient } from "../client.js";
export function getLavinmqTools(lavinmq: LavinmqClient): ToolDefinition[] {
return [
{
name: "amqp_list_vhosts",
description: "List the lavinmq virtual hosts — one per consumer that was granted a broker.",
input: {},
run: async () => ({ vhosts: await lavinmq.listVhosts() }),
},
{
name: "amqp_list_queues",
description: "List the queues on a vhost with message and consumer counts. Omit vhost for the root '/'.",
input: { vhost: { type: "string", description: "the vhost to list, e.g. a consumer's login; defaults to '/'" } },
run: async (args) => ({ queues: await lavinmq.listQueues(args.vhost ? String(args.vhost) : undefined) }),
},
];
}
// The tools exist only when the server can be reached from the environment; without it, lavinmq
// contributes none rather than failing the whole tool runtime.
registerModuleTools("lavinmq", (env) => {
try {
return getLavinmqTools(LavinmqClient.fromEnv(env));
} catch {
return [];
}
});
+12
View File
@@ -0,0 +1,12 @@
{
"compilerOptions": {
"target": "ES2022",
"module": "NodeNext",
"moduleResolution": "NodeNext",
"strict": true,
"esModuleInterop": true,
"skipLibCheck": true,
"noEmit": true
},
"include": ["client.ts", "index.ts", "provisioner/index.ts", "tools/index.ts", "bootstrap/index.ts"]
}
-1
View File
@@ -24,7 +24,6 @@
},
"listens": [
{
"name": "web",
"port": 8283,
"protocol": "tcp",
"from": "mesh",
+2 -2
View File
@@ -40,12 +40,12 @@ async function pollQueue(lidarr: LidarrClient): Promise<void> {
if (primed) {
// Entered the queue since last look — Lidarr grabbed a release.
for (const [id, item] of now) {
if (!inQueue.has(id)) await emit("album.grabbed", { title: item.title, status: item.status });
if (!inQueue.has(id)) await emit("module.lidarr.album.grabbed", { title: item.title, status: item.status });
}
// Left the queue — imported and done, unless it was last seen failing.
for (const [id, item] of inQueue) {
if (!now.has(id) && !FAILED_STATUSES.has(item.status)) {
await emit("download.completed", { title: item.title });
await emit("module.lidarr.download.completed", { title: item.title });
}
}
}
+4 -5
View File
@@ -5,8 +5,8 @@
"container-runtime"
],
"emits": [
"album.grabbed",
"download.completed"
"module.lidarr.album.grabbed",
"module.lidarr.download.completed"
],
"consumes": [],
"own-secrets": {
@@ -14,7 +14,6 @@
},
"listens": [
{
"name": "web",
"port": 8686,
"protocol": "tcp",
"from": "mesh",
@@ -75,7 +74,7 @@
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_LIDARR_URL": "http://127.0.0.1:${port:8686}",
"MESH_LIDARR_URL": "http://127.0.0.1:8686",
"MESH_LIDARR_CONFIG_DIR": "/var/lib/lidarr/config"
},
"artifact": "runtime"
@@ -87,7 +86,7 @@
"contributes": {
"route": {
"label": "lidarr",
"endpoint": "web"
"port": 8686
}
},
"binds": {
+2 -2
View File
@@ -17,7 +17,7 @@ FROM ${BUILD_BASE} AS build
# resolved away.
WORKDIR /app/modules/mailu
COPY . .
RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts tools/index.ts provisioner/index.ts \
RUN node /app/node_modules/typescript/bin/tsc client.ts index.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
FROM ${RUNTIME_BASE}
@@ -27,4 +27,4 @@ COPY --from=build /app/modules/mailu/dist /app/modules/mailu/dist
# the convention novox/hq issues 060/061 settled. A container that instead ran only its
# provisioner (`run`) served no tools and emitted no events; a container that named no command
# ran no provisioner at all.
ENV MESH_TOOL_MODULES=/app/modules/mailu/dist/index.js,/app/modules/mailu/dist/tools/index.js,/app/modules/mailu/dist/provisioner/index.js
ENV MESH_TOOL_MODULES=/app/modules/mailu/dist/index.js,/app/modules/mailu/dist/tools/index.js
-44
View File
@@ -1,44 +0,0 @@
# automx2 — the autoconfig/autodiscover responder, carried by the mailu module as its own
# artifact: it is a config-baked sidecar of this mail server, not a standalone application
# (novox/hq ADR 0015 draws that line at applications).
#
# The base is named rather than pinned (novox/hq issue 044): declared in module.json's
# `build.on`. The build context is the module's own directory; every ADD says so.
ARG PYTHON_BASE
FROM ${PYTHON_BASE}
RUN apk add --no-cache bash sqlite
WORKDIR /automx2
ADD automx/files/setupvenv.sh /automx2/setupvenv.sh
ADD automx/files/start /automx2/start
ADD automx/files/setup /automx2/setup
ADD automx/files/setup-db /automx2/setup-db
ADD automx/files/add-domains /automx2/add-domains
RUN chmod u+x setupvenv.sh start add-domains setup setup-db
RUN ./setupvenv.sh \
&& . .venv/bin/activate \
&& pip install automx2==2021.6
# The launcher `start` expects. In the predecessor's image this wrapper appeared during a build
# step that never made it into the files this module carries — the image worked and the recipe
# could not reproduce it. Written here explicitly, verbatim from the proven image, so the build
# is the whole truth about the image again.
RUN mkdir -p .venv/scripts && printf '%s\n' \
'#!/usr/bin/env bash' \
'set -euo pipefail' \
'. .venv/bin/activate' \
"export FLASK_ENV='production'" \
"export FLASK_APP='automx2.server:app'" \
'flask "$@"' > .venv/scripts/flask.sh && chmod +x .venv/scripts/flask.sh
ENV AUTOMX2_CONF=/etc/automx2.conf
ADD automx/files/automx2.conf /etc/automx2.conf
# VOLUME deliberately absent: the anonymous /data volume is exactly what lost db.sqlite on
# every recreate (measured on novox 2026-08-10). The manifest binds a real directory instead.
ENTRYPOINT ["/bin/sh"]
CMD ["./start"]
EXPOSE 4243
-49
View File
@@ -1,49 +0,0 @@
#!/usr/bin/env bash
set -e
echo "${MAIL_DOMAINS}"
# Split domains into array
IFS=', ' read -r -a array <<< "${AMX_MAIL_DOMAINS}"
# User configurable section -- START
PROVIDER_ID=001
SQL_CMD="";
# Iterate domains resulting from split on second arg
for element in "${array[@]}"
do
# Set vars
DOMAIN=$element
PROVIDER_NAME=$DOMAIN
PROVIDER_SHORTNAME=$DOMAIN
# Optional LDAP server
#LDAP_SERVER="ldap.${DOMAIN}"
# User configurable section -- END
s1_id=$((PROVIDER_ID + 1))
s2_id=$((PROVIDER_ID + 2))
s3_id=$((PROVIDER_ID + 3))
dom_id=$((PROVIDER_ID + 4))
s3_id='NULL'
SQL_CMD=$(cat <<EOT
$SQL_CMD
INSERT INTO provider(id, name, short_name) VALUES(${PROVIDER_ID}, '${PROVIDER_NAME}', '${PROVIDER_SHORTNAME}');
INSERT INTO server(id, port, type, name, socket_type, user_name, authentication)
VALUES(${s1_id}, ${AMX_IMAP_PORT}, 'imap', '${AMX_IMAP_SERVER}', 'STARTTLS', '%EMAILADDRESS%', 'password-cleartext');
INSERT INTO server(id, port, type, name, socket_type, user_name, authentication)
VALUES(${s2_id}, ${AMX_SMTP_PORT}, 'smtp', '${AMX_SMTP_ADDRESS}', 'STARTTLS', '%EMAILADDRESS%', 'password-cleartext');
INSERT INTO domain(id, name, provider_id, ldapserver_id) VALUES(${dom_id}, '${DOMAIN}', ${PROVIDER_ID}, ${s3_id});
INSERT INTO server_domain(server_id, domain_id) VALUES(${s1_id}, ${dom_id});
INSERT INTO server_domain(server_id, domain_id) VALUES(${s2_id}, ${dom_id});
EOT
)
PROVIDER_ID=$((PROVIDER_ID+10))
done
echo -e ${SQL_CMD}
echo -e ${SQL_CMD} | sqlite3 /data/db.sqlite
-21
View File
@@ -1,21 +0,0 @@
[automx2]
# A typical production setup would use loglevel = WARNING
loglevel = WARNING
# Echo SQL commands into log? Used for debugging.
db_echo = false
# In-memory SQLite database
# db_uri = sqlite:///:memory:
# SQLite database in a UNIX-like file system
db_uri = sqlite:////data/db.sqlite
# MySQL database on a remote server. This example does not use an encrypted
# connection and is therefore *not* recommended for production use.
#db_uri = mysql://username:password@server.example.com/db
# Number of proxy servers between automx2 and the client (default: 0).
# If your logs only show 127.0.0.1 or ::1 as the source IP for incoming
# connections, proxy_count probably needs to be changed.
proxy_count = 1
-12
View File
@@ -1,12 +0,0 @@
#!/usr/bin/env bash
set -e
if [ ! -e /data/db.sqlite ]; then
# DB SETUP
echo "SETTING UP DB"
./setup-db
echo "ADDING DOMAINS"
# Add the domains
./add-domains
fi
-84
View File
@@ -1,84 +0,0 @@
#!/usr/bin/env bash
set -e
# LDAP-Server
LDAP=$(cat <<EOT
CREATE TABLE ldapserver(
id INT PRIMARY KEY NOT NULL,
name TEXT NOT NULL,
port INT NOT NULL,
use_ssl INT NOT NULL,
search_base TEXT NOT NULL,
search_filter TEXT NOT NULL,
attr_uid TEXT NOT NULL,
attr_cn TEXT NOT NULL,
bind_password TEXT NOT NULL,
bind_user TEXT NOT NULL
);
EOT
)
# Provider
PROVIDER=$(cat <<EOT
CREATE TABLE provider(
id INT PRIMARY KEY NOT NULL,
name TEXT NOT NULL,
short_name TEXT NOT NULL
);
EOT
)
# Server
SERVER=$(cat <<EOT
CREATE TABLE server(
id INT PRIMARY KEY NOT NULL,
prio INT NOT NULL DEFAULT 10,
name TEXT NOT NULL,
port INT NOT NULL,
type TEXT NOT NULL,
socket_type TEXT NOT NULL,
user_name TEXT NOT NULL,
authentication TEXT NOT NULL
);
EOT
)
# Domain
DOMAIN=$(cat <<EOT
CREATE TABLE domain(
id INT PRIMARY KEY NOT NULL,
name TEXT NOT NULL,
provider_id INT NOT NULL,
ldapserver_id INT NULL,
FOREIGN KEY(ldapserver_id) REFERENCES ldapserver(id),
FOREIGN KEY(provider_id) REFERENCES provider(id)
);
CREATE UNIQUE INDEX domain_name ON domain(name);
EOT
)
# Server-Domain
SERVER_DOMAIN=$(cat <<EOT
CREATE TABLE server_domain(
server_id INT NOT NULL,
domain_id INT NOT NULL,
FOREIGN KEY(server_id) REFERENCES server(id),
FOREIGN KEY(domain_id) REFERENCES domain(id)
);
EOT
)
## TODO Foreign keys
SQL_CMD=$(cat <<EOT
$LDAP
$PROVIDER
$SERVER
$DOMAIN
$SERVER_DOMAIN
EOT
)
echo -e ${SQL_CMD}
echo -e ${SQL_CMD} | sqlite3 /data/db.sqlite
-38
View File
@@ -1,38 +0,0 @@
#!/usr/bin/env bash
# vim:ts=4:sw=4:noet
#
# Creates a Python 3 virtual environment. The target directory can be passed
# as a parameter. The default path is '.venv' in the current directory.
dir="${1:-.venv}"
echo "Setup dir $dir"
set -e
if [ -d "${dir}" ]; then
echo >&2 "Directory '${dir}' already exists, exiting."
exit 1
fi
python3 -m venv "${dir}"
source "${dir}/bin/activate"
set +e
pip install -U pip setuptools wheel || true
#set -e
## vim:tabstop=4:noexpandtab
##
## Creates a Python 3 virtual environment. The target directory can be passed
## as a parameter. The default path is 'venv' in the current directory.
#
#dir="${1:-venv}"
#
#set -e
#if [ -d "${dir}" ]; then
# echo "Directory '${dir}' already exists, exiting." >&2
# exit 1
#fi
#python3 -m venv "${dir}"
#. "${dir}/bin/activate"
#
#set +e
#pip install -U pip setuptools || true
-8
View File
@@ -1,8 +0,0 @@
#!/usr/bin/env bash
set -e
# Setup
./setup
# Start
./.venv/scripts/flask.sh run --host=0.0.0.0 --port=4243
-29
View File
@@ -130,39 +130,10 @@ export class MailuClient {
await this.api("PATCH", `/user/${encodeURIComponent(email)}`, { raw_password: password });
}
/**
* Set the mesh's password on a mailbox the mesh provisions, and enable it. A disabled mailbox is
* what the provisioner's check reports as lost, so applying again must enable it, or the two would
* disagree for ever. Separate from changePassword, which an operator's tool uses and which must
* not re-enable a mailbox someone disabled.
*/
async applyProvisioned(email: string, password: string): Promise<void> {
await this.api("PATCH", `/user/${encodeURIComponent(email)}`, { raw_password: password, enabled: true });
}
async deleteUser(email: string): Promise<void> {
await this.api("DELETE", `/user/${encodeURIComponent(email)}`);
}
/**
* Whether a mailbox exists and is enabled. Read-only, through the admin API.
*
* **The password is not checked.** Mailu authenticates in its admin service, behind the front;
* the imap server's own password database accepts any password from Mailu's subnet, so asking it
* (`doveadm auth test`) proves nothing, or refuses everyone. A lost or disabled mailbox is caught;
* a password changed by hand is not (novox/hq issue 120).
*/
async holdsUser(email: string): Promise<boolean> {
const res = await fetch(`${this.baseUrl}/user/${encodeURIComponent(email)}`, {
headers: { Authorization: this.apiKey, Accept: "application/json" },
});
if (res.status === 404) return false;
if (!res.ok) throw new Error(`Mailu API GET /user/${email}: ${res.status} ${await res.text()}`);
const user = (await res.json()) as { enabled?: boolean };
return user.enabled !== false;
}
async listAliases(): Promise<MailuAlias[]> {
const aliases = await this.api<any[]>("GET", "/alias");
return (aliases ?? []).map((a) => ({
+2 -2
View File
@@ -36,8 +36,8 @@ function watcher(created: string, deleted: string): (keys: string[]) => Promise<
};
}
const watchUsers = watcher("user.created", "user.deleted");
const watchAliases = watcher("alias.created", "alias.deleted");
const watchUsers = watcher("module.mailu.user.created", "module.mailu.user.deleted");
const watchAliases = watcher("module.mailu.alias.created", "module.mailu.alias.deleted");
async function pollUsers(): Promise<void> {
await watchUsers((await mailu.listUsers()).map((u) => u.email));
+94 -233
View File
@@ -14,53 +14,30 @@
"name": "mailu"
},
"route": {
"web": {
"label": "mail",
"endpoint": "web-tls",
"scheme": "https",
"insecure": true
},
"acme": {
"label": "mail",
"path": "/.well-known/acme-challenge",
"endpoint": "web",
"priority": 100
},
"autoconfig": {
"label": "autoconfig",
"endpoint": "autoconfig"
},
"autodiscover": {
"label": "autodiscover",
"endpoint": "autoconfig"
},
"automx": {
"label": "automx",
"endpoint": "autoconfig"
}
"label": "mail",
"port": 7080
}
},
"binds": {
"postgres-database": "${dir:state}/database.json",
"route": "${dir:state}/route.json"
"postgres-database": "/var/lib/mailu/database.json",
"route": "/var/lib/mailu/route.json"
},
"secrets": {
"postgres-database": "${dir:state}/database.secret",
"postgres-database": "/var/lib/mailu/database.secret",
"secret": {
"secret-key": "${dir:state}/secret-key.secret",
"admin": "${dir:state}/admin.secret",
"api-token": "${dir:state}/api-token.secret"
"secret-key": "/var/lib/mailu/secret-key.secret",
"admin": "/var/lib/mailu/admin.secret",
"api-token": "/var/lib/mailu/api-token.secret"
}
},
"emits": [
"user.created",
"user.deleted",
"alias.created",
"alias.deleted"
"module.mailu.user.created",
"module.mailu.user.deleted",
"module.mailu.alias.created",
"module.mailu.alias.deleted"
],
"listens": [
{
"name": "smtp",
"port": 25,
"protocol": "tcp",
"from": "anywhere",
@@ -68,23 +45,6 @@
"fixed": true
},
{
"name": "pop3",
"port": 110,
"protocol": "tcp",
"from": "anywhere",
"why": "POP3, kept at parity with the predecessor; pruning legacy protocols is its own deliberate change",
"fixed": true
},
{
"name": "imap",
"port": 143,
"protocol": "tcp",
"from": "anywhere",
"why": "IMAP with STARTTLS, kept at parity",
"fixed": true
},
{
"name": "smtps",
"port": 465,
"protocol": "tcp",
"from": "anywhere",
@@ -92,15 +52,13 @@
"fixed": true
},
{
"name": "submission",
"port": 587,
"protocol": "tcp",
"from": "anywhere",
"why": "submission; also what the smtp provision serves consumers",
"why": "submission",
"fixed": true
},
{
"name": "imaps",
"port": 993,
"protocol": "tcp",
"from": "anywhere",
@@ -108,33 +66,10 @@
"fixed": true
},
{
"name": "pop3s",
"port": 995,
"protocol": "tcp",
"from": "anywhere",
"why": "POP3 over TLS, kept at parity",
"fixed": true
},
{
"name": "web",
"port": 7080,
"protocol": "tcp",
"from": "mesh",
"why": "the web front over http; only the ACME HTTP-01 passthrough is routed here \u2014 everything else 301s to https and would loop a proxy"
},
{
"name": "web-tls",
"port": 7443,
"protocol": "tcp",
"from": "mesh",
"why": "the web front over its own TLS (admin, webmail, API); the public name mail.novox.be is a route grant reaching it here"
},
{
"name": "autoconfig",
"port": 4243,
"protocol": "tcp",
"from": "mesh",
"why": "automx: mail client autoconfiguration; autoconfig/autodiscover/automx.novox.be are route grants reaching it here"
"why": "the web interface (admin, webmail, admin API), behind the route proxy"
}
],
"own-secrets": {
@@ -150,125 +85,125 @@
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
},
{
"id": "grants",
"type": "directory",
"mode": "0700"
},
{
"id": "data-automx",
"type": "directory",
"path": "/var/lib/mailu",
"mode": "0700"
},
{
"id": "config-env",
"type": "file",
"path": "${dir:state}/mailu.env",
"path": "/var/lib/mailu/mailu.env",
"mode": "0644",
"content": "ADMIN_ADDRESS=mailu-admin\nANTISPAM_ADDRESS=mailu-antispam\nANTIVIRUS_ADDRESS=mailu-antivirus\nIMAP_ADDRESS=mailu-imap\nSMTP_ADDRESS=mailu-smtp\nFRONT_ADDRESS=mailu-front\nWEBMAIL_ADDRESS=mailu-webmail\nWEBDAV_ADDRESS=mailu-webdav\nREDIS_ADDRESS=mailu-redis\nPORTS=25,80,443,465,993,995,4190,110,143,587\nDOMAIN=novox.be\nHOSTNAMES=mail.novox.be\nPOSTMASTER=admin\nSITENAME=Novox\nWEBSITE=https://novox.be\nTLS_FLAVOR=letsencrypt\nSUBNET=192.168.203.0/24\nCOMPOSE_PROJECT_NAME=mailu\nANTIVIRUS=clamav\nWEBMAIL=roundcube\nWEBDAV=radicale\nFETCHMAIL_ENABLED=True\nFETCHMAIL_DELAY=600\nADMIN=true\nWEB_ADMIN=/admin\nWEB_WEBMAIL=/webmail\nWEBROOT_REDIRECT=/webmail\nAPI=true\nWEB_API=/api\nAUTH_RATELIMIT_IP=6000/hour\nAUTH_RATELIMIT_USER=1000/day\nCREDENTIAL_ROUNDS=12\nPASSWORD_SCHEME=PBKDF2\nDISABLE_STATISTICS=True\nMESSAGE_SIZE_LIMIT=50000000\nMESSAGE_RATELIMIT=200/day\nRECIPIENT_DELIMITER=+\nPOSTFIX_MYNETWORKS=127.0.0.0/8 [::1]/128\nRELAYNETS=\nRELAYHOST=\nREJECT_UNLISTED_RECIPIENT=\nDB_FLAVOR=postgresql\nINITIAL_ADMIN_ACCOUNT=admin\nINITIAL_ADMIN_DOMAIN=novox.be\nINITIAL_ADMIN_MODE=ifmissing\nSMTP_PORT=25\nSMTPS_PORT=465\nSUBMISSION_PORT=587\nPOP3_PORT=110\nPOP3S_PORT=995\nIMAP_PORT=143\nIMAPS_PORT=993\nHTTP_PORT=7080\nHTTPS_PORT=7443\nAUTOMX_PORT=4243\nAMX_SMTP_ADDRESS=mail.novox.be\nAMX_SMTP_PORT=587\nAMX_IMAP_ADDRESS=mail.novox.be\nAMX_IMAP_PORT=143\nAMX_MAIL_DOMAINS=novox.be\nDMARC_RUA=admin\nDMARC_RUF=admin\nLETSENCRYPT_SHORTCHAIN=True\nTZ=Etc/UTC\nLOG_LEVEL=INFO\nWELCOME=false\nREAL_IP_HEADER=X-Real-IP\nREAL_IP_FROM=142.132.152.141\nCOMPRESSION=\nCOMPRESS_LEVEL=\nCOMPRESSION_LEVEL=\nBIND_ADDRESS4=127.0.0.1\nBIND_ADDRESS6=::1\nMAILU_VERSION=1.9\nDOCKER_ORG=mailu\nDOCKER_PREFIX=\nWELCOME_SUBJECT=Welcome to your new email account\nWELCOME_BODY=Welcome to your new email account, if you can read this, then it is configured properly!\n"
"content": "DOMAIN=novox.be\nHOSTNAMES=mail.novox.be\nPOSTMASTER=admin\nSITENAME=Novox\nWEBSITE=https://novox.be\nTLS_FLAVOR=cert\nSUBNET=192.168.203.0/24\nCOMPOSE_PROJECT_NAME=mailu\nHOST_ADMIN=mailu-admin\nHOST_ANTISPAM=mailu-antispam:11332\nHOST_IMAP=mailu-imap\nHOST_SMTP=mailu-smtp\nHOST_WEBMAIL=mailu-webmail\nHOST_WEBDAV=mailu-webdav:5232\nHOST_REDIS=mailu-redis\nHOST_FRONT=mailu-front\nREDIS_ADDRESS=mailu-redis\nANTIVIRUS=clamav\nWEBMAIL=roundcube\nWEBDAV=radicale\nFETCHMAIL_ENABLED=True\nFETCHMAIL_DELAY=600\nADMIN=true\nWEB_ADMIN=/admin\nWEB_WEBMAIL=/webmail\nWEBROOT_REDIRECT=/webmail\nWEBMAIL_ADDRESS=webmail\nAPI=true\nWEB_API=/api\nAUTH_RATELIMIT_IP=6000/hour\nAUTH_RATELIMIT_USER=1000/day\nCREDENTIAL_ROUNDS=12\nPASSWORD_SCHEME=PBKDF2\nDISABLE_STATISTICS=True\nMESSAGE_SIZE_LIMIT=50000000\nMESSAGE_RATELIMIT=200/day\nRECIPIENT_DELIMITER=+\nPOSTFIX_MYNETWORKS=127.0.0.0/8 [::1]/128\nRELAYNETS=\nRELAYHOST=\nREJECT_UNLISTED_RECIPIENT=\nDB_FLAVOR=postgresql\nINITIAL_ADMIN_ACCOUNT=admin\nINITIAL_ADMIN_DOMAIN=novox.be\nINITIAL_ADMIN_MODE=ifmissing\nSMTP_PORT=25\nSMTPS_PORT=465\nSUBMISSION_PORT=587\nPOP3_PORT=110\nPOP3S_PORT=995\nIMAP_PORT=143\nIMAPS_PORT=993\nHTTP_PORT=7080\nHTTPS_PORT=7443\nAUTOMX_PORT=4243\nAMX_SMTP_ADDRESS=mail.novox.be\nAMX_SMTP_PORT=587\nAMX_IMAP_ADDRESS=mail.novox.be\nAMX_IMAP_PORT=143\nAMX_MAIL_DOMAINS=novox.be\nDMARC_RUA=admin\nDMARC_RUF=admin\nLETSENCRYPT_SHORTCHAIN=True\nTZ=Etc/UTC\nLOG_LEVEL=INFO\nWELCOME=false\nREAL_IP_HEADER=X-Real-IP\nREAL_IP_FROM=\nCOMPRESSION=\nCOMPRESS_LEVEL=\nCOMPRESSION_LEVEL=\nBIND_ADDRESS4=127.0.0.1\nBIND_ADDRESS6=::1\nMAILU_VERSION=1.9\nDOCKER_ORG=mailu\nDOCKER_PREFIX=\n"
},
{
"id": "secret-env",
"type": "file",
"path": "${dir:state}/secret.env",
"path": "/var/lib/mailu/secret.env",
"mode": "0600",
"content": "SECRET_KEY=${secret:secret-key}\n"
},
{
"id": "database-env",
"type": "file",
"path": "${dir:state}/database.env",
"path": "/var/lib/mailu/database.env",
"mode": "0600",
"content": "DB_FLAVOR=postgresql\nDB_HOST=${bound:postgres-database:at}:${bound:postgres-database:port}\nDB_USER=${bound:postgres-database:as}\nDB_NAME=${bound:postgres-database:as}\nDB_PW=${secret:postgres-database}\n"
},
{
"id": "admin-env",
"type": "file",
"path": "${dir:state}/admin.env",
"path": "/var/lib/mailu/admin.env",
"mode": "0600",
"content": "INITIAL_ADMIN_PW=${secret:admin}\nAPI_TOKEN=${secret:api-token}\n"
},
{
"id": "data-certs",
"type": "directory",
"path": "/services/mailu/data/certs",
"mode": "0700"
},
{
"id": "data-data",
"type": "directory",
"path": "/services/mailu/data/data",
"mode": "0700"
},
{
"id": "data-dkim",
"type": "directory",
"path": "/services/mailu/data/dkim",
"mode": "0700"
},
{
"id": "data-mail",
"type": "directory",
"path": "/services/mailu/data/mail",
"mode": "0700"
},
{
"id": "data-mailqueue",
"type": "directory",
"mode": "0755"
"path": "/services/mailu/data/mailqueue",
"mode": "0700"
},
{
"id": "data-filter",
"type": "directory",
"mode": "0700"
},
{
"id": "data-clamav",
"type": "directory",
"path": "/services/mailu/data/filter",
"mode": "0700"
},
{
"id": "data-redis",
"type": "directory",
"path": "/services/mailu/data/redis",
"mode": "0700"
},
{
"id": "data-webmail",
"type": "directory",
"path": "/services/mailu/data/webmail",
"mode": "0700"
},
{
"id": "data-dav",
"type": "directory",
"path": "/services/mailu/data/dav",
"mode": "0700"
},
{
"id": "data-fetchmail",
"type": "directory",
"path": "/services/mailu/data/data/fetchmail",
"mode": "0700"
},
{
"id": "data-overrides-nginx",
"type": "directory",
"path": "/services/mailu/data/overrides/nginx",
"mode": "0700"
},
{
"id": "data-overrides-dovecot",
"type": "directory",
"path": "/services/mailu/data/overrides/dovecot",
"mode": "0700"
},
{
"id": "data-overrides-postfix",
"type": "directory",
"path": "/services/mailu/data/overrides/postfix",
"mode": "0700"
},
{
"id": "data-overrides-rspamd",
"type": "directory",
"path": "/services/mailu/data/overrides/rspamd",
"mode": "0700"
},
{
"id": "data-overrides-roundcube",
"type": "directory",
"path": "/services/mailu/data/overrides/roundcube",
"mode": "0700"
},
{
@@ -280,191 +215,164 @@
"id": "resolver",
"type": "container",
"name": "mailu-resolver",
"image": "ghcr.io/mailu/unbound@sha256:3a0fdfb364a63f4f9259526e013c1ef40f5f14de3621ce1560804b3a5909584a",
"image": "ghcr.io/mailu/unbound@sha256:142aaad82ad1b0d5b59a5f1303778dba61a3e0a540f5d969c48862bcc99f6f5d",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env",
"${dir:state}/secret.env"
"/var/lib/mailu/mailu.env",
"/var/lib/mailu/secret.env"
],
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified",
"ip": "192.168.203.254"
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified"
},
{
"id": "redis",
"type": "container",
"name": "mailu-redis",
"image": "redis@sha256:4bed291aa5efb9f0d77b76ff7d4ab71eee410962965d052552db1fb80576431d",
"image": "redis@sha256:1db42ccef14898aa29bae778452d567534b59c107129cbc1163fb552de184d3c",
"network": "mailu",
"volumes": [
"${dir:data-redis}:/data"
"/services/mailu/data/redis:/data"
]
},
{
"id": "admin",
"type": "container",
"name": "mailu-admin",
"image": "ghcr.io/mailu/admin@sha256:6dbfdadc4a9590dcb7652357b505200115b689b74008653bbf369e4599a3be5a",
"image": "ghcr.io/mailu/admin@sha256:dcac20e9cbdad560faef9653b1b5ac0d9266f4098dc00f0e7f0d35f4e70ed8f1",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env",
"${dir:state}/secret.env",
"${dir:state}/database.env",
"${dir:state}/admin.env"
"/var/lib/mailu/mailu.env",
"/var/lib/mailu/secret.env",
"/var/lib/mailu/database.env",
"/var/lib/mailu/admin.env"
],
"volumes": [
"${dir:data-data}:/data",
"${dir:data-dkim}:/dkim"
"/services/mailu/data/data:/data",
"/services/mailu/data/dkim:/dkim"
],
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified",
"dns": [
"192.168.203.254"
]
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified"
},
{
"id": "imap",
"type": "container",
"name": "mailu-imap",
"image": "ghcr.io/mailu/dovecot@sha256:7f0ed5db996fbdc00adc5c5e38a08492e04f7eb4a9fbd66a03aa9a28ddf23993",
"image": "ghcr.io/mailu/dovecot@sha256:46d18ba51032be8ebd6841aa49c1ef8762c729038c5fd86a081b5b884d478af9",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env"
"/var/lib/mailu/mailu.env"
],
"volumes": [
"${dir:data-mail}:/mail",
"${dir:data-overrides-dovecot}:/overrides:ro"
],
"dns": [
"192.168.203.254"
"/services/mailu/data/mail:/mail",
"/services/mailu/data/overrides/dovecot:/overrides:ro"
]
},
{
"id": "smtp",
"type": "container",
"name": "mailu-smtp",
"image": "ghcr.io/mailu/postfix@sha256:e2e49f39e53b80eac9e7a2f18d9df11edeb4914fd62dbba89b3155e8e034f62e",
"image": "ghcr.io/mailu/postfix@sha256:bbf882880f68849511710b35237a933f3fe80c4b28bf48ff20205dbd1f1433d7",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env"
"/var/lib/mailu/mailu.env"
],
"volumes": [
"${dir:data-mailqueue}:/queue",
"${dir:data-overrides-postfix}:/overrides:ro"
],
"dns": [
"192.168.203.254"
"/services/mailu/data/mailqueue:/queue",
"/services/mailu/data/overrides/postfix:/overrides:ro"
]
},
{
"id": "antispam",
"type": "container",
"name": "mailu-antispam",
"image": "ghcr.io/mailu/rspamd@sha256:ff3666d8a61f17d309c5c6f6bcf4d40470b82299ca706ac650301175bb1a079d",
"image": "ghcr.io/mailu/rspamd@sha256:e87ab93dd252cc69499caa5317dd10d445fd4291a7ecf6bca09793c7d475a0c8",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env"
"/var/lib/mailu/mailu.env"
],
"volumes": [
"${dir:data-filter}:/var/lib/rspamd",
"${dir:data-overrides-rspamd}:/etc/rspamd/override.d:ro"
],
"dns": [
"192.168.203.254"
"/services/mailu/data/filter:/var/lib/rspamd",
"/services/mailu/data/overrides/rspamd:/etc/rspamd/override.d:ro"
]
},
{
"id": "antivirus",
"type": "container",
"name": "mailu-antivirus",
"image": "clamav/clamav-debian@sha256:b12ef8fefddbba7d88de59bea8a32622f365339154adf02d38fd089112e6745a",
"image": "ghcr.io/mailu/clamav@sha256:01d30483e4a8a20a54566addb1f9b00ebb51e8a103f9226602379c412cf5fb62",
"network": "mailu",
"volumes": [
"${dir:data-clamav}:/var/lib/clamav"
"env-file": [
"/var/lib/mailu/mailu.env",
"/var/lib/mailu/secret.env"
],
"dns": [
"192.168.203.254"
]
"volumes": [
"/services/mailu/data/filter:/data"
],
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified"
},
{
"id": "webmail",
"type": "container",
"name": "mailu-webmail",
"image": "ghcr.io/mailu/webmail@sha256:bdbee44cdb05a4658f0e3b62cc448de55ca8f8aea172279fda594826144c04f6",
"image": "ghcr.io/mailu/roundcube@sha256:19ccc9c21b2420dabb893ffa707ef90785c785e53dcb6bb9f98da01598412c43",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env",
"${dir:state}/secret.env"
"/var/lib/mailu/mailu.env",
"/var/lib/mailu/secret.env"
],
"volumes": [
"${dir:data-webmail}:/data",
"${dir:data-overrides-roundcube}:/overrides:ro"
"/services/mailu/data/webmail:/data",
"/services/mailu/data/overrides/roundcube:/overrides:ro"
],
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified",
"dns": [
"192.168.203.254"
]
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified"
},
{
"id": "webdav",
"type": "container",
"name": "mailu-webdav",
"image": "ghcr.io/mailu/radicale@sha256:690ed6edf189dfef100a5a8b37c195ebf5d9241ac5f23f2f44b8b7b75726e3de",
"image": "ghcr.io/mailu/radicale@sha256:e13cbad3791c0a6841b5d387e57e49a117808dcef87b8c9969f671ae9c3b67c0",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env",
"${dir:state}/secret.env"
"/var/lib/mailu/mailu.env",
"/var/lib/mailu/secret.env"
],
"volumes": [
"${dir:data-dav}:/data"
"/services/mailu/data/dav:/data"
],
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified",
"dns": [
"192.168.203.254"
]
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified"
},
{
"id": "fetchmail",
"type": "container",
"name": "mailu-fetchmail",
"image": "ghcr.io/mailu/fetchmail@sha256:f881c8412d3bbe73d638469b48321558d6403a9d45bfa043c1e52c752103d42d",
"image": "ghcr.io/mailu/fetchmail@sha256:7dcd1392882925d612ab2d0230d437f0c660989d572283c48b0d0f2d491adce7",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env",
"${dir:state}/secret.env"
"/var/lib/mailu/mailu.env",
"/var/lib/mailu/secret.env"
],
"volumes": [
"${dir:data-fetchmail}:/data"
"/services/mailu/data/data/fetchmail:/data"
],
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified",
"dns": [
"192.168.203.254"
]
"secrets-in-environment": "mailu-admin honours SECRET_KEY_FILE, DB_PW_FILE and API_TOKEN_FILE (configuration.py) but INITIAL_ADMIN_PW is env-only (start.py); the remaining containers' need for SECRET_KEY is unverified"
},
{
"id": "front",
"type": "container",
"name": "mailu-front",
"image": "ghcr.io/mailu/nginx@sha256:36f98897cd1bc9d27628bbb4e04bdf60147af2ec7507d6da77f002c4f256896d",
"image": "ghcr.io/mailu/nginx@sha256:09f28ab6d36367fcacc7994f7021f132ac845bdc05f04bf80906102d11aaa057",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env"
"/var/lib/mailu/mailu.env"
],
"ports": [
"25",
"110",
"143",
"465",
"587",
"993",
"995",
"7080:80",
"7443:443"
"7080:80"
],
"volumes": [
"${dir:data-certs}:/certs",
"${dir:data-overrides-nginx}:/overrides:ro"
],
"dns": [
"192.168.203.254"
"/services/mailu/data/certs:/certs",
"/services/mailu/data/overrides/nginx:/overrides:ro"
]
},
{
@@ -482,40 +390,21 @@
"network": "mailu",
"volumes": [
"/var/lib/mesh/mailu/broker:/run/secrets/broker:ro",
"${dir:state}/api-token.secret:/run/secrets/api-token:ro",
"${dir:grants}:${dir:grants}:ro",
"/var/lib/mailu/api-token.secret:/run/secrets/api-token:ro",
"/var/lib/mesh/mailu/config.json:/run/config/config.json:ro",
"/var/run/docker.sock:/var/run/docker.sock"
],
"env": {
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_MAILU_URL": "http://mailu-admin:8080/api/v1",
"MESH_MAILU_URL": "http://mailu-admin/api/v1",
"MESH_MAILU_API_KEY_FILE": "/run/secrets/api-token",
"MESH_MAILU_IMAP_CONTAINER": "mailu-imap",
"MESH_MAILU_CONFIG_FILE": "/run/config/config.json",
"MESH_MAILU_DOMAIN": "novox.be",
"MESH_RECEIVES": "${dir:grants}/mesh.json"
"MESH_MAILU_CONFIG_FILE": "/run/config/config.json"
},
"restart-on": [
"runtime-config"
],
"artifact": "runtime"
},
{
"id": "automx",
"type": "container",
"name": "mailu-automx",
"artifact": "automx",
"network": "mailu",
"env-file": [
"${dir:state}/mailu.env"
],
"ports": [
"4243"
],
"volumes": [
"${dir:data-automx}:/data"
]
}
],
"build": {
@@ -529,10 +418,6 @@
"arg": "RUNTIME_BASE",
"module": "mesh-tools",
"artifact": "runtime"
},
{
"arg": "PYTHON_BASE",
"image": "python@sha256:25f3cfeaceca14921366af4d1240b56457ef46273bdb508c7b0e8f469f6fd228"
}
],
"artifacts": [
@@ -540,31 +425,7 @@
"name": "runtime",
"kind": "image",
"from": "Dockerfile"
},
{
"name": "automx",
"kind": "image",
"from": "automx/Dockerfile"
}
]
},
"provides": [
{
"name": "smtp",
"scope": "mesh"
}
],
"serves": {
"smtp": {
"port": 587,
"domain": "novox.be",
"name": "mail.novox.be"
}
},
"receives": {
"smtp": "${dir:grants}/mesh.json"
},
"grants": {
"smtp": "${dir:grants}"
}
}
+1 -1
View File
@@ -5,7 +5,7 @@
"type": "module",
"private": true,
"dependencies": {
"@novox/mesh-sdk": "^0.1.1"
"@novox/mesh-sdk": "^0.1.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
-70
View File
@@ -1,70 +0,0 @@
// mailu's provisioner — the adapter that makes mailu a provider of the mesh `smtp` interface.
// The reconcile loop, the contributions file, and reading the mesh's minted password are the sdk
// harness's; this writes only the per-service half: how mailu creates and removes a consumer's
// sending account (novox/hq ADR 0048/0076, gitea's package-registry provisioner is the sibling).
//
// The `smtp` interface: a consumer authenticates to submission (port 587, STARTTLS) as a real
// mailbox this provisioner creates. The address is `<account>@<domain>`: the local part is the
// consumer's `account` contribution — the name it wants to send as — falling back to the mesh's
// own login for a consumer that named none; the domain is the mail server's, which is this
// module's fact, not the consumer's.
//
// **The password is the mesh's, not the provisioner's (ADR 0048).** The mesh mints it and hands
// it to both ends; mailu sets exactly that password every run — so a rotation takes — and seals
// nothing: the consumer already has its copy through the mesh's own channel.
import { runProvisioner, type Provision } from "@novox/mesh-sdk/provisioner";
import { MailuClient } from "../client.js";
const mailu = MailuClient.fromEnv();
// The mail server's own domain. From the environment the manifest composes, because the client's
// config file carries the admin API's coordinates, not the mail domain.
function domain(): string {
const named = (process.env.MESH_MAILU_DOMAIN ?? "").trim();
if (named === "") {
throw new Error("MESH_MAILU_DOMAIN is not set, so a consumer's address cannot be composed");
}
return named;
}
// The address one consumer sends as. The local part is refused rather than sanitised when it is
// not a plain mailbox name — a rewritten name is an address nobody asked for.
function addressOf(p: { as: string; values?: Readonly<Record<string, unknown>> }): string {
const contributed = typeof p.values?.["account"] === "string" ? (p.values["account"] as string).trim() : "";
const local = contributed !== "" ? contributed : p.as;
if (!/^[a-z0-9][a-z0-9._-]*$/.test(local)) {
throw new Error(`${JSON.stringify(local)} is not a usable mailbox name`);
}
return `${local}@${domain()}`;
}
runProvisioner("smtp", {
async create(p: Provision): Promise<void> {
const email = addressOf(p);
// Create if absent, and set exactly the minted password either way so a rotation takes.
// Mailu's create refuses a duplicate address, which is the signal to fall through to the
// password set — the same found-then-apply shape gitea's ensureUser settled on.
try {
await mailu.createUser(email, p.password);
} catch {
await mailu.applyProvisioned(email, p.password);
}
},
async remove(p: { as: string }): Promise<void> {
// The withdrawal only knows the mesh login, never the contributed local part — so accounts
// that contributed one are removed when the address matching the login is absent? No: the
// harness hands remove only `as`, and an address composed from a contribution cannot be
// recomputed from it. The account is therefore removed by its login-shaped address when one
// exists, and left otherwise — a mailbox holding mail is the one thing a background loop
// must not guess about (this module's own events file says the same). Withdrawal of a
// named-account consumer is an operator action until the harness carries values here.
await mailu.deleteUser(`${p.as}@${domain()}`).catch(() => {});
},
// Asked every minute by the harness: whether the backend still holds this consumer exactly as
// the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120).
async holds(p: Provision): Promise<boolean> {
return mailu.holdsUser(addressOf(p));
},
});
-1
View File
@@ -6,7 +6,6 @@
],
"listens": [
{
"name": "api",
"port": 59125,
"protocol": "tcp",
"from": "mesh",
+1 -6
View File
@@ -22,7 +22,7 @@ COPY . .
# The compiler is invoked by its real path rather than through node_modules/.bin, whose entries are
# symlinks to a launcher that requires its library relatively — resolved away when the base image
# was assembled.
RUN node /app/node_modules/typescript/bin/tsc pg.d.ts store.ts index.ts tools/index.ts prepare/index.ts \
RUN node /app/node_modules/typescript/bin/tsc pg.d.ts store.ts index.ts tools/index.ts \
--module NodeNext --moduleResolution NodeNext --target ES2022 --outDir dist
# **A module may need something the base image does not carry.** The base holds what every module
@@ -48,8 +48,3 @@ COPY --from=build /deps/node_modules /app/modules/mesh-catalog/node_modules
# to listen for what the builder announces. Serve binds the broker first, then imports these, so
# `on()` has something to subscribe to.
ENV MESH_TOOL_MODULES=/app/modules/mesh-catalog/dist/index.js,/app/modules/mesh-catalog/dist/tools/index.js
# And what prepares this module's state, for the runtime's `prepare` mode (novox/hq ADR 0135). Named
# here, beside the entrypoints above, because the module knows which of its files prepares its state
# and nothing else could: the mesh asks one word and this says what answers it.
ENV MESH_PREPARE=/app/modules/mesh-catalog/dist/prepare/index.js
+11 -25
View File
@@ -1,6 +1,6 @@
// mesh-catalog's entrypoint — the module graph's consumer (novox/hq ADR 0070, ADR 0072).
//
// The build-machine role announces what it built; this places it in the graph and announces what that means.
// The builder announces what it built; this places it in the graph and announces what that means.
// The control plane hooks the *meaning* — a module was upgraded — rather than the build output, so
// it never has to interpret an artifact or ask this module anything.
//
@@ -14,12 +14,10 @@ import { Graph, type Made } from "./store.js";
const graph = Graph.fromEnv();
// The schema is not brought up here. The mesh prepares this module's state before it starts this
// version, and does not start it if that failed (novox/hq ADR 0135) — see prepare/index.ts. Doing it
// at start made a schema that could not be reached a crash loop instead of a stop, with the graph
// keeping a gap and nothing saying so. The reason it used to be here — that a step blocking the apply
// would block the very apply that brings the overlay up — stopped being true when a step's failure
// became this module's business and not the machine's (ADR 0136).
// Before subscribing, and idempotent. The runtime is restarted until its store is reachable, which
// is the same arrangement model-usage uses: a schema step that had to reach the provider over the
// overlay would block the very apply that brings the overlay up.
await graph.migrate();
/** What the builder says when it has built something. */
interface Built {
@@ -49,15 +47,7 @@ interface Built {
replay?: boolean;
}
/**
* What a build means for the graph, wherever it came from.
*
* Two emitters say the same thing and neither is a mistake: the build machine says it as it happens,
* and the control plane says what it already held when this module asks what it missed
* (novox/hq ADR 0134). A replay is marked as one in its body, so nothing acts on a module that moved
* months ago — see `replay` above.
*/
const placeTheBuild = async (event: { body: unknown }): Promise<void> => {
await on("module.builder.built", async (event) => {
const body = event.body as Built;
if (!body.module || !body.commit) {
// Said rather than dropped: a build that announced itself without saying what it built is a
@@ -79,7 +69,7 @@ const placeTheBuild = async (event: { body: unknown }): Promise<void> => {
// it was missing, and the mesh is told nothing happened, because nothing did.
if (body.replay) return;
await emit("registered", {
await emit("module.mesh-catalog.registered", {
module: body.module, commit: body.commit, upgraded,
});
@@ -87,23 +77,19 @@ const placeTheBuild = async (event: { body: unknown }): Promise<void> => {
// through modules that did not change, forever (ADR 0072).
if (!upgraded) return;
await emit("upgraded", {
await emit("module.mesh-catalog.upgraded", {
module: body.module, commit: body.commit, previous,
});
// What can be built now — stale, and waiting on nothing that is itself stale.
for (const next of await graph.buildable()) {
await emit("rebuild-needed", {
await emit("module.mesh-catalog.rebuild-needed", {
module: next.module,
builtAt: next.commit,
because: next.because,
});
}
};
// As it happens, and what the mesh already held when this module asked what it missed.
await on("mesh-build-machine.built", placeTheBuild);
await on("mesh-controller.built-before", placeTheBuild);
});
// **And ask for what was built before this catalogue existed** (novox/hq 04-ISSUES/050).
//
@@ -115,4 +101,4 @@ await on("mesh-controller.built-before", placeTheBuild);
// Asked on every start, not only the first. A catalogue cannot tell whether it has a gap, and the
// answer is idempotent: registering a build already held changes nothing and announces nothing.
// Asked AFTER subscribing, so a build arriving during the replay is not lost between the two.
await emit("catching-up", {});
await emit("module.mesh-catalog.catching-up", {});
+5 -8
View File
@@ -7,7 +7,7 @@
],
"claims": [
{
"name": "mesh-catalog",
"name": "the-catalogue",
"scope": "mesh"
}
],
@@ -29,16 +29,13 @@
"broker": "/var/lib/mesh/mesh-catalog/broker"
},
"consumes": [
"mesh-build-machine.built",
"mesh-controller.built-before"
"module.builder.built"
],
"emits": [
"registered",
"upgraded",
"rebuild-needed",
"catching-up"
"module.mesh-catalog.registered",
"module.mesh-catalog.upgraded",
"module.mesh-catalog.rebuild-needed"
],
"prepares": true,
"resources": [
{
"id": "mesh-state",
-17
View File
@@ -1,17 +0,0 @@
// The catalogue's state, brought to the shape this version needs (novox/hq ADR 0135).
//
// **The mesh runs this before the version that needs it, and does not start that version if it
// fails** — and the refusal reaches this module and nothing else on the machine
// (novox/hq ADR 0136). That is the whole difference from where this used to happen: at start, inside
// the runtime, a schema that could not be brought up was a crash loop, the graph kept a gap, and
// nothing anywhere said so.
//
// Nothing here connects to the broker. Preparation runs before the version that would use it, so
// there is nothing yet to talk to; the runtime's `prepare` mode imports this and awaits it, and this
// process exiting non-zero is how the host knows not to start the runtime.
import { Graph } from "../store.js";
const graph = Graph.fromEnv();
await graph.migrate();
console.log("[mesh-catalog] the module graph's schema is what this version needs");
await graph.close();
+1 -2
View File
@@ -12,7 +12,6 @@
"pg.d.ts",
"store.ts",
"index.ts",
"tools/index.ts",
"prepare/index.ts"
"tools/index.ts"
]
}
+3 -3
View File
@@ -16,15 +16,15 @@ interface SecretEvent {
rotations?: number;
}
await on<SecretEvent>("secret.provisioned", async (e) => {
await on<SecretEvent>("module.mesh-vault.secret.provisioned", async (e) => {
console.log(`[mesh-vault] secret provisioned for ${e.body.as} on ${e.body.consumer} (${e.body.fingerprint})`);
});
await on<SecretEvent>("secret.rotated", async (e) => {
await on<SecretEvent>("module.mesh-vault.secret.rotated", async (e) => {
console.log(`[mesh-vault] secret rotated for ${e.body.as} — rotation ${e.body.rotations} (${e.body.fingerprint})`);
});
await on<SecretEvent>("secret.deprovisioned", async (e) => {
await on<SecretEvent>("module.mesh-vault.secret.deprovisioned", async (e) => {
console.log(`[mesh-vault] secret withdrawn from ${e.body.as}`);
});
+6 -6
View File
@@ -11,14 +11,14 @@
"container-runtime"
],
"emits": [
"secret.provisioned",
"secret.rotated",
"secret.deprovisioned"
"module.mesh-vault.secret.provisioned",
"module.mesh-vault.secret.rotated",
"module.mesh-vault.secret.deprovisioned"
],
"consumes": [
"mesh-vault.secret.provisioned",
"mesh-vault.secret.rotated",
"mesh-vault.secret.deprovisioned"
"module.mesh-vault.secret.provisioned",
"module.mesh-vault.secret.rotated",
"module.mesh-vault.secret.deprovisioned"
],
"receives": {
"secret": "/var/lib/mesh-vault/grants/mesh.json"
+1 -1
View File
@@ -43,6 +43,6 @@ runProvisioner("secret", {
async remove(p: { as: string }): Promise<void> {
if (!ledger.withdraw(p.as)) return;
console.log(`[mesh-vault] withdrawn: ${p.as}`);
await announce("secret.deprovisioned", { as: p.as });
await announce("module.mesh-vault.secret.deprovisioned", { as: p.as });
},
});
+4 -17
View File
@@ -125,18 +125,6 @@ export class MinioClient {
throw new Error(`minio bucketExists ${bucket}: ${status}`);
}
/**
* Whether a consumer's access key, with exactly this secret, reaches its bucket: a HEAD of the
* bucket signed as the consumer, the way it signs. Read-only. `false` when the key is unknown, the
* secret wrong, access denied or the bucket gone; any other answer rejects (novox/hq issue 120).
*/
async canReachAs(bucket: string, accessKey: string, secretKey: string): Promise<boolean> {
const { status } = await this.request("HEAD", `/${bucket}`, {}, { accessKey, secretKey });
if (status === 200) return true;
if (status === 403 || status === 404) return false;
throw new Error(`minio HEAD ${bucket} as ${accessKey}: ${status}`);
}
async createBucket(bucket: string): Promise<void> {
const { status, text } = await this.request("PUT", `/${bucket}`);
// 200 created; 409 BucketAlreadyOwnedByYou — idempotent, a re-provision must not fail.
@@ -263,7 +251,6 @@ export class MinioClient {
method: string,
path: string,
query: Record<string, string> = {},
as: { accessKey: string; secretKey: string } = { accessKey: this.rootUser, secretKey: this.rootPassword },
): Promise<{ status: number; headers: Headers; text: string }> {
const { amzDate, dateStamp } = this.stamp();
const host = new URL(this.baseUrl).host;
@@ -275,8 +262,8 @@ export class MinioClient {
const canonicalRequest = [method, encodedPath, canonicalQuery, canonicalHeaders, signedHeaders, payloadHash].join("\n");
const scope = `${dateStamp}/${this.region}/s3/aws4_request`;
const stringToSign = ["AWS4-HMAC-SHA256", amzDate, scope, sha256hex(canonicalRequest)].join("\n");
const signature = hmac(this.signingKey(dateStamp, as.secretKey), stringToSign).toString("hex");
const authorization = `AWS4-HMAC-SHA256 Credential=${as.accessKey}/${scope}, SignedHeaders=${signedHeaders}, Signature=${signature}`;
const signature = hmac(this.signingKey(dateStamp), stringToSign).toString("hex");
const authorization = `AWS4-HMAC-SHA256 Credential=${this.rootUser}/${scope}, SignedHeaders=${signedHeaders}, Signature=${signature}`;
const url = `${this.baseUrl}${encodedPath}${canonicalQuery ? `?${canonicalQuery}` : ""}`;
const res = await fetch(url, {
@@ -288,8 +275,8 @@ export class MinioClient {
return { status: res.status, headers: res.headers, text };
}
private signingKey(dateStamp: string, secretKey: string = this.rootPassword): Buffer {
const kDate = hmac(`AWS4${secretKey}`, dateStamp);
private signingKey(dateStamp: string): Buffer {
const kDate = hmac(`AWS4${this.rootPassword}`, dateStamp);
const kRegion = hmac(kDate, this.region);
const kService = hmac(kRegion, "s3");
return hmac(kService, "aws4_request");
+14 -66
View File
@@ -7,48 +7,25 @@
"scope": "mesh"
}
],
"requires": [
"route"
],
"contributes": {
"route": {
"api": {
"label": "files-api",
"endpoint": "s3"
},
"console": {
"label": "files",
"endpoint": "console"
}
}
},
"capabilities": [
"container-runtime"
],
"emits": [
"bucket.created",
"bucket.removed"
"module.minio.bucket.created",
"module.minio.bucket.removed"
],
"listens": [
{
"name": "s3",
"port": 9000,
"protocol": "tcp",
"from": "mesh",
"why": "the S3 endpoint"
},
{
"name": "console",
"port": 9001,
"protocol": "tcp",
"from": "mesh",
"why": "the admin console"
}
],
"serves": {
"s3-bucket": {
"scheme": "http",
"region": "eu-west",
"region": "us-east-1",
"port": 9000
}
},
@@ -91,20 +68,20 @@
{
"id": "data",
"type": "directory",
"path": "/var/lib/minio-store",
"path": "/services/minio/data/data1-1",
"mode": "0700"
},
{
"id": "net",
"type": "network",
"name": "minio-net"
"name": "minio"
},
{
"id": "server",
"type": "container",
"name": "minio",
"image": "docker.io/pgsty/minio@sha256:b6bfe7239bfc83fb90d31612d9704d86039dd714f7904b3f1ad68f211e602372",
"network": "minio-net",
"image": "quay.io/minio/minio@sha256:14cea493d9a34af32f524e538b8346cf79f3321eff8e708c1e2960462bd8936e",
"network": "minio",
"args": [
"server",
"/data",
@@ -115,24 +92,22 @@
"/var/lib/minio/root.env"
],
"ports": [
"9000",
"9001"
"9000"
],
"volumes": [
"/var/lib/minio-store:/data",
"/services/minio/data/data1-1:/data",
"/var/lib/minio/root.secret:/run/secrets/root:ro"
],
"env": {
"MINIO_ROOT_PASSWORD_FILE": "/run/secrets/root",
"MINIO_BROWSER_REDIRECT_URL": "https://files.novox.be",
"MINIO_REGION": "eu-west"
"MINIO_ROOT_PASSWORD_FILE": "/run/secrets/root"
}
},
{
"id": "runtime",
"type": "container",
"name": "mesh-minio",
"network": "minio-net",
"image": "mesh-runtime-minio@sha256:0000000000000000000000000000000000000000000000000000000000000000",
"network": "minio",
"volumes": [
"/var/lib/mesh/minio/broker:/run/secrets/broker:ro",
"/var/lib/minio/grants:/var/lib/minio/grants:ro",
@@ -142,36 +117,9 @@
"MESH_MINIO_ENDPOINT": "http://minio:9000",
"MESH_MINIO_ROOT_USER": "meshroot",
"MESH_MINIO_ROOT_PASSWORD_FILE": "/run/secrets/root",
"MESH_MINIO_REGION": "eu-west",
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_RECEIVES": "/var/lib/minio/grants/mesh.json"
},
"artifact": "runtime"
}
}
],
"build": {
"on": [
{
"arg": "BUILD_BASE",
"module": "mesh-tools",
"artifact": "build"
},
{
"arg": "RUNTIME_BASE",
"module": "mesh-tools",
"artifact": "runtime"
},
{
"arg": "MC_CLI",
"image": "docker.io/pgsty/minio@sha256:b6bfe7239bfc83fb90d31612d9704d86039dd714f7904b3f1ad68f211e602372"
}
],
"artifacts": [
{
"name": "runtime",
"kind": "image",
"from": "Dockerfile"
}
]
}
]
}
+1 -1
View File
@@ -5,7 +5,7 @@
"type": "module",
"private": true,
"dependencies": {
"@novox/mesh-sdk": "^0.1.1"
"@novox/mesh-sdk": "^0.1.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
+2 -8
View File
@@ -30,7 +30,7 @@ runProvisioner("s3-bucket", {
try { await minio.removeAccessKey(accessKeyId); } catch { /* none yet — first provision */ }
await minio.createAccessKey(bucket, accessKeyId, p.password);
await announce("bucket.created", {
await announce("module.minio.bucket.created", {
bucket,
consumer: p.consumer ?? "",
accessKey: accessKeyId,
@@ -51,13 +51,7 @@ runProvisioner("s3-bucket", {
console.error(`[minio] bucket ${bucket} not removed (likely non-empty), access revoked: ${err}`);
}
await announce("bucket.removed", { bucket, accessKey: p.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).
async holds(p: Provision): Promise<boolean> {
return minio.canReachAs(bucketFor(p.as), p.as, p.password);
await announce("module.minio.bucket.removed", { bucket, accessKey: p.as });
},
});
+1 -1
View File
@@ -22,7 +22,7 @@ const store = UsageStore.fromEnv();
// to reach the provider over the overlay would block the very apply that brings the overlay up.
await store.migrate();
await on("*.usage.*", async (event) => {
await on("module.*.usage.*", async (event) => {
const body = event.body as { rows?: UsageRow[]; raw?: unknown };
for (const row of body.rows ?? []) {
try {
+1 -1
View File
@@ -20,7 +20,7 @@
"postgres-database": "/var/lib/model-usage/database.secret"
},
"consumes": [
"*.usage.*"
"module.*.usage.*"
],
"own-secrets": {
"broker": "/var/lib/mesh/model-usage/broker"
-28
View File
@@ -109,34 +109,6 @@ print(EJSON.stringify({ ok: 1 }));
await this.evalJs<{ ok: number }>(js);
}
/**
* Whether `user` authenticates against `database` with exactly `password` and holds `dbOwner`
* there: checked by connecting as the consumer, the way it connects. Read-only. `false` only on an
* authentication failure or a missing role; an unreachable server rejects (novox/hq issue 120).
*/
async canAuthenticateAs(database: string, user: string, password: string): Promise<boolean> {
// Connected without credentials, then authenticated inside the eval from the environment, so
// the consumer's password is neither on argv nor in the message of a failed command.
const uri = `mongodb://${this.conn.host}:${this.conn.port}/?serverSelectionTimeoutMS=10000`;
const js =
"const t = db.getSiblingDB(process.env.MESH_HOLDS_DB);" +
"t.auth(process.env.MESH_HOLDS_USER, process.env.MESH_HOLDS_PW);" +
"print(EJSON.stringify(t.runCommand({ connectionStatus: 1 }).authInfo.authenticatedUserRoles))";
let stdout: string;
try {
({ stdout } = await run("mongosh", [uri, "--quiet", "--eval", js], {
env: { ...process.env, MESH_HOLDS_DB: database, MESH_HOLDS_USER: user, MESH_HOLDS_PW: password },
timeout: 30_000,
}));
} catch (err) {
const text = `${(err as { stderr?: string }).stderr ?? ""}${(err as { stdout?: string }).stdout ?? ""}`;
if (/Authentication failed|AuthenticationFailed/i.test(text)) return false;
throw new Error(`mongosh could not check ${user}: ${text.trim().slice(0, 500) || String((err as Error).message).split("\n")[0]}`);
}
const roles = JSON.parse(stdout.trim()) as { role: string; db: string }[];
return roles.some((r) => r.role === "dbOwner" && r.db === database);
}
/** Drop a database and its owning user, idempotently. Dropping the database evicts its data; the
* user is removed first so a re-grant of the same login starts clean. */
async dropDatabaseAndUser(database: string, user: string): Promise<void> {
+2 -2
View File
@@ -14,11 +14,11 @@ interface DatabaseEvent {
user?: string;
}
await on<DatabaseEvent>("database.provisioned", async (e) => {
await on<DatabaseEvent>("module.mongodb.database.provisioned", async (e) => {
console.log(`[mongodb] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
});
await on<DatabaseEvent>("database.deprovisioned", async (e) => {
await on<DatabaseEvent>("module.mongodb.database.deprovisioned", async (e) => {
console.log(`[mongodb] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
});
+18 -17
View File
@@ -11,16 +11,15 @@
"container-runtime"
],
"emits": [
"database.provisioned",
"database.deprovisioned"
"module.mongodb.database.provisioned",
"module.mongodb.database.deprovisioned"
],
"consumes": [
"mongodb.database.provisioned",
"mongodb.database.deprovisioned"
"module.mongodb.database.provisioned",
"module.mongodb.database.deprovisioned"
],
"listens": [
{
"name": "database",
"port": 27017,
"protocol": "tcp",
"from": "mesh",
@@ -33,13 +32,13 @@
}
},
"receives": {
"mongodb-database": "${dir:grants}/mesh.json"
"mongodb-database": "/var/lib/mongodb/grants/mesh.json"
},
"grants": {
"mongodb-database": "${dir:grants}"
"mongodb-database": "/var/lib/mongodb/grants"
},
"own-secrets": {
"root": "${dir:state}/root.secret",
"root": "/var/lib/mongodb/root.secret",
"broker": "/var/lib/mesh/mongodb/broker"
},
"secrets-owner": "999:999",
@@ -53,17 +52,19 @@
{
"id": "state",
"type": "directory",
"mode": "0700",
"place": "."
"path": "/var/lib/mongodb",
"mode": "0700"
},
{
"id": "grants",
"type": "directory",
"path": "/var/lib/mongodb/grants",
"mode": "0700"
},
{
"id": "data",
"type": "directory",
"path": "/services/mongodb/db-data",
"mode": "0700"
},
{
@@ -74,7 +75,7 @@
{
"id": "server",
"type": "container",
"name": "mongodb-server",
"name": "mongo",
"image": "mongo@sha256:e3fa459b4f4b72f3257c67a23c145e250b8b5700f033860392c68539b998bbe3",
"network": "mongodb",
"env": {
@@ -85,8 +86,8 @@
"27017"
],
"volumes": [
"${dir:data}:/data/db",
"${dir:state}/root.secret:/run/secrets/root:ro"
"/services/mongodb/db-data:/data/db",
"/var/lib/mongodb/root.secret:/run/secrets/root:ro"
]
},
{
@@ -96,14 +97,14 @@
"network": "mongodb",
"volumes": [
"/var/lib/mesh/mongodb/broker:/run/secrets/broker:ro",
"${dir:grants}:${dir:grants}:ro",
"${dir:state}/root.secret:/run/secrets/root:ro"
"/var/lib/mongodb/grants:/var/lib/mongodb/grants:ro",
"/var/lib/mongodb/root.secret:/run/secrets/root:ro"
],
"env": {
"MESH_PROVISION_MONGODB": "mongodb://root@mongodb-server:27017/admin?authSource=admin",
"MESH_PROVISION_MONGODB": "mongodb://root@mongo:27017/admin?authSource=admin",
"MESH_PROVISION_PASSWORD_FILE": "/run/secrets/root",
"MESH_BROKER_FILE": "/run/secrets/broker",
"MESH_RECEIVES": "${dir:grants}/mesh.json"
"MESH_RECEIVES": "/var/lib/mongodb/grants/mesh.json"
},
"artifact": "runtime"
}
+1 -1
View File
@@ -5,7 +5,7 @@
"type": "module",
"private": true,
"dependencies": {
"@novox/mesh-sdk": "^0.1.1"
"@novox/mesh-sdk": "^0.1.0"
},
"devDependencies": {
"@types/node": "^22.0.0",
+2 -7
View File
@@ -34,7 +34,7 @@ runProvisioner("mongodb-database", {
// Database and owning user share the consumer's login, so the consumer owns exactly its own.
const database = p.as;
await mongo.createDatabaseAndUser(database, p.as, p.password);
await announce("database.provisioned", {
await announce("module.mongodb.database.provisioned", {
consumer: p.consumer ?? "",
database,
user: p.as,
@@ -43,11 +43,6 @@ runProvisioner("mongodb-database", {
async remove(p: { as: string }): Promise<void> {
await mongo.dropDatabaseAndUser(p.as, p.as);
await announce("database.deprovisioned", { database: p.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).
async holds(p: Provision): Promise<boolean> {
return mongo.canAuthenticateAs(p.as, p.as, p.password);
await announce("module.mongodb.database.deprovisioned", { database: p.as });
},
});
+3 -94
View File
@@ -15,7 +15,6 @@
// The one cost dynsec carries is the bootstrap file; see initBootstrapFile() and the module README.
import { randomBytes } from "node:crypto";
import { connect as tcpConnect } from "node:net";
import { readFileSync } from "node:fs";
import { execFile } from "node:child_process";
import { promisify } from "node:util";
@@ -88,20 +87,9 @@ export class MosquittoClient {
"-u", this.conn.adminUser,
"-P", this.conn.adminPassword,
];
let stdout: string;
let stderr: string;
try {
({ stdout, stderr } = await run("mosquitto_ctrl", [...base, "dynsec", ...args], {
maxBuffer: 16 << 20,
timeout: 30_000,
}));
} catch (err) {
// A failed run's message repeats its argv, the admin password (-P) included; say what failed
// without it.
const e = err as { code?: unknown; signal?: unknown; stderr?: string; stdout?: string };
const detail = `${e.stderr ?? ""}${e.stdout ?? ""}`.trim().slice(0, 500);
throw new Error(`mosquitto_ctrl dynsec ${args[0] ?? ""} could not run (${e.code ?? e.signal ?? "error"}): ${detail}`);
}
const { stdout, stderr } = await run("mosquitto_ctrl", [...base, "dynsec", ...args], {
maxBuffer: 16 << 20,
});
const failure = ctlError(`${stdout}\n${stderr}`);
if (failure) {
throw new Error(`mosquitto_ctrl dynsec ${args[0] ?? ""} failed: ${failure}`);
@@ -152,11 +140,6 @@ export class MosquittoClient {
if (await this.clientExists(username)) {
await this.ctl("setClientPassword", username, password);
// A disabled client is refused like a wrong password, so the check the provisioner runs
// reports it lost; applying again must enable it, or the two would disagree for ever.
if (/Disabled:\s*true/i.test(await this.ctl("getClient", username))) {
await this.ctl("enableClient", username);
}
} else {
await this.ctl("createClient", username, "-p", password);
}
@@ -180,28 +163,6 @@ export class MosquittoClient {
}
}
/**
* Whether a consumer's client accepts exactly this password and still carries its own role.
* Read-only. The password is checked the way the consumer is checked, by an MQTT CONNECT as it,
* and the broker's CONNACK code is the answer: 0 accepted, 4 bad credentials, 5 not authorised.
* Nothing rides on argv. An unreachable broker rejects (novox/hq issue 120).
*/
async holdsClient(username: string, password: string): Promise<boolean> {
const code = await mqttConnack(this.conn.host, this.conn.port, username, password);
if (code === 4 || code === 5) return false;
if (code !== 0) throw new Error(`mosquitto refused ${username} with CONNACK ${code}`);
// The role, asked directly: only "not found" means absent. Any other failure to ask rejects,
// unlike clientHasRole, which reads every failure as "no role".
let out: string;
try {
out = await this.ctl("getClient", username);
} catch (err) {
if (/not\s*found|does not exist|no such/i.test(String(err))) return false;
throw err;
}
return new RegExp(`(^|\\s)${escapeRegExp(username)}\\s+\\(priority`, "m").test(out);
}
/** Remove a client and the per-client role created for it, idempotently. */
async deleteScopedClient(username: string): Promise<void> {
await ignoreMissing(this.ctl("deleteClient", username));
@@ -288,55 +249,3 @@ function readSecretFile(path: string | undefined): string | undefined {
return undefined;
}
}
/**
* Connect once over MQTT 3.1.1 with a username and password, return the broker's CONNACK return code,
* and disconnect. A clean session under a throwaway client id, so no consumer session is taken over.
*/
function mqttConnack(host: string, port: number, username: string, password: string): Promise<number> {
const str = (v: string): Buffer => {
const b = Buffer.from(v, "utf8");
const len = Buffer.alloc(2);
len.writeUInt16BE(b.length);
return Buffer.concat([len, b]);
};
const variable = Buffer.concat([str("MQTT"), Buffer.from([4, 0xc2, 0, 10])]); // level 4; user+pass+clean; keepalive 10s
const payload = Buffer.concat([str(`mesh-holds-${randomBytes(6).toString("hex")}`), str(username), str(password)]);
let remaining = variable.length + payload.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);
const packet = Buffer.concat([Buffer.from([0x10, ...lenBytes]), variable, payload]);
return new Promise((resolve, reject) => {
const socket = tcpConnect({ host, port });
let buf = Buffer.alloc(0);
const timer = setTimeout(() => {
socket.destroy();
reject(new Error(`no CONNACK from ${host}:${port} within 10s`));
}, 10_000);
socket.on("connect", () => socket.write(packet));
socket.on("data", (chunk) => {
buf = Buffer.concat([buf, chunk]);
if (buf.length < 4) return;
clearTimeout(timer);
if (buf[0] !== 0x20) {
socket.destroy();
reject(new Error(`unexpected MQTT packet 0x${buf[0].toString(16)} instead of CONNACK`));
return;
}
const code = buf[3];
if (code === 0) socket.end(Buffer.from([0xe0, 0])); // DISCONNECT
else socket.destroy();
resolve(code);
});
socket.on("error", (err) => {
clearTimeout(timer);
reject(err);
});
});
}
+2 -2
View File
@@ -14,11 +14,11 @@ interface TopicEvent {
topicPrefix?: string;
}
await on<TopicEvent>("topic.provisioned", async (e) => {
await on<TopicEvent>("module.mosquitto.topic.provisioned", async (e) => {
console.log(`[mosquitto] topic provisioned for ${e.body.consumer} (client ${e.body.username})`);
});
await on<TopicEvent>("topic.deprovisioned", async (e) => {
await on<TopicEvent>("module.mosquitto.topic.deprovisioned", async (e) => {
console.log(`[mosquitto] topic deprovisioned for ${e.body.consumer} (client ${e.body.username})`);
});

Some files were not shown because too many files have changed in this diff Show More