Files
mesh-catalog/modules/nzbget/index.ts
T

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("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<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("module.nzbget.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");