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 `#`.
65 lines
2.6 KiB
TypeScript
65 lines
2.6 KiB
TypeScript
// nzbget's events. The tool runtime imports this once the broker is bound. It watches the download
|
|
// queue and the history and turns their comings and goings into mesh events.
|
|
//
|
|
// Emits (novox/hq ADR 0041/0042):
|
|
// module.nzbget.download.added — an NZB entered the queue
|
|
// module.nzbget.download.completed — an NZB finished successfully (left the queue, landed in
|
|
// history as SUCCESS). This exact routing key is what the
|
|
// plex module consumes (module.*.download.completed) to
|
|
// rescan, so the new file becomes a visible item.
|
|
// Consumes: none.
|
|
//
|
|
// Two diffs, each primed silently on the first look (like plex's and sonarr's index.ts) so a
|
|
// restart mid-download does not re-announce everything already in flight or already finished. The
|
|
// queue tells us what was grabbed; history — not the queue's disappearance — tells us what actually
|
|
// succeeded, since a failed or deleted download also leaves the queue.
|
|
|
|
import { emit } from "@novox/mesh-sdk/events";
|
|
import { NzbgetClient } from "./client.js";
|
|
|
|
const nzbget = NzbgetClient.fromEnv();
|
|
|
|
const inQueue = new Set<number>();
|
|
let queuePrimed = false;
|
|
async function pollQueue(): Promise<void> {
|
|
const items = await nzbget.getQueue();
|
|
const now = new Set(items.map((i) => i.id));
|
|
if (queuePrimed) {
|
|
for (const item of items) {
|
|
if (!inQueue.has(item.id)) {
|
|
await emit("download.added", { name: item.name, category: item.category, sizeMB: item.sizeMB });
|
|
}
|
|
}
|
|
}
|
|
inQueue.clear();
|
|
for (const id of now) inQueue.add(id);
|
|
queuePrimed = true;
|
|
}
|
|
|
|
const seenHistory = new Set<number>();
|
|
let historyPrimed = false;
|
|
async function pollHistory(): Promise<void> {
|
|
const items = await nzbget.getHistory(50);
|
|
for (const item of items) {
|
|
if (!seenHistory.has(item.id)) {
|
|
// 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("download.completed", { name: item.name, category: item.category, sizeMB: item.sizeMB });
|
|
}
|
|
seenHistory.add(item.id);
|
|
}
|
|
}
|
|
historyPrimed = true;
|
|
}
|
|
|
|
const tick = (fn: () => Promise<void>, everyMs: number): void => {
|
|
const run = (): void => void fn().catch((err) => console.error(`[nzbget] ${err}`));
|
|
setInterval(run, everyMs);
|
|
run();
|
|
};
|
|
tick(pollQueue, 20_000);
|
|
tick(pollHistory, 30_000);
|
|
|
|
console.log("[nzbget] watching the queue and history, emitting adds and completions");
|