Author SHA1 Message Date
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
85 changed files with 353 additions and 181 deletions
+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": [
{
+2 -2
View File
@@ -14,11 +14,11 @@ interface AmqpEvent {
vhost?: string;
}
await on<AmqpEvent>("module.lavinmq.amqp.provisioned", async (e) => {
await on<AmqpEvent>("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) => {
await on<AmqpEvent>("amqp.deprovisioned", async (e) => {
console.log(`[lavinmq] broker deprovisioned (user ${e.body.user})`);
});
+4 -4
View File
@@ -17,12 +17,12 @@
"container-runtime"
],
"emits": [
"module.lavinmq.amqp.provisioned",
"module.lavinmq.amqp.deprovisioned"
"amqp.provisioned",
"amqp.deprovisioned"
],
"consumes": [
"module.lavinmq.amqp.provisioned",
"module.lavinmq.amqp.deprovisioned"
"lavinmq.amqp.provisioned",
"lavinmq.amqp.deprovisioned"
],
"serves": {
"amqp": {
+2 -2
View File
@@ -37,7 +37,7 @@ runProvisioner("amqp", {
// 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", {
await announce("amqp.provisioned", {
consumer: p.consumer ?? "",
user: p.as,
vhost: p.as,
@@ -46,7 +46,7 @@ runProvisioner("amqp", {
async remove(p: { as: string }): Promise<void> {
await lavinmq.removeConsumer(p.as);
await announce("module.lavinmq.amqp.deprovisioned", { user: p.as, vhost: p.as });
await announce("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).
+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).
+18
View File
@@ -0,0 +1,18 @@
# 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.
FROM nats@sha256:b83efabe3e7def1e0a4a31ec6e078999bb17c80363f881df35edc70fcb6bb927
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"
+84
View File
@@ -0,0 +1,84 @@
{
"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": {
"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": {
+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"