Author SHA1 Message Date
jschoubben 3f0a174392 nats declares the upstream image it is built from
The recipe started FROM the upstream server's digest directly, and the build machine
refuses that: every base is declared under build.on and copied into the mesh's own
store before a build, so a build never reaches out to a registry the mesh does not
run (novox/hq ADR 0097). Found the first time the module was built on a real mesh.

Same digest, now declared as NATS_BASE and arriving as a build argument; the
Dockerfile says where it comes from and why the digest is the index's.
2026-09-27 23:40:16 +02:00
jschoubben c0edebefb9 Merge pull request 'lavinmq, amqp-ping and amqp-email-forwarder leave the catalogue' (#120) from feat/amqp-leaves-the-catalogue into main 2026-09-27 21:30:19 +00:00
jschoubben b9605a6e3b lavinmq, amqp-ping and amqp-email-forwarder leave the catalogue
AMQP is not a provision (novox/hq ADR 0131, design 28 task 5.4). These were the
only three manifests that named it: the broker that provided it, a proof that a
grant worked end to end, and a forwarder reading mail off a queue. Removed, not
converted — a module that wants messaging wants the mesh's bus, reached through
the sdk and named by the mesh-broker seat, and either of the last two is re-done
against that if wanted, as a new module under the record.

The controller refuses a manifest naming amqp from its next release, so these
could not be re-registered anyway. Nothing else in the catalogue referenced them.
2026-09-27 23:27:24 +02:00
jschoubben dca84d3bb4 Merge pull request 'lavinmq holds mesh-broker again, now the check reads the store' (#119) from restore/broker-claim into main 2026-09-27 20:24:22 +00:00
jschoubben 9f9ce92d0f lavinmq holds mesh-broker again, now the check reads the store
The claim comes back for the third and last time. The store's row says the bus seat
answers for `amqp`, lavinmq provides `amqp`, and with mesh-controller#89 the check
that judges a claim reads that row instead of a copy compiled into the build
machine. So this is accepted for the reason it should have been all along.

Restores the holder the controller composes its own bus address through, which is
what ends tonight's crash loop. nats takes the seat over when the cutover is done
deliberately, not because the seat emptied itself.
2026-09-27 22:23:37 +02:00
jschoubben 6e5b2557ba Merge pull request 'Undo: claiming mesh-broker made lavinmq unbuildable' (#118) from revert/broker-seat-claim into main 2026-09-27 19:53:49 +00:00
jschoubben f97b7dd544 Undo: lavinmq cannot hold mesh-broker, and claiming it makes lavinmq unbuildable
Putting the claim back was wrong on its own terms. `mesh-broker` delivers
`mesh-bus`, and a seat that delivers a provision may only be held by a module that
provides it — so the claim is refused at registration:

  lavinmq claims mesh-broker, whose holder answers for "mesh-bus",
  and lavinmq does not provide "mesh-bus" at mesh scope

Which means the merged claim does not restore the holder, it stops lavinmq being
built at all. Removed again.

The seat being empty is still the live fault, and it has only one valid answer: the
holder must provide `mesh-bus`, and the module that does is nats. Recorded against
the rollout, because it moves a step that was optional into the critical path.
2026-09-27 21:27:55 +02:00
jschoubben 5f5798ec8a Merge pull request 'lavinmq keeps mesh-broker until something else can take it' (#117) from fix/broker-seat-must-stay-held into main 2026-09-27 19:17:45 +00:00
jschoubben 42550dbe43 lavinmq keeps mesh-broker until something else can take it
Taking the claim off made the seat unheld, and the controller dereferences that
seat to find its own bus (to-be 26, "the one exception is the controller itself").
Unheld, the composed address fell back to a default port nothing serves, and the
control plane crash-looped: "cannot reach the broker named in MESH_BROKER_AMQP:
dial tcp 127.0.0.1:5672". The broker itself never stopped — it is healthy on the
port the mesh actually assigned it.

lavinmq becoming an ordinary provider is right, and it is still a provider of amqp
here. What was wrong is the order: the seat has to pass from one holder to the next,
and it cannot be empty in between, because the thing that reads it is the thing that
would have to fix it.
2026-09-27 21:13:13 +02:00
jschoubben 51713dd631 Merge pull request 'The Go base has to be 1.26 for what compiles the controller's code' (#116) from fix/go-126-base into main 2026-09-27 19:00:38 +00:00
jschoubben 4cda964a43 The Go base has to be 1.26 for what compiles the controller's code
builder and route-proxy both build from the mesh-controller repository's context,
so its go.mod is theirs, and `nats.go v1.54.0` puts that at `go >= 1.26`. Pinned at
1.25.14 they cannot compile it: the build machine's own build failed with "go.mod
requires go >= 1.26.0 (running go 1.25.14)".

Each moves to the 1.26.8 digest of the flavour it already used — alpine for
builder, debian for route-proxy — so nothing changes but the compiler version.
2026-09-27 20:58:09 +02:00
jschoubben adb02da136 Merge pull request 'The nats module, and every manifest's event names made local' (#115) from feat/nats-genesis into main 2026-09-27 17:31:04 +00:00
jschoubben 06954b5a70 Merge main: the trunk's seat names, this branch's event names
Two lines of work renamed the same seats differently. The trunk named them for their
scope — node-scoped ones `node-*`, leaving `the-artifact-store`, `npm-package-registry`
and `git` as they were — and this branch had renamed ten of them to `mesh-*`. The trunk's
set is what the live controller loads and what the live seats were actually renamed to, so
a manifest claiming this branch's name is one the running mesh refuses. Three of them
needed reverting by hand: git had auto-merged this branch's names where the trunk had not
touched those lines, which is the quiet kind of merge result.

Event names are this branch's, because the trunk has not converted them and they are what
issue 127 was about.

Verdaccio goes with the trunk's removal of it. The template work on dnsmasq's roster fact
is the trunk's, sitting beside this branch's local event names in the same file — the one
hunk where both changes landed together.

75 manifests, all parsing, no claim outside the trunk's set and no event name left in the
old bus's form.
2026-09-27 18:25:57 +02:00
jschoubben 3d7d896014 Merge pull request 'fail2ban never bans a tunnel peer: ignoreip names the mesh range' (#113) from fix/fail2ban-ignores-the-mesh-range into main 2026-09-27 14:55:55 +00:00
jschoubben b58a3b487d A build's outcome belongs to the role, not to the module holding it
ADR 0121. The builder declared `built` as its own event, so every consumer depended
on which module happens to be the build machine today. It is the build-machine
role's event now: the builder declares none of its own, and the catalogue listens
for `mesh-build-machine.built` rather than `builder.built`.

Nothing changes about what reaches the catalogue. What changes is that it survives
the build machine being a different module, which is the whole reason the mesh has a
word for a role.
2026-09-27 15:37:21 +02:00
jschoubben 7b06a7a408 Event names are local now, in the manifests and in the code
Every module named its events the way the old bus spelled a routing key —
`module.<module>.<verb>`. Design 29 says a module names an event locally and the
mesh works out where it lands, so all 37 were stale against a rule already
decided. On the new bus that derives into a namespace belonging to a module
called "module", so no cross-module subscription in the mesh matched anything:
nothing failed, nothing reacted (novox/hq 04-ISSUES/127).

36 manifests converted, and 43 files of module code with them. The code mattered
as much as the manifests: the runtime builds the subject from what `emit()` is
handed, so a converted manifest with unconverted code would have had the
permission and the subject disagree.

Three things the new check found on the way:

- `photos` emitted an event its manifest never declared, which the new bus refuses
  outright. Declared.
- `showcase` waited for an event nothing emits, so its demo could never be
  triggered — only `showcase` may publish under its own name. It emits both halves
  now.
- `distribution` declared an event named after a different module. It emits
  `image.pushed` under its own name. An event about a *role* belongs on the seat,
  where the name outlives whoever holds it, but the sdk has no way to publish on a
  seat yet, so that stays recorded rather than declared.

The audit logger's "everything" pattern is `**` rather than the old bus's `#`.
2026-09-27 14:42:28 +02:00
jschoubben f0a6ce8d4a Merge pull request 'Rename seat claims to mesh-*/node-*; retire verdaccio (ADR 0121)' (#112) from feat/system-seats-named-by-scope into main 2026-09-27 12:32:29 +00:00
jschoubben a093c88c32 nats declares its own server settings, and where the mesh's users go
The split the controller now makes, from this side. The module's own
configuration — ports, TLS, JetStream — is a declared file resource, because those
are properties of this container and change when its image does. `bus-users` names
where the mesh writes every account and permission, in the same directory, and the
module's configuration includes it.

**Both files in one directory because they have to be.** An absolute include path
is resolved relative to the including file's directory: nats-server given
`include /etc/nats/accounts.conf` from /etc/nats-server/nats.conf looks for
/etc/nats-server/etc/nats/accounts.conf and refuses to start. Verified against the
server, and recorded in the configuration itself where somebody moving a file will
read it.

**`verify: true` is gone, and it was refusing every connection in the mesh.** It
makes the server demand a client certificate; a host pins this server's exact
certificate and authenticates with the password the mesh minted, and presents none.
Found by building this image and connecting to it as a host would.

The entrypoint now waits for both files and watches the mesh's half: the module's
own does not change without a new declaration, and that recreates the container
anyway. Verified end to end against this image — the mesh's user list rewritten,
the module noticing and reloading the server itself with no signal from outside,
and the connection the mesh already had still working afterwards.
2026-09-27 02:50:35 +02:00
jschoubben f67f0ca9bc Merge pull request 'dnsmasq owns its resolver format: node-zones is a template (ADR 0120)' (#111) from feat/roster-facts-are-templates into main 2026-09-26 23:51:29 +00:00
jschoubben ae99204a8c Claim the renamed seats (novox/hq ADR 0118)
Ten manifests claim mesh-* names now. What they PROVIDE is unchanged: gitea
still provides git and npm-package-registry, and a consumer requires the
interface, not the seat.
2026-09-26 23:07:09 +02:00
jschoubben 86882cacd4 nats provides mesh-bus (novox/hq ADR 0120) 2026-09-26 21:17:21 +02:00
jschoubben ea7f6796e8 lavinmq is a provider, not foundation
It claims no seat: mesh-broker is the NATS server's (novox/hq ADR 0119).
The amqp interface stays exactly as it is — a backing service a module may
require, like a database.
2026-09-26 21:08:40 +02:00
jschoubben 9b063a77b2 nats: the module, and an image that reloads in place
Step 1.1 and 1.2 of novox/hq ADR 0116. The server is a built artifact rather
than the upstream image directly, because it needs an entrypoint of its own:
the host can only recreate a container, and recreating the bus for every
permission change drops every connection and every in-flight ack. nats-server
reloads on SIGHUP by itself, so the config is mounted as a directory (not
digest-tracked, hq issue 103) and the entrypoint watches the one file.

Verified against the real server, not assumed: a user added to the config
connects, a revoked one is refused, both within one poll interval, with the
container's PID and restart count unchanged and "Reloaded: accounts" in its
log.

Two corrections found by checking rather than reading:
- the seat delivers nothing now (hq ADR 0117), and the controller's parser
  refused the manifest until it did — "nats claims mesh-broker, whose holder
  answers for amqp, and nats does not provide amqp"
- pinned to the multi-arch index digest; the first pin was the amd64
  manifest, which builds here and fails on any other architecture
2026-09-26 19:34:14 +02:00
99 changed files with 354 additions and 1205 deletions
-56
View File
@@ -1,56 +0,0 @@
{
"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
@@ -1,31 +0,0 @@
# 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
@@ -1,191 +0,0 @@
// 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
@@ -1,53 +0,0 @@
// 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
@@ -1,89 +0,0 @@
{
"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
@@ -1,14 +0,0 @@
{
"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
@@ -1,12 +0,0 @@
{
"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": [
"module.anthropic-consumer.usage.session"
"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", "module.anthropic-consumer.usage.session", JSON.stringify(body)],
[main, "emit", "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": [
"module.anthropic-manager.usage.read"
"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", "module.anthropic-manager.usage.read", JSON.stringify(body)], {
const child = spawn(process.execPath, [main, "emit", "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("module.umami.site.created", { domain: "my-app" });
await emit("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), ["module.umami.site.created", "node.anchor.joined"]);
assert.deepEqual(lines.map((l) => l.type), ["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 -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("module.bazarr.subtitle.downloaded", {
await emit("subtitle.downloaded", {
kind: entry.kind,
title: entry.title,
language: entry.language,
+1 -1
View File
@@ -5,7 +5,7 @@
"container-runtime"
],
"emits": [
"module.bazarr.subtitle.downloaded"
"subtitle.downloaded"
],
"own-secrets": {
"broker": "/var/lib/mesh/bazarr/broker",
+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("module.bookshelf.book.grabbed", { title: item.title, status: item.status });
if (!inQueue.has(id)) await emit("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("module.bookshelf.download.completed", { title: item.title });
await emit("download.completed", { title: item.title });
}
}
}
+2 -2
View File
@@ -6,8 +6,8 @@
"container-runtime"
],
"emits": [
"module.bookshelf.book.grabbed",
"module.bookshelf.download.completed"
"book.grabbed",
"download.completed"
],
"consumes": [],
"own-secrets": {
+1 -4
View File
@@ -20,9 +20,6 @@
"secrets": {
"npm-package-registry": "/var/lib/mesh/builder/package-registry.secret"
},
"emits": [
"module.builder.built"
],
"own-secrets": {
"broker": "/var/lib/mesh/builder/broker"
},
@@ -80,7 +77,7 @@
"on": [
{
"arg": "GO_BASE",
"image": "golang@sha256:1ae0735f00daffa3aaf1363a5184c0d2dc55c78e3db4ec70241cdac97bf84b59"
"image": "golang@sha256:8ac98ca534ac3f51e1f420a1dd2c15e74c75cfa0f23f3ad27eb5d7236c349a0c"
},
{
"arg": "ALPINE_BASE",
+2 -2
View File
@@ -22,8 +22,8 @@
"broker": "/var/lib/mesh/cloudflare-dns/broker"
},
"emits": [
"module.cloudflare-dns.record.created",
"module.cloudflare-dns.record.removed"
"record.created",
"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("module.cloudflare-dns.record.created", {
await announce("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("module.cloudflare-dns.record.removed", { name: fqdn, consumer: p.as });
await announce("record.removed", { name: fqdn, consumer: p.as });
},
});
+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("module.registry.image.pushed", { repo, tag });
if (primed) await emit("image.pushed", { repo, tag });
seen.add(id);
}
}
+1 -1
View File
@@ -17,7 +17,7 @@
"container-runtime"
],
"emits": [
"module.registry.image.pushed"
"image.pushed"
],
"own-secrets": {
"broker": "/var/lib/mesh/registry/broker"
+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("module.dnsmasq.name.added", { name, address });
if (!known.has(name)) await emit("name.added", { name, address });
}
for (const [name] of known) {
if (!now.has(name)) await emit("module.dnsmasq.name.removed", { name });
if (!now.has(name)) await emit("name.removed", { name });
}
}
known.clear();
+2 -2
View File
@@ -8,8 +8,8 @@
"mesh-addressing"
],
"emits": [
"module.dnsmasq.name.added",
"module.dnsmasq.name.removed"
"name.added",
"name.removed"
],
"own-secrets": {
"broker": "/var/lib/mesh/dnsmasq/broker"
+1 -1
View File
@@ -35,7 +35,7 @@ async function pollRepos(client: GiteaClient): Promise<void> {
for (const repo of repos) {
if (!seen.has(repo.full_name)) {
if (primed) {
await emit("module.gitea.repo.created", {
await emit("repo.created", {
full_name: repo.full_name,
owner: repo.owner,
name: repo.name,
+3 -3
View File
@@ -38,9 +38,9 @@
"container-runtime"
],
"emits": [
"module.gitea.repo.created",
"module.gitea.issue.opened",
"module.gitea.pull.merged"
"repo.created",
"issue.opened",
"pull.merged"
],
"listens": [
{
+2 -2
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("module.gitea.issue.opened", {
await emit("issue.opened", {
owner,
repo,
number: issue.number,
@@ -231,7 +231,7 @@ 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);
await emit("module.gitea.pull.merged", {
await emit("pull.merged", {
owner,
repo,
number,
+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("module.grafana.alert.firing", { name: a.name, labels: a.labels, activeAt: a.activeAt });
await emit("alert.firing", { name: a.name, labels: a.labels, activeAt: a.activeAt });
}
}
}
+1 -1
View File
@@ -2,7 +2,7 @@
"module": "grafana",
"version": "1",
"emits": [
"module.grafana.alert.firing"
"alert.firing"
],
"own-secrets": {
"admin": "/var/lib/grafana-module/admin.secret",
+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("module.home-assistant.state.changed", {
await emit("state.changed", {
entity: s.entity_id,
name: nameOf(s),
from: prev,
+1 -1
View File
@@ -6,7 +6,7 @@
"container-runtime"
],
"emits": [
"module.home-assistant.state.changed"
"state.changed"
],
"own-secrets": {
"broker": "/var/lib/mesh/home-assistant/broker",
+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("module.icecast.stream.started", {
await emit("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("module.icecast.stream.stopped", { mount, name: m.name });
await emit("stream.stopped", { mount, name: m.name });
}
}
}
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime"
],
"emits": [
"module.icecast.stream.started",
"module.icecast.stream.stopped"
"stream.started",
"stream.stopped"
],
"own-secrets": {
"broker": "/var/lib/mesh/icecast/broker"
+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("module.keycloak.user.created", { realm, username, ...(email ? { email } : {}) }),
announce("user.created", { realm, username, ...(email ? { email } : {}) }),
userDeleted: (realm: string, userId: string) =>
announce("module.keycloak.user.deleted", { realm, userId }),
announce("user.deleted", { realm, userId }),
passwordReset: (realm: string, userId: string) =>
announce("module.keycloak.password.reset", { realm, userId }),
announce("password.reset", { realm, userId }),
clientCreated: (realm: string, clientId: string, name?: string) =>
announce("module.keycloak.client.created", { realm, clientId, ...(name ? { name } : {}) }),
announce("client.created", { realm, clientId, ...(name ? { name } : {}) }),
groupCreated: (realm: string, name: string) =>
announce("module.keycloak.group.created", { realm, name }),
announce("group.created", { realm, name }),
roleCreated: (realm: string, name: string) =>
announce("module.keycloak.role.created", { realm, name }),
announce("role.created", { realm, name }),
};
console.log("[keycloak] event surface ready — identity, client, group and role changes are announced");
+6 -6
View File
@@ -25,12 +25,12 @@
"container-runtime"
],
"emits": [
"module.keycloak.user.created",
"module.keycloak.user.deleted",
"module.keycloak.password.reset",
"module.keycloak.client.created",
"module.keycloak.group.created",
"module.keycloak.role.created"
"user.created",
"user.deleted",
"password.reset",
"client.created",
"group.created",
"role.created"
],
"listens": [
{
-42
View File
@@ -1,42 +0,0 @@
# 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
@@ -1,45 +0,0 @@
// 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)`);
-183
View File
@@ -1,183 +0,0 @@
// 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: ".*" });
}
/**
* Whether a consumer's user exists with exactly this password and full permissions on its own
* vhost. Read-only: the stored hash is salted SHA-256, the scheme `rabbitHash` writes, so the
* password is checked by hashing it with the stored salt rather than by logging in. `false` when
* the user or its permission is gone or the password differs; an unreachable API rejects
* (novox/hq issue 120).
*/
async holdsConsumer(login: string, password: string): Promise<boolean> {
const v = encodeURIComponent(login);
const u = encodeURIComponent(login);
const user = await this.getOrNull<{ password_hash?: string; hashing_algorithm?: string }>(`/users/${u}`);
if (!user?.password_hash) return false;
if (user.hashing_algorithm && !/sha256/i.test(user.hashing_algorithm)) {
throw new Error(`lavinmq user ${login} is hashed with ${user.hashing_algorithm}, which this check cannot verify`);
}
const stored = Buffer.from(user.password_hash, "base64");
if (stored.length < 5 || rabbitHash(password, stored.subarray(0, 4)) !== user.password_hash) return false;
const perm = await this.getOrNull<{ configure?: string; write?: string; read?: string }>(`/permissions/${v}/${u}`);
return perm?.configure === ".*" && perm?.write === ".*" && perm?.read === ".*";
}
/** A GET that answers null for a 404 and rejects on anything else that is not 2xx. */
private async getOrNull<T>(path: string): Promise<T | null> {
const resp = await fetch(`${this.conn.base}/api${path}`, {
headers: {
Authorization: "Basic " + Buffer.from(`${this.conn.adminUser}:${this.conn.adminPassword}`).toString("base64"),
},
});
if (resp.status === 404) return null;
if (!resp.ok) throw new Error(`lavinmq management API GET ${path} -> ${resp.status}: ${await resp.text()}`);
return (await resp.json()) as T;
}
/** 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
@@ -1,25 +0,0 @@
// 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");
-151
View File
@@ -1,151 +0,0 @@
{
"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": "broker-data",
"type": "directory",
"path": "/var/lib/mesh-broker",
"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": [
"/var/lib/mesh-broker:/var/lib/lavinmq",
"/var/lib/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"
}
]
},
"accesses": [
{
"path": "/var/lib/mesh-broker-tls",
"mode": "read"
}
]
}
-14
View File
@@ -1,14 +0,0 @@
{
"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.1"
},
"devDependencies": {
"@types/node": "^22.0.0",
"typescript": "^5.6.0"
}
}
-56
View File
@@ -1,56 +0,0 @@
// 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 });
},
// 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 lavinmq.holdsConsumer(p.as, p.password);
},
});
-32
View File
@@ -1,32 +0,0 @@
// 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
@@ -1,12 +0,0 @@
{
"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"]
}
+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("module.lidarr.album.grabbed", { title: item.title, status: item.status });
if (!inQueue.has(id)) await emit("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("module.lidarr.download.completed", { title: item.title });
await emit("download.completed", { title: item.title });
}
}
}
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime"
],
"emits": [
"module.lidarr.album.grabbed",
"module.lidarr.download.completed"
"album.grabbed",
"download.completed"
],
"consumes": [],
"own-secrets": {
+2 -2
View File
@@ -36,8 +36,8 @@ function watcher(created: string, deleted: string): (keys: string[]) => Promise<
};
}
const watchUsers = watcher("module.mailu.user.created", "module.mailu.user.deleted");
const watchAliases = watcher("module.mailu.alias.created", "module.mailu.alias.deleted");
const watchUsers = watcher("user.created", "user.deleted");
const watchAliases = watcher("alias.created", "alias.deleted");
async function pollUsers(): Promise<void> {
await watchUsers((await mailu.listUsers()).map((u) => u.email));
+4 -4
View File
@@ -53,10 +53,10 @@
}
},
"emits": [
"module.mailu.user.created",
"module.mailu.user.deleted",
"module.mailu.alias.created",
"module.mailu.alias.deleted"
"user.created",
"user.deleted",
"alias.created",
"alias.deleted"
],
"listens": [
{
+6 -6
View File
@@ -1,6 +1,6 @@
// mesh-catalog's entrypoint — the module graph's consumer (novox/hq ADR 0070, ADR 0072).
//
// The builder announces what it built; this places it in the graph and announces what that means.
// The build-machine role 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.
//
@@ -47,7 +47,7 @@ interface Built {
replay?: boolean;
}
await on("module.builder.built", async (event) => {
await on("mesh-build-machine.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
@@ -69,7 +69,7 @@ await on("module.builder.built", async (event) => {
// it was missing, and the mesh is told nothing happened, because nothing did.
if (body.replay) return;
await emit("module.mesh-catalog.registered", {
await emit("registered", {
module: body.module, commit: body.commit, upgraded,
});
@@ -77,13 +77,13 @@ await on("module.builder.built", async (event) => {
// through modules that did not change, forever (ADR 0072).
if (!upgraded) return;
await emit("module.mesh-catalog.upgraded", {
await emit("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("module.mesh-catalog.rebuild-needed", {
await emit("rebuild-needed", {
module: next.module,
builtAt: next.commit,
because: next.because,
@@ -101,4 +101,4 @@ await on("module.builder.built", async (event) => {
// 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("module.mesh-catalog.catching-up", {});
await emit("catching-up", {});
+4 -4
View File
@@ -29,12 +29,12 @@
"broker": "/var/lib/mesh/mesh-catalog/broker"
},
"consumes": [
"module.builder.built"
"mesh-build-machine.built"
],
"emits": [
"module.mesh-catalog.registered",
"module.mesh-catalog.upgraded",
"module.mesh-catalog.rebuild-needed"
"registered",
"upgraded",
"rebuild-needed"
],
"resources": [
{
+3 -3
View File
@@ -16,15 +16,15 @@ interface SecretEvent {
rotations?: number;
}
await on<SecretEvent>("module.mesh-vault.secret.provisioned", async (e) => {
await on<SecretEvent>("secret.provisioned", async (e) => {
console.log(`[mesh-vault] secret provisioned for ${e.body.as} on ${e.body.consumer} (${e.body.fingerprint})`);
});
await on<SecretEvent>("module.mesh-vault.secret.rotated", async (e) => {
await on<SecretEvent>("secret.rotated", async (e) => {
console.log(`[mesh-vault] secret rotated for ${e.body.as} — rotation ${e.body.rotations} (${e.body.fingerprint})`);
});
await on<SecretEvent>("module.mesh-vault.secret.deprovisioned", async (e) => {
await on<SecretEvent>("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": [
"module.mesh-vault.secret.provisioned",
"module.mesh-vault.secret.rotated",
"module.mesh-vault.secret.deprovisioned"
"secret.provisioned",
"secret.rotated",
"secret.deprovisioned"
],
"consumes": [
"module.mesh-vault.secret.provisioned",
"module.mesh-vault.secret.rotated",
"module.mesh-vault.secret.deprovisioned"
"mesh-vault.secret.provisioned",
"mesh-vault.secret.rotated",
"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("module.mesh-vault.secret.deprovisioned", { as: p.as });
await announce("secret.deprovisioned", { as: p.as });
},
});
+2 -2
View File
@@ -26,8 +26,8 @@
"container-runtime"
],
"emits": [
"module.minio.bucket.created",
"module.minio.bucket.removed"
"bucket.created",
"bucket.removed"
],
"listens": [
{
+2 -2
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("module.minio.bucket.created", {
await announce("bucket.created", {
bucket,
consumer: p.consumer ?? "",
accessKey: accessKeyId,
@@ -51,7 +51,7 @@ runProvisioner("s3-bucket", {
console.error(`[minio] bucket ${bucket} not removed (likely non-empty), access revoked: ${err}`);
}
await announce("module.minio.bucket.removed", { bucket, accessKey: p.as });
await announce("bucket.removed", { bucket, accessKey: p.as });
},
// Asked every minute by the harness: whether the backend still holds this consumer exactly 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("module.*.usage.*", async (event) => {
await on("*.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": [
"module.*.usage.*"
"*.usage.*"
],
"own-secrets": {
"broker": "/var/lib/mesh/model-usage/broker"
+2 -2
View File
@@ -14,11 +14,11 @@ interface DatabaseEvent {
user?: string;
}
await on<DatabaseEvent>("module.mongodb.database.provisioned", async (e) => {
await on<DatabaseEvent>("database.provisioned", async (e) => {
console.log(`[mongodb] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
});
await on<DatabaseEvent>("module.mongodb.database.deprovisioned", async (e) => {
await on<DatabaseEvent>("database.deprovisioned", async (e) => {
console.log(`[mongodb] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
});
+4 -4
View File
@@ -11,12 +11,12 @@
"container-runtime"
],
"emits": [
"module.mongodb.database.provisioned",
"module.mongodb.database.deprovisioned"
"database.provisioned",
"database.deprovisioned"
],
"consumes": [
"module.mongodb.database.provisioned",
"module.mongodb.database.deprovisioned"
"mongodb.database.provisioned",
"mongodb.database.deprovisioned"
],
"listens": [
{
+2 -2
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("module.mongodb.database.provisioned", {
await announce("database.provisioned", {
consumer: p.consumer ?? "",
database,
user: p.as,
@@ -43,7 +43,7 @@ runProvisioner("mongodb-database", {
async remove(p: { as: string }): Promise<void> {
await mongo.dropDatabaseAndUser(p.as, p.as);
await announce("module.mongodb.database.deprovisioned", { database: 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).
+2 -2
View File
@@ -14,11 +14,11 @@ interface TopicEvent {
topicPrefix?: string;
}
await on<TopicEvent>("module.mosquitto.topic.provisioned", async (e) => {
await on<TopicEvent>("topic.provisioned", async (e) => {
console.log(`[mosquitto] topic provisioned for ${e.body.consumer} (client ${e.body.username})`);
});
await on<TopicEvent>("module.mosquitto.topic.deprovisioned", async (e) => {
await on<TopicEvent>("topic.deprovisioned", async (e) => {
console.log(`[mosquitto] topic deprovisioned for ${e.body.consumer} (client ${e.body.username})`);
});
+4 -4
View File
@@ -12,12 +12,12 @@
"container-runtime"
],
"emits": [
"module.mosquitto.topic.provisioned",
"module.mosquitto.topic.deprovisioned"
"topic.provisioned",
"topic.deprovisioned"
],
"consumes": [
"module.mosquitto.topic.provisioned",
"module.mosquitto.topic.deprovisioned"
"mosquitto.topic.provisioned",
"mosquitto.topic.deprovisioned"
],
"serves": {
"mqtt-topic": {}
+2 -2
View File
@@ -32,7 +32,7 @@ runProvisioner("mqtt-topic", {
// The topic subtree is scoped to the consumer's own login, so one cannot read another's topics.
const topicPrefix = p.as;
await mosquitto.createScopedClient(p.as, p.password, topicPrefix);
await announce("module.mosquitto.topic.provisioned", {
await announce("topic.provisioned", {
consumer: p.consumer ?? "",
username: p.as,
topicPrefix,
@@ -41,7 +41,7 @@ runProvisioner("mqtt-topic", {
async remove(p: { as: string }): Promise<void> {
await mosquitto.deleteScopedClient(p.as);
await announce("module.mosquitto.topic.deprovisioned", { username: p.as });
await announce("topic.deprovisioned", { username: 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).
+2 -2
View File
@@ -14,11 +14,11 @@ interface DatabaseEvent {
user?: string;
}
await on<DatabaseEvent>("module.mssql.database.provisioned", async (e) => {
await on<DatabaseEvent>("database.provisioned", async (e) => {
console.log(`[mssql] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
});
await on<DatabaseEvent>("module.mssql.database.deprovisioned", async (e) => {
await on<DatabaseEvent>("database.deprovisioned", async (e) => {
console.log(`[mssql] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
});
+4 -4
View File
@@ -11,12 +11,12 @@
"container-runtime"
],
"emits": [
"module.mssql.database.provisioned",
"module.mssql.database.deprovisioned"
"database.provisioned",
"database.deprovisioned"
],
"consumes": [
"module.mssql.database.provisioned",
"module.mssql.database.deprovisioned"
"mssql.database.provisioned",
"mssql.database.deprovisioned"
],
"listens": [
{
+2 -2
View File
@@ -33,7 +33,7 @@ runProvisioner("mssql-database", {
// Database, login and user share the consumer's name, so the consumer owns exactly its own.
const database = p.as;
await mssql.createDatabaseAndLogin(database, p.as, p.password);
await announce("module.mssql.database.provisioned", {
await announce("database.provisioned", {
consumer: p.consumer ?? "",
database,
user: p.as,
@@ -42,7 +42,7 @@ runProvisioner("mssql-database", {
async remove(p: { as: string }): Promise<void> {
await mssql.dropDatabaseAndLogin(p.as, p.as);
await announce("module.mssql.database.deprovisioned", { database: 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).
+21
View File
@@ -0,0 +1,21 @@
# nats's server image: the upstream server, plus an entrypoint that reloads it in place when the
# mesh rewrites its configuration. See entrypoint.sh for why that belongs here and not in the host.
#
# **Pinned to the multi-architecture index digest, not a platform's.** `docker manifest inspect`
# reports a platform manifest per architecture and the index that lists them; pinning a platform's
# digest builds on this workstation and fails on any node of another architecture, with an error
# that names a manifest rather than the mistake. This is the index — `docker pull` reports the same
# one, and `RepoDigests` confirms it.
#
# Unlike every other module's Dockerfile, this builds no TypeScript and uses no mesh base image:
# the module's code is the server, which upstream already built. There is no BUILD_BASE here on
# purpose — nothing is compiled. The upstream image is declared in the manifest under build.on and
# arrives as NATS_BASE, like every other base the mesh copies into its own store before a build
# (novox/hq ADR 0097); the digest above is the index one for the reason given.
ARG NATS_BASE
FROM ${NATS_BASE}
COPY entrypoint.sh /usr/local/bin/mesh-nats-entrypoint
RUN chmod 0755 /usr/local/bin/mesh-nats-entrypoint
ENTRYPOINT ["/usr/local/bin/mesh-nats-entrypoint"]
+69
View File
@@ -0,0 +1,69 @@
#!/bin/sh
# nats's entrypoint: run the server, and reload it in place when the mesh rewrites its
# configuration.
#
# **Why this exists inside the module** (novox/hq design 25 §5). The controller composes every
# account and permission into one file, and that file changes whenever a module is added,
# reassigned, or a person's access is granted or revoked — which is often, and on the one server
# everything else depends on. The host has no way to say "reload this container": a
# container resource has `restart-on` and nothing else, and a container's `restart-on` means
# *recreate* — every connection dropped and every in-flight JetStream ack lost, mid-flight, for a
# permission change. `reload-on` is real but it is a *service* field, not a container's.
#
# nats-server already reloads its own configuration on SIGHUP — accounts, permissions, everything
# the mesh composes — without dropping a connection. That is the server's own documented
# capability, not something built for the mesh. So the configuration is mounted as a directory
# (a directory's contents are not digest-tracked the way a directly-mounted file's are, novox/hq
# issue 103), and this watches the one file inside it and signals the server itself. The host's
# only job is what it already does for any directory: keep the file's content current. Nothing
# here is declared `restart-on` or `reload-on`.
set -eu
# **Two files, and only one of them is the mesh's** (novox/hq design 25 §4, task 1.7). CONF is this
# module's own — ports, TLS, JetStream — declared in its manifest, because those are properties of
# the container this module raises. USERS is every account and permission, composed by the
# controller, and CONF includes it. So what is watched here is the mesh's half: the module's own
# does not change without a new declaration, and that recreates the container anyway.
CONF="${MESH_NATS_CONF:-/etc/nats/nats.conf}"
USERS="${MESH_NATS_USERS:-/etc/nats/accounts.conf}"
POLL="${MESH_NATS_CONF_POLL_SECONDS:-5}"
# Both are written as part of the same declaration that creates this container, but none of the
# three are ordered against each other. Waiting is correct and starting without them is not:
# nats-server given a configuration whose include is missing refuses to start, and one given no
# configuration at all comes up with its compiled-in defaults — no TLS, no accounts, every subject
# open to anyone who can reach the port. A bus that is briefly open to everything is not a bus that
# is briefly wrong; it is an open bus.
for needed in "$CONF" "$USERS"; do
while [ ! -s "$needed" ]; do
echo "[nats] waiting for the mesh to write $needed"
sleep 1
done
done
digest() { sha256sum "$USERS" 2>/dev/null | cut -d' ' -f1; }
nats-server --config "$CONF" "$@" &
server=$!
# Forward a stop to the server and let it drain, rather than dying and leaving it orphaned as
# PID 1's child.
stop() { kill -TERM "$server" 2>/dev/null || true; }
trap stop TERM INT
last=$(digest)
while kill -0 "$server" 2>/dev/null; do
sleep "$POLL"
now=$(digest)
# An empty digest means the file is mid-write or briefly gone. Reloading on that would hand the
# server a truncated configuration; the next tick sees the finished one.
[ -n "$now" ] || continue
if [ "$now" != "$last" ]; then
last=$now
echo "[nats] the mesh's user list changed; reloading in place"
kill -HUP "$server" || true
fi
done
# `wait` on an already-exited child still yields its status, which becomes this container's.
wait "$server"
+90
View File
@@ -0,0 +1,90 @@
{
"module": "nats",
"version": "1",
"provides": [
{
"name": "mesh-bus",
"scope": "mesh"
}
],
"claims": [
{
"name": "mesh-broker",
"scope": "mesh"
}
],
"bus-users": "/var/lib/nats-module/conf/accounts.conf",
"capabilities": [
"container-runtime"
],
"emits": [],
"consumes": [],
"listens": [
{
"port": 4222,
"protocol": "tcp",
"from": "mesh",
"why": "the mesh bus \u2014 every link the mesh has, over TLS, reached across the overlay"
}
],
"guards": [
8222
],
"resources": [
{
"id": "jetstream-data",
"type": "directory",
"path": "/var/lib/mesh-broker-nats",
"mode": "0700"
},
{
"id": "conf-dir",
"type": "directory",
"path": "/var/lib/nats-module/conf",
"mode": "0700"
},
{
"id": "server-conf",
"type": "file",
"path": "/var/lib/nats-module/conf/nats.conf",
"content": "# The nats module's own server settings. Declared by the module, because a port, a TLS path\n# and a store directory are properties of the container this module raises: they live in its\n# image and its mounts and change when it does.\n#\n# The mesh writes accounts.conf beside this one and nothing else. A controller that wrote the\n# whole file would have to be kept in step with a Dockerfile it never sees.\n\nport: 4222\nhttp: 127.0.0.1:8222\n\ntls {\n cert_file: \"/tls/tls.crt\"\n key_file: \"/tls/tls.key\"\n ca_file: \"/tls/ca.crt\"\n}\n\n# **No `verify`, deliberately, and it was `verify: true` until a probe ran this image.** That\n# setting makes the server demand a *client* certificate, and nothing in the mesh presents one: a\n# host pins this server's exact certificate and authenticates with the password the mesh minted\n# (novox/hq ADR 0004, design 25 \u00a74), and so does a module's runtime. With it on, every connection\n# in the mesh is refused at the TLS handshake, before any password is looked at \u2014 and the error is\n# \"client didn't provide a certificate\", which reads as a client fault.\n#\n# TLS is still required: a tls block is what makes it required, and verify only decides whether\n# client certificates are checked. What is given up is a second factor the mesh has no machinery\n# to issue or rotate \u2014 a certificate per module per node \u2014 and what is kept is stronger than a\n# name check in both directions: an exact pin outward, a per-user password inward.\n\njetstream {\n store_dir: \"/data\"\n}\n\n# Every user of the mesh, composed by the controller and rewritten whenever a module is\n# assigned, a node enrols or a person's access changes.\n#\n# **Relative, and in this same directory, because it has to be.** An absolute include path is\n# resolved relative to the including file's directory, not from the root: nats-server given\n# `include /etc/nats/accounts.conf` from /etc/nats-server/nats.conf looks for\n# /etc/nats-server/etc/nats/accounts.conf and refuses to start. Verified against the server.\ninclude accounts.conf\n",
"mode": "0644"
},
{
"id": "server",
"type": "container",
"name": "mesh-broker-nats",
"ports": [
"4222:4222",
"127.0.0.1:8222:8222"
],
"volumes": [
"/var/lib/mesh-broker-nats:/data",
"/var/lib/nats-module/conf:/etc/nats:ro",
"/var/lib/mesh-broker-nats-tls:/tls:ro"
],
"artifact": "server"
}
],
"accesses": [
{
"path": "/var/lib/mesh-broker-nats-tls",
"mode": "read"
}
],
"build": {
"on": [
{
"arg": "NATS_BASE",
"image": "nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927"
}
],
"artifacts": [
{
"name": "server",
"kind": "image",
"from": "Dockerfile"
}
]
}
}
+2 -2
View File
@@ -26,7 +26,7 @@ async function pollUsers(client: NextcloudClient): Promise<void> {
const users = client.listUsers();
for (const u of users) {
if (knownUsers.has(u.uid)) continue;
if (usersPrimed) await emit("module.nextcloud.user.created", { uid: u.uid, displayName: u.displayName });
if (usersPrimed) await emit("user.created", { uid: u.uid, displayName: u.displayName });
knownUsers.add(u.uid);
}
usersPrimed = true;
@@ -38,7 +38,7 @@ async function pollShares(client: NextcloudClient): Promise<void> {
const shares = await client.listShares();
for (const s of shares) {
if (knownShares.has(s.id)) continue;
if (sharesPrimed) await emit("module.nextcloud.share.created", { id: s.id, path: s.path, shareType: s.shareType, shareWith: s.shareWith, owner: s.owner });
if (sharesPrimed) await emit("share.created", { id: s.id, path: s.path, shareType: s.shareType, shareWith: s.shareWith, owner: s.owner });
knownShares.add(s.id);
}
sharesPrimed = true;
+2 -2
View File
@@ -26,8 +26,8 @@
"s3-bucket": "${dir:state}/store.secret"
},
"emits": [
"module.nextcloud.user.created",
"module.nextcloud.share.created"
"user.created",
"share.created"
],
"own-secrets": {
"admin": "${dir:state}/admin.secret",
+1 -1
View File
@@ -2,7 +2,7 @@
"module": "nodered",
"version": "1",
"emits": [
"module.nodered.flows.deployed"
"flows.deployed"
],
"own-secrets": {
"broker": "/var/lib/mesh/nodered/broker"
+1 -1
View File
@@ -43,7 +43,7 @@ export function getNodeRedTools(nodered: NodeRedClient): ToolDefinition[] {
const result = await nodered.deployFlows(flows, type);
// Best-effort announcement — a deploy must not fail because the broker is unbound here.
try {
await emit("module.nodered.flows.deployed", { rev: result.rev, nodeCount: result.nodeCount, type });
await emit("flows.deployed", { rev: result.rev, nodeCount: result.nodeCount, type });
} catch (err) {
console.error(`[nodered] deployed but could not emit: ${err}`);
}
+2 -2
View File
@@ -27,7 +27,7 @@ async function pollQueue(): Promise<void> {
if (queuePrimed) {
for (const item of items) {
if (!inQueue.has(item.id)) {
await emit("module.nzbget.download.added", { name: item.name, category: item.category, sizeMB: item.sizeMB });
await emit("download.added", { name: item.name, category: item.category, sizeMB: item.sizeMB });
}
}
}
@@ -45,7 +45,7 @@ async function pollHistory(): Promise<void> {
// A newly-appeared history entry is a completion only if it actually succeeded; a failure or
// a manual delete lands in history too, and neither is a "download.completed".
if (historyPrimed && item.success) {
await emit("module.nzbget.download.completed", { name: item.name, category: item.category, sizeMB: item.sizeMB });
await emit("download.completed", { name: item.name, category: item.category, sizeMB: item.sizeMB });
}
seenHistory.add(item.id);
}
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime"
],
"emits": [
"module.nzbget.download.added",
"module.nzbget.download.completed"
"download.added",
"download.completed"
],
"consumes": [],
"own-secrets": {
+2 -2
View File
@@ -31,7 +31,7 @@ async function pollRequests(): Promise<void> {
const key = keyOf(r);
const known = approvedState.has(key);
if (primed && !known) {
await emit("module.ombi.request.created", {
await emit("request.created", {
kind: r.kind,
id: r.id,
title: r.title,
@@ -41,7 +41,7 @@ async function pollRequests(): Promise<void> {
}
// Approval: the flag went from false to true for a request we already knew about.
if (primed && known && r.approved && approvedState.get(key) === false) {
await emit("module.ombi.request.approved", { kind: r.kind, id: r.id, title: r.title, tmdbId: r.tmdbId });
await emit("request.approved", { kind: r.kind, id: r.id, title: r.title, tmdbId: r.tmdbId });
}
approvedState.set(key, r.approved);
}
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime"
],
"emits": [
"module.ombi.request.created",
"module.ombi.request.approved"
"request.created",
"request.approved"
],
"own-secrets": {
"broker": "/var/lib/mesh/ombi/broker",
+1 -1
View File
@@ -20,7 +20,7 @@ async function pollRecent(): Promise<void> {
for (const asset of items) {
if (!seen.has(asset.id)) {
if (primed) {
await emit("module.photos.item.added", {
await emit("item.added", {
id: asset.id,
fileName: asset.fileName,
kind: asset.type,
+3
View File
@@ -4,6 +4,9 @@
"capabilities": [
"container-runtime"
],
"emits": [
"item.added"
],
"requires": [
"s3-bucket",
"mongodb-database",
+4 -4
View File
@@ -24,10 +24,10 @@ async function pollSessions(): Promise<void> {
const now = new Map(sessions.map((s) => [s.key, s]));
if (playbackPrimed) {
for (const [key, s] of now) {
if (!active.has(key)) await emit("module.plex.playback.started", { title: s.title, user: s.user, player: s.player, kind: s.type });
if (!active.has(key)) await emit("playback.started", { title: s.title, user: s.user, player: s.player, kind: s.type });
}
for (const [key, s] of active) {
if (!now.has(key)) await emit("module.plex.playback.stopped", { title: s.title, user: s.user, player: s.player });
if (!now.has(key)) await emit("playback.stopped", { title: s.title, user: s.user, player: s.player });
}
}
active.clear();
@@ -44,7 +44,7 @@ async function pollRecent(): Promise<void> {
for (const item of items) {
const id = `${item.title}@${item.addedAt ?? ""}`;
if (!seen.has(id)) {
if (itemsPrimed) await emit("module.plex.item.added", item);
if (itemsPrimed) await emit("item.added", item);
seen.add(id);
}
}
@@ -53,7 +53,7 @@ async function pollRecent(): Promise<void> {
// A downloader finished somewhere on the mesh: rescan, so what it fetched becomes a visible item
// rather than a file Plex has not noticed. Idempotent — a rescan too many costs a little disk I/O.
await on("module.*.download.completed", async () => {
await on("*.download.completed", async () => {
await plex.refreshAll();
});
+4 -4
View File
@@ -5,12 +5,12 @@
"container-runtime"
],
"emits": [
"module.plex.playback.started",
"module.plex.playback.stopped",
"module.plex.item.added"
"playback.started",
"playback.stopped",
"item.added"
],
"consumes": [
"module.*.download.completed"
"*.download.completed"
],
"own-secrets": {
"broker": "/var/lib/mesh/plex/broker",
+2 -2
View File
@@ -14,11 +14,11 @@ interface DatabaseEvent {
user?: string;
}
await on<DatabaseEvent>("module.postgres.database.provisioned", async (e) => {
await on<DatabaseEvent>("database.provisioned", async (e) => {
console.log(`[postgres] database provisioned for ${e.body.consumer} (db ${e.body.database})`);
});
await on<DatabaseEvent>("module.postgres.database.deprovisioned", async (e) => {
await on<DatabaseEvent>("database.deprovisioned", async (e) => {
console.log(`[postgres] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
});
+4 -4
View File
@@ -17,12 +17,12 @@
"container-runtime"
],
"emits": [
"module.postgres.database.provisioned",
"module.postgres.database.deprovisioned"
"database.provisioned",
"database.deprovisioned"
],
"consumes": [
"module.postgres.database.provisioned",
"module.postgres.database.deprovisioned"
"postgres.database.provisioned",
"postgres.database.deprovisioned"
],
"listens": [
{
+2 -2
View File
@@ -34,7 +34,7 @@ runProvisioner("postgres-database", {
// Database and owning role share the consumer's login, so the consumer owns exactly its own.
const database = p.as;
await postgres.createDatabaseAndRole(database, p.as, p.password);
await announce("module.postgres.database.provisioned", {
await announce("database.provisioned", {
consumer: p.consumer ?? "",
database,
user: p.as,
@@ -43,7 +43,7 @@ runProvisioner("postgres-database", {
async remove(p: { as: string }): Promise<void> {
await postgres.dropDatabaseAndRole(p.as, p.as);
await announce("module.postgres.database.deprovisioned", { database: 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).
+2 -2
View File
@@ -30,9 +30,9 @@ async function pollTorrents(): Promise<void> {
for (const [hash, t] of now) {
const before = progressByHash.get(hash);
if (before === undefined) {
await emit("module.qbittorrent.download.added", { name: t.name, category: t.category, sizeBytes: t.sizeBytes });
await emit("download.added", { name: t.name, category: t.category, sizeBytes: t.sizeBytes });
} else if (before < 1 && t.progress >= 1) {
await emit("module.qbittorrent.download.completed", { name: t.name, category: t.category, sizeBytes: t.sizeBytes });
await emit("download.completed", { name: t.name, category: t.category, sizeBytes: t.sizeBytes });
}
}
}
+2 -2
View File
@@ -6,8 +6,8 @@
"container-runtime"
],
"emits": [
"module.qbittorrent.download.added",
"module.qbittorrent.download.completed"
"download.added",
"download.completed"
],
"consumes": [],
"own-secrets": {
+2 -2
View File
@@ -40,12 +40,12 @@ async function pollQueue(radarr: RadarrClient): Promise<void> {
if (primed) {
// Entered the queue since last look — Radarr grabbed a release.
for (const [id, item] of now) {
if (!inQueue.has(id)) await emit("module.radarr.movie.grabbed", { title: item.title, status: item.status });
if (!inQueue.has(id)) await emit("movie.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("module.radarr.download.completed", { title: item.title });
await emit("download.completed", { title: item.title });
}
}
}
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime"
],
"emits": [
"module.radarr.movie.grabbed",
"module.radarr.download.completed"
"movie.grabbed",
"download.completed"
],
"consumes": [],
"own-secrets": {
+2 -2
View File
@@ -14,11 +14,11 @@ interface CacheEvent {
keyspacePrefix?: string;
}
await on<CacheEvent>("module.redis.cache.provisioned", async (e) => {
await on<CacheEvent>("cache.provisioned", async (e) => {
console.log(`[redis] cache provisioned for ${e.body.consumer} (user ${e.body.username})`);
});
await on<CacheEvent>("module.redis.cache.deprovisioned", async (e) => {
await on<CacheEvent>("cache.deprovisioned", async (e) => {
console.log(`[redis] cache deprovisioned for ${e.body.consumer} (user ${e.body.username})`);
});
+4 -4
View File
@@ -14,12 +14,12 @@
"container-runtime"
],
"emits": [
"module.redis.cache.provisioned",
"module.redis.cache.deprovisioned"
"cache.provisioned",
"cache.deprovisioned"
],
"consumes": [
"module.redis.cache.provisioned",
"module.redis.cache.deprovisioned"
"redis.cache.provisioned",
"redis.cache.deprovisioned"
],
"serves": {
"redis-cache": {
+2 -2
View File
@@ -32,7 +32,7 @@ runProvisioner("redis-cache", {
// The keyspace is scoped to the consumer's own login, so one cannot read another's keys.
const keyspacePrefix = p.as;
await redis.createAclUser(p.as, p.password, keyspacePrefix);
await announce("module.redis.cache.provisioned", {
await announce("cache.provisioned", {
consumer: p.consumer ?? "",
username: p.as,
keyspacePrefix,
@@ -41,7 +41,7 @@ runProvisioner("redis-cache", {
async remove(p: { as: string }): Promise<void> {
await redis.deleteAclUser(p.as);
await announce("module.redis.cache.deprovisioned", { username: p.as });
await announce("cache.deprovisioned", { username: p.as });
},
// This server keeps its ACL users in memory only, so a restart of it forgets every consumer while
+1 -1
View File
@@ -173,7 +173,7 @@
"on": [
{
"arg": "GO_BASE",
"image": "golang@sha256:699337d620559a59b4a2bb298ad59611e535d2ee755a34cf2d2a98f37578dc80"
"image": "golang@sha256:6c2a5538f964f1c82f97ad14988bf05de100d922d159d0e398b54c7b0ca0c6c9"
},
{
"arg": "ALPINE_BASE",
+2 -2
View File
@@ -5,8 +5,8 @@
// It is here so the module exercises the shape rather than describing it.
import { on, emit } from "@novox/mesh-sdk/events";
await on<{ who?: string }>("module.showcase.greeted", async (event) => {
await on<{ who?: string }>("greeted", async (event) => {
console.log(`[showcase] greeted ${event.body.who ?? "somebody"}`);
// A consumer may emit, which is what makes an event graph rather than a list of sinks.
await emit("module.showcase.acknowledged", { who: event.body.who ?? "somebody" });
await emit("acknowledged", { who: event.body.who ?? "somebody" });
});
+4 -3
View File
@@ -35,10 +35,11 @@
}
],
"emits": [
"module.showcase.acknowledged"
"greeted",
"acknowledged"
],
"consumes": [
"module.showcase.greeted"
"showcase.greeted"
],
"listens": [
{
@@ -172,7 +173,7 @@
{
"id": "tools",
"type": "container",
"name": "mesh-showcase",
"name": "the-showcase",
"artifact": "helper",
"network": "showcase",
"volumes": [
+2 -2
View File
@@ -40,12 +40,12 @@ async function pollQueue(sonarr: SonarrClient): Promise<void> {
if (primed) {
// Entered the queue since last look — Sonarr grabbed a release.
for (const [id, item] of now) {
if (!inQueue.has(id)) await emit("module.sonarr.episode.grabbed", { title: item.title, status: item.status });
if (!inQueue.has(id)) await emit("episode.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("module.sonarr.download.completed", { title: item.title });
await emit("download.completed", { title: item.title });
}
}
}
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime"
],
"emits": [
"module.sonarr.episode.grabbed",
"module.sonarr.download.completed"
"episode.grabbed",
"download.completed"
],
"consumes": [],
"own-secrets": {
-26
View File
@@ -1,26 +0,0 @@
{
"module": "ssh-client",
"version": "1",
"resources": [
{
"id": "openssh",
"type": "package",
"package": "openssh"
},
{
"id": "ssh-dir",
"type": "directory",
"path": "${machine:account-home}/.ssh",
"owner": "${machine:account}",
"mode": "0700"
}
],
"facts": {
"ssh-config": {
"path": ".ssh/config",
"home": true,
"shared": true,
"template": "# The mesh's Host blocks — every other node, so `ssh <node>` reaches it as the\n# right account. This region is replaced whenever a node joins, leaves or is\n# renamed; the rest of this file is yours and is kept untouched.\n{{range .Machines}}{{if ne .Name $.Node}}\nHost {{.Name}} {{.FQDN}}\n HostName {{.FQDN}}\n{{if .Account}} User {{.Account}}\n{{end}}{{end}}{{end}}"
}
}
}
+1 -1
View File
@@ -33,7 +33,7 @@ async function pollHistory(client: TautulliClient): Promise<void> {
}
async function emitWatch(w: TautulliWatch): Promise<void> {
await emit("module.tautulli.watch.recorded", {
await emit("watch.recorded", {
title: w.title, user: w.user, mediaType: w.mediaType,
watchedStatus: w.watchedStatus, percentComplete: w.percentComplete, at: w.date,
});
+1 -1
View File
@@ -2,7 +2,7 @@
"module": "tautulli",
"version": "1",
"emits": [
"module.tautulli.watch.recorded"
"watch.recorded"
],
"own-secrets": {
"broker": "/var/lib/mesh/tautulli/broker"