// 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(); let queuePrimed = false; async function pollQueue(): Promise { 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("module.nzbget.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(); let historyPrimed = false; async function pollHistory(): Promise { 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("module.nzbget.download.completed", { name: item.name, category: item.category, sizeMB: item.sizeMB }); } seenHistory.add(item.id); } } historyPrimed = true; } const tick = (fn: () => Promise, 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");