The nats module, and every manifest's event names made local #115

Merged
jschoubben merged 8 commits from feat/nats-genesis into main 2026-09-27 17:31:04 +00:00
84 changed files with 351 additions and 185 deletions
+1 -1
View File
@@ -18,7 +18,7 @@
"broker": "/var/lib/mesh/anthropic-consumer/broker" "broker": "/var/lib/mesh/anthropic-consumer/broker"
}, },
"emits": [ "emits": [
"module.anthropic-consumer.usage.session" "usage.session"
], ],
"resources": [ "resources": [
{ {
+1 -1
View File
@@ -123,7 +123,7 @@ async function emitUsage(body: Record<string, unknown>): Promise<void> {
await new Promise<void>((resolve) => { await new Promise<void>((resolve) => {
const child = spawn( const child = spawn(
process.execPath, process.execPath,
[main, "emit", "module.anthropic-consumer.usage.session", JSON.stringify(body)], [main, "emit", "usage.session", JSON.stringify(body)],
{ stdio: "inherit" }, { stdio: "inherit" },
); );
child.on("exit", () => resolve()); child.on("exit", () => resolve());
+1 -1
View File
@@ -18,7 +18,7 @@
"broker": "/var/lib/mesh/anthropic-manager/broker" "broker": "/var/lib/mesh/anthropic-manager/broker"
}, },
"emits": [ "emits": [
"module.anthropic-manager.usage.read" "usage.read"
], ],
"resources": [ "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 main = process.env.MESH_TOOLS_MAIN ?? "/app/dist/main.js";
const { spawn } = await import("node:child_process"); const { spawn } = await import("node:child_process");
await new Promise<void>((resolve) => { 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", stdio: "inherit",
}); });
child.on("exit", () => resolve()); child.on("exit", () => resolve());
+1 -1
View File
@@ -3,7 +3,7 @@
"version": "1", "version": "1",
"slug": "audit", "slug": "audit",
"consumes": [ "consumes": [
"#" "**"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/audit-logger/broker" "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"); const path = join(dir, "audit.log");
// The audit-logger's whole behaviour: consume everything, record it. // 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_MODULE = "umami";
process.env.MESH_NODE = "anchor"; 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 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)); const lines = (await readFile(path, "utf8")).trim().split("\n").map((l) => JSON.parse(l));
assert.equal(lines.length, 2); 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].source, "umami");
assert.equal(lines[0].node, "anchor"); assert.equal(lines[0].node, "anchor");
assert.equal(lines[0].body.domain, "my-app"); 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) { for (const entry of entries) {
if (seen.has(entry.id)) continue; if (seen.has(entry.id)) continue;
if (primed) { if (primed) {
await emit("module.bazarr.subtitle.downloaded", { await emit("subtitle.downloaded", {
kind: entry.kind, kind: entry.kind,
title: entry.title, title: entry.title,
language: entry.language, language: entry.language,
+1 -1
View File
@@ -5,7 +5,7 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.bazarr.subtitle.downloaded" "subtitle.downloaded"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/bazarr/broker", "broker": "/var/lib/mesh/bazarr/broker",
+2 -2
View File
@@ -45,12 +45,12 @@ async function pollQueue(bookshelf: BookshelfClient): Promise<void> {
if (primed) { if (primed) {
// Entered the queue since last look — Bookshelf grabbed a release. // Entered the queue since last look — Bookshelf grabbed a release.
for (const [id, item] of now) { 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. // Left the queue — imported and done, unless it was last seen failing.
for (const [id, item] of inQueue) { for (const [id, item] of inQueue) {
if (!now.has(id) && !FAILED_STATUSES.has(item.status)) { 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" "container-runtime"
], ],
"emits": [ "emits": [
"module.bookshelf.book.grabbed", "book.grabbed",
"module.bookshelf.download.completed" "download.completed"
], ],
"consumes": [], "consumes": [],
"own-secrets": { "own-secrets": {
-3
View File
@@ -20,9 +20,6 @@
"secrets": { "secrets": {
"npm-package-registry": "/var/lib/mesh/builder/package-registry.secret" "npm-package-registry": "/var/lib/mesh/builder/package-registry.secret"
}, },
"emits": [
"module.builder.built"
],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/builder/broker" "broker": "/var/lib/mesh/builder/broker"
}, },
+2 -2
View File
@@ -22,8 +22,8 @@
"broker": "/var/lib/mesh/cloudflare-dns/broker" "broker": "/var/lib/mesh/cloudflare-dns/broker"
}, },
"emits": [ "emits": [
"module.cloudflare-dns.record.created", "record.created",
"module.cloudflare-dns.record.removed" "record.removed"
], ],
"resources": [ "resources": [
{ {
+2 -2
View File
@@ -20,7 +20,7 @@ runProvisioner("public-dns", {
async create(p: Provision): Promise<void> { async create(p: Provision): Promise<void> {
const fqdn = cloudflare.nameFor(p.as); const fqdn = cloudflare.nameFor(p.as);
await cloudflare.upsert(fqdn); await cloudflare.upsert(fqdn);
await announce("module.cloudflare-dns.record.created", { await announce("record.created", {
name: fqdn, name: fqdn,
target: cloudflare.ingress, target: cloudflare.ingress,
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
@@ -30,7 +30,7 @@ runProvisioner("public-dns", {
async remove(p: { as: string }): Promise<void> { async remove(p: { as: string }): Promise<void> {
const fqdn = cloudflare.nameFor(p.as); const fqdn = cloudflare.nameFor(p.as);
await cloudflare.remove(fqdn); 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) { for (const tag of tags) {
const id = `${repo}:${tag}`; const id = `${repo}:${tag}`;
if (!seen.has(id)) { 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); seen.add(id);
} }
} }
+1 -1
View File
@@ -17,7 +17,7 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.registry.image.pushed" "image.pushed"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/registry/broker" "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])); const now = new Map((await dnsmasq.answeredNames()).map((a) => [a.name, a.address]));
if (primed) { if (primed) {
for (const [name, address] of now) { 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) { 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(); known.clear();
+2 -2
View File
@@ -8,8 +8,8 @@
"mesh-addressing" "mesh-addressing"
], ],
"emits": [ "emits": [
"module.dnsmasq.name.added", "name.added",
"module.dnsmasq.name.removed" "name.removed"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/dnsmasq/broker" "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) { for (const repo of repos) {
if (!seen.has(repo.full_name)) { if (!seen.has(repo.full_name)) {
if (primed) { if (primed) {
await emit("module.gitea.repo.created", { await emit("repo.created", {
full_name: repo.full_name, full_name: repo.full_name,
owner: repo.owner, owner: repo.owner,
name: repo.name, name: repo.name,
+3 -3
View File
@@ -38,9 +38,9 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.gitea.repo.created", "repo.created",
"module.gitea.issue.opened", "issue.opened",
"module.gitea.pull.merged" "pull.merged"
], ],
"listens": [ "listens": [
{ {
+2 -2
View File
@@ -124,7 +124,7 @@ export function getGiteaTools(gitea: GiteaClient): ToolDefinition[] {
labels: labelIds, labels: labelIds,
}); });
// The mesh just opened an issue — announce it the moment it exists. // The mesh just opened an issue — announce it the moment it exists.
await emit("module.gitea.issue.opened", { await emit("issue.opened", {
owner, owner,
repo, repo,
number: issue.number, 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. // 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); const pull = await gitea.getPullRequest(owner, repo, number);
await gitea.mergePullRequest(owner, repo, number, method, deleteBranch); await gitea.mergePullRequest(owner, repo, number, method, deleteBranch);
await emit("module.gitea.pull.merged", { await emit("pull.merged", {
owner, owner,
repo, repo,
number, number,
+1 -1
View File
@@ -40,7 +40,7 @@ async function pollAlerts(client: GrafanaClient): Promise<void> {
for (const key of now) { for (const key of now) {
if (!firing.has(key)) { if (!firing.has(key)) {
const a = byKey.get(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", "module": "grafana",
"version": "1", "version": "1",
"emits": [ "emits": [
"module.grafana.alert.firing" "alert.firing"
], ],
"own-secrets": { "own-secrets": {
"admin": "/var/lib/grafana-module/admin.secret", "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) { for (const s of states) {
const prev = lastState.get(s.entity_id); const prev = lastState.get(s.entity_id);
if (primed && prev !== undefined && prev !== s.state) { if (primed && prev !== undefined && prev !== s.state) {
await emit("module.home-assistant.state.changed", { await emit("state.changed", {
entity: s.entity_id, entity: s.entity_id,
name: nameOf(s), name: nameOf(s),
from: prev, from: prev,
+1 -1
View File
@@ -6,7 +6,7 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.home-assistant.state.changed" "state.changed"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/home-assistant/broker", "broker": "/var/lib/mesh/home-assistant/broker",
+2 -2
View File
@@ -23,7 +23,7 @@ async function pollMounts(): Promise<void> {
if (primed) { if (primed) {
for (const [mount, m] of now) { for (const [mount, m] of now) {
if (!live.has(mount)) { if (!live.has(mount)) {
await emit("module.icecast.stream.started", { await emit("stream.started", {
mount, mount,
name: m.name, name: m.name,
description: m.description, description: m.description,
@@ -33,7 +33,7 @@ async function pollMounts(): Promise<void> {
} }
for (const [mount, m] of live) { for (const [mount, m] of live) {
if (!now.has(mount)) { 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" "container-runtime"
], ],
"emits": [ "emits": [
"module.icecast.stream.started", "stream.started",
"module.icecast.stream.stopped" "stream.stopped"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/icecast/broker" "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 = { export const events = {
userCreated: (realm: string, username: string, email?: string) => 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) => userDeleted: (realm: string, userId: string) =>
announce("module.keycloak.user.deleted", { realm, userId }), announce("user.deleted", { realm, userId }),
passwordReset: (realm: string, userId: string) => passwordReset: (realm: string, userId: string) =>
announce("module.keycloak.password.reset", { realm, userId }), announce("password.reset", { realm, userId }),
clientCreated: (realm: string, clientId: string, name?: string) => 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) => groupCreated: (realm: string, name: string) =>
announce("module.keycloak.group.created", { realm, name }), announce("group.created", { realm, name }),
roleCreated: (realm: string, name: string) => 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"); console.log("[keycloak] event surface ready — identity, client, group and role changes are announced");
+6 -6
View File
@@ -25,12 +25,12 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.keycloak.user.created", "user.created",
"module.keycloak.user.deleted", "user.deleted",
"module.keycloak.password.reset", "password.reset",
"module.keycloak.client.created", "client.created",
"module.keycloak.group.created", "group.created",
"module.keycloak.role.created" "role.created"
], ],
"listens": [ "listens": [
{ {
+2 -2
View File
@@ -14,11 +14,11 @@ interface AmqpEvent {
vhost?: string; 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})`); 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})`); console.log(`[lavinmq] broker deprovisioned (user ${e.body.user})`);
}); });
+4 -10
View File
@@ -7,22 +7,16 @@
"scope": "mesh" "scope": "mesh"
} }
], ],
"claims": [
{
"name": "mesh-broker",
"scope": "mesh"
}
],
"capabilities": [ "capabilities": [
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.lavinmq.amqp.provisioned", "amqp.provisioned",
"module.lavinmq.amqp.deprovisioned" "amqp.deprovisioned"
], ],
"consumes": [ "consumes": [
"module.lavinmq.amqp.provisioned", "lavinmq.amqp.provisioned",
"module.lavinmq.amqp.deprovisioned" "lavinmq.amqp.deprovisioned"
], ],
"serves": { "serves": {
"amqp": { "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. // The vhost and the user share the consumer's login, so one cannot reach another's broker.
await lavinmq.waitReady(); await lavinmq.waitReady();
await lavinmq.createConsumer(p.as, p.password); await lavinmq.createConsumer(p.as, p.password);
await announce("module.lavinmq.amqp.provisioned", { await announce("amqp.provisioned", {
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
user: p.as, user: p.as,
vhost: p.as, vhost: p.as,
@@ -46,7 +46,7 @@ runProvisioner("amqp", {
async remove(p: { as: string }): Promise<void> { async remove(p: { as: string }): Promise<void> {
await lavinmq.removeConsumer(p.as); 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 // Asked every minute by the harness: whether the backend still holds this consumer exactly as
// the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120).
+2 -2
View File
@@ -40,12 +40,12 @@ async function pollQueue(lidarr: LidarrClient): Promise<void> {
if (primed) { if (primed) {
// Entered the queue since last look — Lidarr grabbed a release. // Entered the queue since last look — Lidarr grabbed a release.
for (const [id, item] of now) { 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. // Left the queue — imported and done, unless it was last seen failing.
for (const [id, item] of inQueue) { for (const [id, item] of inQueue) {
if (!now.has(id) && !FAILED_STATUSES.has(item.status)) { 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" "container-runtime"
], ],
"emits": [ "emits": [
"module.lidarr.album.grabbed", "album.grabbed",
"module.lidarr.download.completed" "download.completed"
], ],
"consumes": [], "consumes": [],
"own-secrets": { "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 watchUsers = watcher("user.created", "user.deleted");
const watchAliases = watcher("module.mailu.alias.created", "module.mailu.alias.deleted"); const watchAliases = watcher("alias.created", "alias.deleted");
async function pollUsers(): Promise<void> { async function pollUsers(): Promise<void> {
await watchUsers((await mailu.listUsers()).map((u) => u.email)); await watchUsers((await mailu.listUsers()).map((u) => u.email));
+4 -4
View File
@@ -53,10 +53,10 @@
} }
}, },
"emits": [ "emits": [
"module.mailu.user.created", "user.created",
"module.mailu.user.deleted", "user.deleted",
"module.mailu.alias.created", "alias.created",
"module.mailu.alias.deleted" "alias.deleted"
], ],
"listens": [ "listens": [
{ {
+6 -6
View File
@@ -1,6 +1,6 @@
// mesh-catalog's entrypoint — the module graph's consumer (novox/hq ADR 0070, ADR 0072). // 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 // 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. // it never has to interpret an artifact or ask this module anything.
// //
@@ -47,7 +47,7 @@ interface Built {
replay?: boolean; replay?: boolean;
} }
await on("module.builder.built", async (event) => { await on("mesh-build-machine.built", async (event) => {
const body = event.body as Built; const body = event.body as Built;
if (!body.module || !body.commit) { if (!body.module || !body.commit) {
// Said rather than dropped: a build that announced itself without saying what it built is a // 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. // it was missing, and the mesh is told nothing happened, because nothing did.
if (body.replay) return; if (body.replay) return;
await emit("module.mesh-catalog.registered", { await emit("registered", {
module: body.module, commit: body.commit, upgraded, 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). // through modules that did not change, forever (ADR 0072).
if (!upgraded) return; if (!upgraded) return;
await emit("module.mesh-catalog.upgraded", { await emit("upgraded", {
module: body.module, commit: body.commit, previous, module: body.module, commit: body.commit, previous,
}); });
// What can be built now — stale, and waiting on nothing that is itself stale. // What can be built now — stale, and waiting on nothing that is itself stale.
for (const next of await graph.buildable()) { for (const next of await graph.buildable()) {
await emit("module.mesh-catalog.rebuild-needed", { await emit("rebuild-needed", {
module: next.module, module: next.module,
builtAt: next.commit, builtAt: next.commit,
because: next.because, 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 // 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. // 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. // 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" "broker": "/var/lib/mesh/mesh-catalog/broker"
}, },
"consumes": [ "consumes": [
"module.builder.built" "mesh-build-machine.built"
], ],
"emits": [ "emits": [
"module.mesh-catalog.registered", "registered",
"module.mesh-catalog.upgraded", "upgraded",
"module.mesh-catalog.rebuild-needed" "rebuild-needed"
], ],
"resources": [ "resources": [
{ {
+3 -3
View File
@@ -16,15 +16,15 @@ interface SecretEvent {
rotations?: number; 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})`); 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})`); 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}`); console.log(`[mesh-vault] secret withdrawn from ${e.body.as}`);
}); });
+6 -6
View File
@@ -11,14 +11,14 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.mesh-vault.secret.provisioned", "secret.provisioned",
"module.mesh-vault.secret.rotated", "secret.rotated",
"module.mesh-vault.secret.deprovisioned" "secret.deprovisioned"
], ],
"consumes": [ "consumes": [
"module.mesh-vault.secret.provisioned", "mesh-vault.secret.provisioned",
"module.mesh-vault.secret.rotated", "mesh-vault.secret.rotated",
"module.mesh-vault.secret.deprovisioned" "mesh-vault.secret.deprovisioned"
], ],
"receives": { "receives": {
"secret": "/var/lib/mesh-vault/grants/mesh.json" "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> { async remove(p: { as: string }): Promise<void> {
if (!ledger.withdraw(p.as)) return; if (!ledger.withdraw(p.as)) return;
console.log(`[mesh-vault] withdrawn: ${p.as}`); 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" "container-runtime"
], ],
"emits": [ "emits": [
"module.minio.bucket.created", "bucket.created",
"module.minio.bucket.removed" "bucket.removed"
], ],
"listens": [ "listens": [
{ {
+2 -2
View File
@@ -30,7 +30,7 @@ runProvisioner("s3-bucket", {
try { await minio.removeAccessKey(accessKeyId); } catch { /* none yet — first provision */ } try { await minio.removeAccessKey(accessKeyId); } catch { /* none yet — first provision */ }
await minio.createAccessKey(bucket, accessKeyId, p.password); await minio.createAccessKey(bucket, accessKeyId, p.password);
await announce("module.minio.bucket.created", { await announce("bucket.created", {
bucket, bucket,
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
accessKey: accessKeyId, accessKey: accessKeyId,
@@ -51,7 +51,7 @@ runProvisioner("s3-bucket", {
console.error(`[minio] bucket ${bucket} not removed (likely non-empty), access revoked: ${err}`); 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 // 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. // to reach the provider over the overlay would block the very apply that brings the overlay up.
await store.migrate(); await store.migrate();
await on("module.*.usage.*", async (event) => { await on("*.usage.*", async (event) => {
const body = event.body as { rows?: UsageRow[]; raw?: unknown }; const body = event.body as { rows?: UsageRow[]; raw?: unknown };
for (const row of body.rows ?? []) { for (const row of body.rows ?? []) {
try { try {
+1 -1
View File
@@ -20,7 +20,7 @@
"postgres-database": "/var/lib/model-usage/database.secret" "postgres-database": "/var/lib/model-usage/database.secret"
}, },
"consumes": [ "consumes": [
"module.*.usage.*" "*.usage.*"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/model-usage/broker" "broker": "/var/lib/mesh/model-usage/broker"
+2 -2
View File
@@ -14,11 +14,11 @@ interface DatabaseEvent {
user?: string; 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})`); 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})`); console.log(`[mongodb] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
}); });
+4 -4
View File
@@ -11,12 +11,12 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.mongodb.database.provisioned", "database.provisioned",
"module.mongodb.database.deprovisioned" "database.deprovisioned"
], ],
"consumes": [ "consumes": [
"module.mongodb.database.provisioned", "mongodb.database.provisioned",
"module.mongodb.database.deprovisioned" "mongodb.database.deprovisioned"
], ],
"listens": [ "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. // Database and owning user share the consumer's login, so the consumer owns exactly its own.
const database = p.as; const database = p.as;
await mongo.createDatabaseAndUser(database, p.as, p.password); await mongo.createDatabaseAndUser(database, p.as, p.password);
await announce("module.mongodb.database.provisioned", { await announce("database.provisioned", {
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
database, database,
user: p.as, user: p.as,
@@ -43,7 +43,7 @@ runProvisioner("mongodb-database", {
async remove(p: { as: string }): Promise<void> { async remove(p: { as: string }): Promise<void> {
await mongo.dropDatabaseAndUser(p.as, p.as); 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 // Asked every minute by the harness: whether the backend still holds this consumer exactly as
// the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120).
+2 -2
View File
@@ -14,11 +14,11 @@ interface TopicEvent {
topicPrefix?: string; 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})`); 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})`); console.log(`[mosquitto] topic deprovisioned for ${e.body.consumer} (client ${e.body.username})`);
}); });
+4 -4
View File
@@ -12,12 +12,12 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.mosquitto.topic.provisioned", "topic.provisioned",
"module.mosquitto.topic.deprovisioned" "topic.deprovisioned"
], ],
"consumes": [ "consumes": [
"module.mosquitto.topic.provisioned", "mosquitto.topic.provisioned",
"module.mosquitto.topic.deprovisioned" "mosquitto.topic.deprovisioned"
], ],
"serves": { "serves": {
"mqtt-topic": {} "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. // The topic subtree is scoped to the consumer's own login, so one cannot read another's topics.
const topicPrefix = p.as; const topicPrefix = p.as;
await mosquitto.createScopedClient(p.as, p.password, topicPrefix); await mosquitto.createScopedClient(p.as, p.password, topicPrefix);
await announce("module.mosquitto.topic.provisioned", { await announce("topic.provisioned", {
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
username: p.as, username: p.as,
topicPrefix, topicPrefix,
@@ -41,7 +41,7 @@ runProvisioner("mqtt-topic", {
async remove(p: { as: string }): Promise<void> { async remove(p: { as: string }): Promise<void> {
await mosquitto.deleteScopedClient(p.as); 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 // Asked every minute by the harness: whether the backend still holds this consumer exactly as
// the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120).
+2 -2
View File
@@ -14,11 +14,11 @@ interface DatabaseEvent {
user?: string; 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})`); 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})`); console.log(`[mssql] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
}); });
+4 -4
View File
@@ -11,12 +11,12 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.mssql.database.provisioned", "database.provisioned",
"module.mssql.database.deprovisioned" "database.deprovisioned"
], ],
"consumes": [ "consumes": [
"module.mssql.database.provisioned", "mssql.database.provisioned",
"module.mssql.database.deprovisioned" "mssql.database.deprovisioned"
], ],
"listens": [ "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. // Database, login and user share the consumer's name, so the consumer owns exactly its own.
const database = p.as; const database = p.as;
await mssql.createDatabaseAndLogin(database, p.as, p.password); await mssql.createDatabaseAndLogin(database, p.as, p.password);
await announce("module.mssql.database.provisioned", { await announce("database.provisioned", {
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
database, database,
user: p.as, user: p.as,
@@ -42,7 +42,7 @@ runProvisioner("mssql-database", {
async remove(p: { as: string }): Promise<void> { async remove(p: { as: string }): Promise<void> {
await mssql.dropDatabaseAndLogin(p.as, p.as); 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 // Asked every minute by the harness: whether the backend still holds this consumer exactly as
// the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120).
+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(); const users = client.listUsers();
for (const u of users) { for (const u of users) {
if (knownUsers.has(u.uid)) continue; 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); knownUsers.add(u.uid);
} }
usersPrimed = true; usersPrimed = true;
@@ -38,7 +38,7 @@ async function pollShares(client: NextcloudClient): Promise<void> {
const shares = await client.listShares(); const shares = await client.listShares();
for (const s of shares) { for (const s of shares) {
if (knownShares.has(s.id)) continue; 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); knownShares.add(s.id);
} }
sharesPrimed = true; sharesPrimed = true;
+2 -2
View File
@@ -26,8 +26,8 @@
"s3-bucket": "${dir:state}/store.secret" "s3-bucket": "${dir:state}/store.secret"
}, },
"emits": [ "emits": [
"module.nextcloud.user.created", "user.created",
"module.nextcloud.share.created" "share.created"
], ],
"own-secrets": { "own-secrets": {
"admin": "${dir:state}/admin.secret", "admin": "${dir:state}/admin.secret",
+1 -1
View File
@@ -2,7 +2,7 @@
"module": "nodered", "module": "nodered",
"version": "1", "version": "1",
"emits": [ "emits": [
"module.nodered.flows.deployed" "flows.deployed"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/nodered/broker" "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); const result = await nodered.deployFlows(flows, type);
// Best-effort announcement — a deploy must not fail because the broker is unbound here. // Best-effort announcement — a deploy must not fail because the broker is unbound here.
try { 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) { } catch (err) {
console.error(`[nodered] deployed but could not emit: ${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) { if (queuePrimed) {
for (const item of items) { for (const item of items) {
if (!inQueue.has(item.id)) { 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 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". // a manual delete lands in history too, and neither is a "download.completed".
if (historyPrimed && item.success) { 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); seenHistory.add(item.id);
} }
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.nzbget.download.added", "download.added",
"module.nzbget.download.completed" "download.completed"
], ],
"consumes": [], "consumes": [],
"own-secrets": { "own-secrets": {
+2 -2
View File
@@ -31,7 +31,7 @@ async function pollRequests(): Promise<void> {
const key = keyOf(r); const key = keyOf(r);
const known = approvedState.has(key); const known = approvedState.has(key);
if (primed && !known) { if (primed && !known) {
await emit("module.ombi.request.created", { await emit("request.created", {
kind: r.kind, kind: r.kind,
id: r.id, id: r.id,
title: r.title, 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. // Approval: the flag went from false to true for a request we already knew about.
if (primed && known && r.approved && approvedState.get(key) === false) { 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); approvedState.set(key, r.approved);
} }
+2 -2
View File
@@ -5,8 +5,8 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.ombi.request.created", "request.created",
"module.ombi.request.approved" "request.approved"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/ombi/broker", "broker": "/var/lib/mesh/ombi/broker",
+1 -1
View File
@@ -20,7 +20,7 @@ async function pollRecent(): Promise<void> {
for (const asset of items) { for (const asset of items) {
if (!seen.has(asset.id)) { if (!seen.has(asset.id)) {
if (primed) { if (primed) {
await emit("module.photos.item.added", { await emit("item.added", {
id: asset.id, id: asset.id,
fileName: asset.fileName, fileName: asset.fileName,
kind: asset.type, kind: asset.type,
+3
View File
@@ -4,6 +4,9 @@
"capabilities": [ "capabilities": [
"container-runtime" "container-runtime"
], ],
"emits": [
"item.added"
],
"requires": [ "requires": [
"s3-bucket", "s3-bucket",
"mongodb-database", "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])); const now = new Map(sessions.map((s) => [s.key, s]));
if (playbackPrimed) { if (playbackPrimed) {
for (const [key, s] of now) { 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) { 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(); active.clear();
@@ -44,7 +44,7 @@ async function pollRecent(): Promise<void> {
for (const item of items) { for (const item of items) {
const id = `${item.title}@${item.addedAt ?? ""}`; const id = `${item.title}@${item.addedAt ?? ""}`;
if (!seen.has(id)) { if (!seen.has(id)) {
if (itemsPrimed) await emit("module.plex.item.added", item); if (itemsPrimed) await emit("item.added", item);
seen.add(id); 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 // 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. // 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(); await plex.refreshAll();
}); });
+4 -4
View File
@@ -5,12 +5,12 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.plex.playback.started", "playback.started",
"module.plex.playback.stopped", "playback.stopped",
"module.plex.item.added" "item.added"
], ],
"consumes": [ "consumes": [
"module.*.download.completed" "*.download.completed"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/plex/broker", "broker": "/var/lib/mesh/plex/broker",
+2 -2
View File
@@ -14,11 +14,11 @@ interface DatabaseEvent {
user?: string; 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})`); 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})`); console.log(`[postgres] database deprovisioned for ${e.body.consumer} (db ${e.body.database})`);
}); });
+4 -4
View File
@@ -17,12 +17,12 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.postgres.database.provisioned", "database.provisioned",
"module.postgres.database.deprovisioned" "database.deprovisioned"
], ],
"consumes": [ "consumes": [
"module.postgres.database.provisioned", "postgres.database.provisioned",
"module.postgres.database.deprovisioned" "postgres.database.deprovisioned"
], ],
"listens": [ "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. // Database and owning role share the consumer's login, so the consumer owns exactly its own.
const database = p.as; const database = p.as;
await postgres.createDatabaseAndRole(database, p.as, p.password); await postgres.createDatabaseAndRole(database, p.as, p.password);
await announce("module.postgres.database.provisioned", { await announce("database.provisioned", {
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
database, database,
user: p.as, user: p.as,
@@ -43,7 +43,7 @@ runProvisioner("postgres-database", {
async remove(p: { as: string }): Promise<void> { async remove(p: { as: string }): Promise<void> {
await postgres.dropDatabaseAndRole(p.as, p.as); 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 // Asked every minute by the harness: whether the backend still holds this consumer exactly as
// the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120). // the mesh gave it, so a login lost behind the provisioner's back is made again (novox/hq issue 120).
+2 -2
View File
@@ -30,9 +30,9 @@ async function pollTorrents(): Promise<void> {
for (const [hash, t] of now) { for (const [hash, t] of now) {
const before = progressByHash.get(hash); const before = progressByHash.get(hash);
if (before === undefined) { 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) { } 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" "container-runtime"
], ],
"emits": [ "emits": [
"module.qbittorrent.download.added", "download.added",
"module.qbittorrent.download.completed" "download.completed"
], ],
"consumes": [], "consumes": [],
"own-secrets": { "own-secrets": {
+2 -2
View File
@@ -40,12 +40,12 @@ async function pollQueue(radarr: RadarrClient): Promise<void> {
if (primed) { if (primed) {
// Entered the queue since last look — Radarr grabbed a release. // Entered the queue since last look — Radarr grabbed a release.
for (const [id, item] of now) { 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. // Left the queue — imported and done, unless it was last seen failing.
for (const [id, item] of inQueue) { for (const [id, item] of inQueue) {
if (!now.has(id) && !FAILED_STATUSES.has(item.status)) { 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" "container-runtime"
], ],
"emits": [ "emits": [
"module.radarr.movie.grabbed", "movie.grabbed",
"module.radarr.download.completed" "download.completed"
], ],
"consumes": [], "consumes": [],
"own-secrets": { "own-secrets": {
+2 -2
View File
@@ -14,11 +14,11 @@ interface CacheEvent {
keyspacePrefix?: string; 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})`); 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})`); console.log(`[redis] cache deprovisioned for ${e.body.consumer} (user ${e.body.username})`);
}); });
+4 -4
View File
@@ -14,12 +14,12 @@
"container-runtime" "container-runtime"
], ],
"emits": [ "emits": [
"module.redis.cache.provisioned", "cache.provisioned",
"module.redis.cache.deprovisioned" "cache.deprovisioned"
], ],
"consumes": [ "consumes": [
"module.redis.cache.provisioned", "redis.cache.provisioned",
"module.redis.cache.deprovisioned" "redis.cache.deprovisioned"
], ],
"serves": { "serves": {
"redis-cache": { "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. // The keyspace is scoped to the consumer's own login, so one cannot read another's keys.
const keyspacePrefix = p.as; const keyspacePrefix = p.as;
await redis.createAclUser(p.as, p.password, keyspacePrefix); await redis.createAclUser(p.as, p.password, keyspacePrefix);
await announce("module.redis.cache.provisioned", { await announce("cache.provisioned", {
consumer: p.consumer ?? "", consumer: p.consumer ?? "",
username: p.as, username: p.as,
keyspacePrefix, keyspacePrefix,
@@ -41,7 +41,7 @@ runProvisioner("redis-cache", {
async remove(p: { as: string }): Promise<void> { async remove(p: { as: string }): Promise<void> {
await redis.deleteAclUser(p.as); 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 // This server keeps its ACL users in memory only, so a restart of it forgets every consumer while
+2 -2
View File
@@ -5,8 +5,8 @@
// It is here so the module exercises the shape rather than describing it. // It is here so the module exercises the shape rather than describing it.
import { on, emit } from "@novox/mesh-sdk/events"; 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"}`); 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. // 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": [ "emits": [
"module.showcase.acknowledged" "greeted",
"acknowledged"
], ],
"consumes": [ "consumes": [
"module.showcase.greeted" "showcase.greeted"
], ],
"listens": [ "listens": [
{ {
@@ -172,7 +173,7 @@
{ {
"id": "tools", "id": "tools",
"type": "container", "type": "container",
"name": "mesh-showcase", "name": "the-showcase",
"artifact": "helper", "artifact": "helper",
"network": "showcase", "network": "showcase",
"volumes": [ "volumes": [
+2 -2
View File
@@ -40,12 +40,12 @@ async function pollQueue(sonarr: SonarrClient): Promise<void> {
if (primed) { if (primed) {
// Entered the queue since last look — Sonarr grabbed a release. // Entered the queue since last look — Sonarr grabbed a release.
for (const [id, item] of now) { 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. // Left the queue — imported and done, unless it was last seen failing.
for (const [id, item] of inQueue) { for (const [id, item] of inQueue) {
if (!now.has(id) && !FAILED_STATUSES.has(item.status)) { 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" "container-runtime"
], ],
"emits": [ "emits": [
"module.sonarr.episode.grabbed", "episode.grabbed",
"module.sonarr.download.completed" "download.completed"
], ],
"consumes": [], "consumes": [],
"own-secrets": { "own-secrets": {
+1 -1
View File
@@ -33,7 +33,7 @@ async function pollHistory(client: TautulliClient): Promise<void> {
} }
async function emitWatch(w: TautulliWatch): 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, title: w.title, user: w.user, mediaType: w.mediaType,
watchedStatus: w.watchedStatus, percentComplete: w.percentComplete, at: w.date, watchedStatus: w.watchedStatus, percentComplete: w.percentComplete, at: w.date,
}); });
+1 -1
View File
@@ -2,7 +2,7 @@
"module": "tautulli", "module": "tautulli",
"version": "1", "version": "1",
"emits": [ "emits": [
"module.tautulli.watch.recorded" "watch.recorded"
], ],
"own-secrets": { "own-secrets": {
"broker": "/var/lib/mesh/tautulli/broker" "broker": "/var/lib/mesh/tautulli/broker"