nzbget, qbittorrent: full nox downloader modules (ADR 0044/0046)
nzbget (usenet, JSON-RPC) and qbittorrent (torrents, WebUI API with manual SID session). Tools: status, queue/torrents, add, pause/resume, delete. Both poll and emit module.<app>.download.added and module.<app>.download.completed — the key plex consumes to rescan; completion is keyed off real success (nzbget history SUCCESS, qbittorrent progress reaching 1), not mere queue disappearance, so a failed or deleted item is not reported as done. Typecheck; manifests parse.
This commit is contained in:
@@ -0,0 +1,160 @@
|
||||
// The NZBGet API client — nzbget's own code, living in the module (novox/hq ADR 0044). Ported from
|
||||
// hal's shared nzbget tools, but self-contained: a change to NZBGet's JSON-RPC now rebuilds only
|
||||
// nzbget and nothing else. Both this module's tools and its events entrypoint import it, and
|
||||
// nothing outside nzbget does. NZBGet speaks JSON-RPC at /jsonrpc, behind HTTP Basic auth.
|
||||
|
||||
export interface NzbgetStatus {
|
||||
/** Bytes/sec — NZBGet reports it split across two 32-bit halves, rejoined here. */
|
||||
speedBytesPerSec: number;
|
||||
remainingMB: number;
|
||||
downloadedTodayMB: number;
|
||||
downloadedMonthMB: number;
|
||||
freeDiskMB: number;
|
||||
paused: boolean;
|
||||
postJobs: number;
|
||||
uptimeSec: number;
|
||||
}
|
||||
|
||||
export interface NzbgetQueueItem {
|
||||
/** The NZBID — stable while the item is queued, so events can diff on it. */
|
||||
id: number;
|
||||
name: string;
|
||||
status: string;
|
||||
category: string;
|
||||
sizeMB: number;
|
||||
remainingMB: number;
|
||||
percent: number;
|
||||
}
|
||||
|
||||
export interface NzbgetHistoryItem {
|
||||
/** The NZBID — the same id the item carried in the queue. */
|
||||
id: number;
|
||||
name: string;
|
||||
/** NZBGet's own status string, e.g. "SUCCESS/ALL", "FAILURE/PAR", "DELETED/MANUAL". */
|
||||
status: string;
|
||||
category: string;
|
||||
sizeMB: number;
|
||||
/** A genuine completion (status starts "SUCCESS") vs a failed or deleted entry — the difference
|
||||
* between something to announce as done and something that merely left the queue. */
|
||||
success: boolean;
|
||||
}
|
||||
|
||||
export class NzbgetClient {
|
||||
readonly rpcUrl: string;
|
||||
private readonly auth: string;
|
||||
|
||||
constructor(url: string, user: string, password: string) {
|
||||
this.rpcUrl = `${url.replace(/\/$/, "")}/jsonrpc`;
|
||||
this.auth = Buffer.from(`${user}:${password}`).toString("base64");
|
||||
}
|
||||
|
||||
/**
|
||||
* Build from the module's resolved environment. URL and password are read from MESH_NZBGET_URL
|
||||
* and MESH_NZBGET_PASSWORD; both must be present — an unconfigured NZBGet throws rather than
|
||||
* pretend to be reachable, so the tools/events simply do not load (the harness treats the throw
|
||||
* as "exposes nothing"). The control username defaults to "nzbget", NZBGet's own default.
|
||||
*/
|
||||
static fromEnv(env: NodeJS.ProcessEnv = process.env): NzbgetClient {
|
||||
const url = env.MESH_NZBGET_URL;
|
||||
const password = env.MESH_NZBGET_PASSWORD;
|
||||
if (!url || !password) {
|
||||
throw new Error("NZBGet not configured — set MESH_NZBGET_URL and MESH_NZBGET_PASSWORD");
|
||||
}
|
||||
const user = env.MESH_NZBGET_USER ?? "nzbget";
|
||||
return new NzbgetClient(url, user, password);
|
||||
}
|
||||
|
||||
private async rpc<T>(method: string, params: unknown[] = []): Promise<T> {
|
||||
const res = await fetch(this.rpcUrl, {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/json", Authorization: `Basic ${this.auth}` },
|
||||
body: JSON.stringify({ method, params, id: 1 }),
|
||||
});
|
||||
if (!res.ok) throw new Error(`NZBGet API ${method}: ${res.status} ${await res.text()}`);
|
||||
const data = (await res.json()) as { result?: T; error?: unknown };
|
||||
if (data.error) throw new Error(`NZBGet RPC ${method}: ${JSON.stringify(data.error)}`);
|
||||
return data.result as T;
|
||||
}
|
||||
|
||||
async getVersion(): Promise<string> {
|
||||
return this.rpc<string>("version");
|
||||
}
|
||||
|
||||
async getStatus(): Promise<NzbgetStatus> {
|
||||
const s = await this.rpc<Record<string, number | boolean>>("status");
|
||||
const lo = Number(s.DownloadRateLo ?? 0);
|
||||
const hi = Number(s.DownloadRateHi ?? 0);
|
||||
return {
|
||||
speedBytesPerSec: lo + hi * 4294967296,
|
||||
remainingMB: Number(s.RemainingSizeMB ?? 0),
|
||||
downloadedTodayMB: Number(s.DaySizeMB ?? 0),
|
||||
downloadedMonthMB: Number(s.MonthSizeMB ?? 0),
|
||||
freeDiskMB: Number(s.FreeDiskSpaceMB ?? 0),
|
||||
paused: Boolean(s.DownloadPaused),
|
||||
postJobs: Number(s.PostJobCount ?? 0),
|
||||
uptimeSec: Number(s.UpTimeSec ?? 0),
|
||||
};
|
||||
}
|
||||
|
||||
async getQueue(): Promise<NzbgetQueueItem[]> {
|
||||
const groups = await this.rpc<Record<string, any>[]>("listgroups", [0]);
|
||||
return groups.map((g) => {
|
||||
const size = Number(g.FileSizeMB ?? 0);
|
||||
const remaining = Number(g.RemainingSizeMB ?? 0);
|
||||
return {
|
||||
id: Number(g.NZBID),
|
||||
name: String(g.NZBName ?? "Unknown"),
|
||||
status: String(g.Status ?? "unknown"),
|
||||
category: String(g.Category ?? ""),
|
||||
sizeMB: size,
|
||||
remainingMB: remaining,
|
||||
percent: size > 0 ? Math.round(((size - remaining) / size) * 100) : 0,
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
async getHistory(limit = 20): Promise<NzbgetHistoryItem[]> {
|
||||
const history = await this.rpc<Record<string, any>[]>("history", [false]);
|
||||
return history.slice(0, limit).map((h) => {
|
||||
const status = String(h.Status ?? "");
|
||||
return {
|
||||
id: Number(h.NZBID),
|
||||
name: String(h.Name ?? "Unknown"),
|
||||
status,
|
||||
category: String(h.Category ?? ""),
|
||||
sizeMB: Number(h.FileSizeMB ?? 0),
|
||||
success: status.startsWith("SUCCESS"),
|
||||
};
|
||||
});
|
||||
}
|
||||
|
||||
/** Queue an NZB by URL. Returns the new NZBID; a non-positive id means NZBGet refused it. */
|
||||
async add(url: string, category = "", priority = 0, paused = false): Promise<number> {
|
||||
const id = await this.rpc<number>("append", [
|
||||
"", url, category, priority, false, paused, "", 0, "SCORE", false, [],
|
||||
]);
|
||||
if (!id || id <= 0) throw new Error("NZBGet refused the NZB (append returned 0)");
|
||||
return id;
|
||||
}
|
||||
|
||||
async pauseAll(): Promise<void> {
|
||||
await this.rpc("pausedownload");
|
||||
}
|
||||
|
||||
async resumeAll(): Promise<void> {
|
||||
await this.rpc("resumedownload");
|
||||
}
|
||||
|
||||
async pauseItem(id: number): Promise<void> {
|
||||
await this.rpc("editqueue", ["GroupPause", "", [id]]);
|
||||
}
|
||||
|
||||
async resumeItem(id: number): Promise<void> {
|
||||
await this.rpc("editqueue", ["GroupResume", "", [id]]);
|
||||
}
|
||||
|
||||
/** Delete an item from the queue or from history. */
|
||||
async delete(id: number, from: "queue" | "history" = "queue"): Promise<void> {
|
||||
await this.rpc("editqueue", [from === "history" ? "HistoryDelete" : "GroupDelete", "", [id]]);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,64 @@
|
||||
// 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 0046/0047):
|
||||
// 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");
|
||||
@@ -4,6 +4,14 @@
|
||||
"capabilities": [
|
||||
"container-runtime"
|
||||
],
|
||||
"emits": [
|
||||
"module.nzbget.download.added",
|
||||
"module.nzbget.download.completed"
|
||||
],
|
||||
"consumes": [],
|
||||
"own-secrets": {
|
||||
"broker": "/var/lib/nzbget/broker"
|
||||
},
|
||||
"listens": [
|
||||
{
|
||||
"port": 6789,
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"name": "@novox/module-nzbget",
|
||||
"version": "0.1.0",
|
||||
"description": "nzbget — Usenet download client. Its API client, tools and events live here (novox/hq ADR 0044).",
|
||||
"type": "module",
|
||||
"private": true,
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.0.0",
|
||||
"typescript": "^5.6.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
// nzbget's tools — ported from the shared hal sdk (novox/hq ADR 0044), importing nzbget's own
|
||||
// client. They return structured data (not the pre-formatted text hal returned); the mesh serves
|
||||
// them through the sdk's tool harness.
|
||||
|
||||
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { NzbgetClient } from "../client.js";
|
||||
|
||||
export function getNzbgetTools(nzbget: NzbgetClient): ToolDefinition[] {
|
||||
return [
|
||||
{
|
||||
name: "nzbget_status",
|
||||
description: "NZBGet server status: download speed, queue remaining, disk free, paused state.",
|
||||
input: {},
|
||||
run: async () => nzbget.getStatus(),
|
||||
},
|
||||
{
|
||||
name: "nzbget_queue",
|
||||
description: "List the current NZBGet download queue — what is downloading and how far along.",
|
||||
input: {},
|
||||
run: async () => {
|
||||
const items = await nzbget.getQueue();
|
||||
return { count: items.length, items };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "nzbget_history",
|
||||
description: "Recent NZBGet download history, newest first — completed, failed and deleted items.",
|
||||
input: { limit: { type: "number", description: "how many entries (default 20)" } },
|
||||
run: async (args) => {
|
||||
const items = await nzbget.getHistory(args.limit ? Number(args.limit) : 20);
|
||||
return { count: items.length, items };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "nzbget_add",
|
||||
description: "Queue an NZB download by URL, optionally into a category.",
|
||||
input: {
|
||||
url: { type: "string", description: "URL to the NZB file" },
|
||||
category: { type: "string", description: "category name (determines download directory)" },
|
||||
priority: { type: "number", description: "-100 very low … 0 normal … 100 very high (default 0)" },
|
||||
paused: { type: "boolean", description: "add in paused state (default false)" },
|
||||
},
|
||||
run: async (args) => {
|
||||
const id = await nzbget.add(
|
||||
String(args.url),
|
||||
args.category ? String(args.category) : "",
|
||||
args.priority ? Number(args.priority) : 0,
|
||||
args.paused === true || args.paused === "true",
|
||||
);
|
||||
return { added: id, category: args.category ? String(args.category) : null };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "nzbget_pause",
|
||||
description: "Pause or resume all NZBGet downloads.",
|
||||
input: { resume: { type: "boolean", description: "true to resume, false to pause (default false)" } },
|
||||
run: async (args) => {
|
||||
const resume = args.resume === true || args.resume === "true";
|
||||
if (resume) await nzbget.resumeAll();
|
||||
else await nzbget.pauseAll();
|
||||
return { paused: !resume };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "nzbget_delete",
|
||||
description: "Delete an NZB from the queue or from history by its NZBID.",
|
||||
input: {
|
||||
id: { type: "number", description: "the NZBID to delete" },
|
||||
from: { type: "string", description: "'queue' (default) or 'history'" },
|
||||
},
|
||||
run: async (args) => {
|
||||
const from = args.from === "history" ? "history" : "queue";
|
||||
await nzbget.delete(Number(args.id), from);
|
||||
return { deleted: Number(args.id), from };
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
// The tools exist only when NZBGet is configured; without a URL and password, nzbget contributes
|
||||
// none rather than failing the whole runtime.
|
||||
registerModuleTools("nzbget", (env) => {
|
||||
try {
|
||||
return getNzbgetTools(NzbgetClient.fromEnv(env));
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
});
|
||||
@@ -0,0 +1,12 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ES2022",
|
||||
"module": "NodeNext",
|
||||
"moduleResolution": "NodeNext",
|
||||
"strict": true,
|
||||
"esModuleInterop": true,
|
||||
"skipLibCheck": true,
|
||||
"noEmit": true
|
||||
},
|
||||
"include": ["client.ts", "index.ts", "tools/index.ts"]
|
||||
}
|
||||
@@ -0,0 +1,170 @@
|
||||
// The qBittorrent API client — qbittorrent's own code, living in the module (novox/hq ADR 0044).
|
||||
// Written against the WebUI API (/api/v2/...), self-contained so a change to it rebuilds only
|
||||
// qbittorrent. Both this module's tools and its events entrypoint import it, and nothing outside
|
||||
// qbittorrent does.
|
||||
//
|
||||
// The WebUI authenticates with a session cookie (SID) obtained by POSTing credentials, and guards
|
||||
// against CSRF by checking the Referer header. Node's fetch keeps no cookie jar, so the SID is
|
||||
// captured on login and carried by hand on every later call, with a single re-login on expiry.
|
||||
|
||||
export interface QbTransferInfo {
|
||||
dlSpeedBytesPerSec: number;
|
||||
upSpeedBytesPerSec: number;
|
||||
dlData: number;
|
||||
upData: number;
|
||||
connectionStatus: string;
|
||||
}
|
||||
|
||||
export interface QbTorrent {
|
||||
hash: string;
|
||||
name: string;
|
||||
/** qBittorrent's state, e.g. "downloading", "stalledUP", "uploading", "pausedUP", "error". */
|
||||
state: string;
|
||||
/** 0..1 — 1 means the download is complete. */
|
||||
progress: number;
|
||||
sizeBytes: number;
|
||||
dlSpeed: number;
|
||||
upSpeed: number;
|
||||
category: string;
|
||||
ratio: number;
|
||||
savePath: string;
|
||||
}
|
||||
|
||||
export class QbittorrentClient {
|
||||
readonly baseUrl: string;
|
||||
private sid: string | null = null;
|
||||
|
||||
constructor(
|
||||
baseUrl: string,
|
||||
private readonly user: string,
|
||||
private readonly password: string,
|
||||
) {
|
||||
this.baseUrl = baseUrl.replace(/\/$/, "");
|
||||
}
|
||||
|
||||
/**
|
||||
* Build from the module's resolved environment. URL and password are read from
|
||||
* MESH_QBITTORRENT_URL and MESH_QBITTORRENT_PASSWORD; both must be present — an unconfigured
|
||||
* qBittorrent throws rather than pretend to be reachable, so the tools/events simply do not load
|
||||
* (the harness treats the throw as "exposes nothing"). The user defaults to "admin".
|
||||
*/
|
||||
static fromEnv(env: NodeJS.ProcessEnv = process.env): QbittorrentClient {
|
||||
const url = env.MESH_QBITTORRENT_URL;
|
||||
const password = env.MESH_QBITTORRENT_PASSWORD;
|
||||
if (!url || !password) {
|
||||
throw new Error("qBittorrent not configured — set MESH_QBITTORRENT_URL and MESH_QBITTORRENT_PASSWORD");
|
||||
}
|
||||
const user = env.MESH_QBITTORRENT_USER ?? "admin";
|
||||
return new QbittorrentClient(url, user, password);
|
||||
}
|
||||
|
||||
private async login(): Promise<void> {
|
||||
const res = await fetch(`${this.baseUrl}/api/v2/auth/login`, {
|
||||
method: "POST",
|
||||
headers: { "Content-Type": "application/x-www-form-urlencoded", Referer: this.baseUrl },
|
||||
body: new URLSearchParams({ username: this.user, password: this.password }),
|
||||
});
|
||||
if (!res.ok) throw new Error(`qBittorrent login: ${res.status} ${await res.text()}`);
|
||||
if ((await res.text()).trim() !== "Ok.") {
|
||||
throw new Error("qBittorrent login rejected — check credentials");
|
||||
}
|
||||
const match = res.headers.get("set-cookie")?.match(/SID=([^;]+)/);
|
||||
if (!match) throw new Error("qBittorrent login returned no SID cookie");
|
||||
this.sid = match[1];
|
||||
}
|
||||
|
||||
private async call(method: "GET" | "POST", path: string, form?: Record<string, string>): Promise<Response> {
|
||||
if (!this.sid) await this.login();
|
||||
const doFetch = (): Promise<Response> => {
|
||||
const headers: Record<string, string> = { Referer: this.baseUrl, Cookie: `SID=${this.sid}` };
|
||||
const init: RequestInit = { method, headers };
|
||||
if (form) {
|
||||
headers["Content-Type"] = "application/x-www-form-urlencoded";
|
||||
init.body = new URLSearchParams(form);
|
||||
}
|
||||
return fetch(`${this.baseUrl}/api/v2/${path}`, init);
|
||||
};
|
||||
let res = await doFetch();
|
||||
if (res.status === 403) {
|
||||
// The SID expired — re-authenticate once and retry, rather than fail a routine call.
|
||||
await this.login();
|
||||
res = await doFetch();
|
||||
}
|
||||
return res;
|
||||
}
|
||||
|
||||
private async getJson<T>(path: string): Promise<T> {
|
||||
const res = await this.call("GET", path);
|
||||
if (!res.ok) throw new Error(`qBittorrent GET ${path}: ${res.status} ${await res.text()}`);
|
||||
return (await res.json()) as T;
|
||||
}
|
||||
|
||||
async getVersion(): Promise<string> {
|
||||
const res = await this.call("GET", "app/version");
|
||||
if (!res.ok) throw new Error(`qBittorrent app/version: ${res.status}`);
|
||||
return (await res.text()).trim();
|
||||
}
|
||||
|
||||
async getTransferInfo(): Promise<QbTransferInfo> {
|
||||
const d = await this.getJson<Record<string, any>>("transfer/info");
|
||||
return {
|
||||
dlSpeedBytesPerSec: Number(d.dl_info_speed ?? 0),
|
||||
upSpeedBytesPerSec: Number(d.up_info_speed ?? 0),
|
||||
dlData: Number(d.dl_info_data ?? 0),
|
||||
upData: Number(d.up_info_data ?? 0),
|
||||
connectionStatus: String(d.connection_status ?? "unknown"),
|
||||
};
|
||||
}
|
||||
|
||||
async getTorrents(filter?: string): Promise<QbTorrent[]> {
|
||||
const path = filter ? `torrents/info?filter=${encodeURIComponent(filter)}` : "torrents/info";
|
||||
const list = await this.getJson<Record<string, any>[]>(path);
|
||||
return list.map((t) => ({
|
||||
hash: String(t.hash),
|
||||
name: String(t.name ?? "Unknown"),
|
||||
state: String(t.state ?? "unknown"),
|
||||
progress: Number(t.progress ?? 0),
|
||||
sizeBytes: Number(t.size ?? 0),
|
||||
dlSpeed: Number(t.dlspeed ?? 0),
|
||||
upSpeed: Number(t.upspeed ?? 0),
|
||||
category: String(t.category ?? ""),
|
||||
ratio: Number(t.ratio ?? 0),
|
||||
savePath: String(t.save_path ?? ""),
|
||||
}));
|
||||
}
|
||||
|
||||
/** Add a torrent by magnet or http(s) .torrent URL, optionally into a category / save path. */
|
||||
async add(url: string, category = "", savepath = "", paused = false): Promise<void> {
|
||||
const form: Record<string, string> = { urls: url, paused: paused ? "true" : "false" };
|
||||
if (category) form.category = category;
|
||||
if (savepath) form.savepath = savepath;
|
||||
const res = await this.call("POST", "torrents/add", form);
|
||||
const text = (await res.text()).trim();
|
||||
if (!res.ok || text.toLowerCase() === "fails.") {
|
||||
throw new Error(`qBittorrent refused the torrent: ${res.status} ${text}`);
|
||||
}
|
||||
}
|
||||
|
||||
// qBittorrent 5.x renamed pause/resume to stop/start; try the modern name and fall back to the
|
||||
// legacy one on a 404, so the client works against both.
|
||||
private async command(modern: string, legacy: string, hashes: string): Promise<void> {
|
||||
let res = await this.call("POST", `torrents/${modern}`, { hashes });
|
||||
if (res.status === 404) res = await this.call("POST", `torrents/${legacy}`, { hashes });
|
||||
if (!res.ok) throw new Error(`qBittorrent torrents/${modern}: ${res.status} ${await res.text()}`);
|
||||
}
|
||||
|
||||
/** Pause torrents — a pipe-separated hash list, or "all" (the default). */
|
||||
async pause(hashes = "all"): Promise<void> {
|
||||
await this.command("stop", "pause", hashes);
|
||||
}
|
||||
|
||||
/** Resume torrents — a pipe-separated hash list, or "all" (the default). */
|
||||
async resume(hashes = "all"): Promise<void> {
|
||||
await this.command("start", "resume", hashes);
|
||||
}
|
||||
|
||||
async delete(hashes: string, deleteFiles = false): Promise<void> {
|
||||
const res = await this.call("POST", "torrents/delete", { hashes, deleteFiles: deleteFiles ? "true" : "false" });
|
||||
if (!res.ok) throw new Error(`qBittorrent torrents/delete: ${res.status} ${await res.text()}`);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,52 @@
|
||||
// qbittorrent's events. The tool runtime imports this once the broker is bound. It watches the
|
||||
// torrent list and turns its comings and goings into mesh events.
|
||||
//
|
||||
// Emits (novox/hq ADR 0046/0047):
|
||||
// module.qbittorrent.download.added — a torrent was added
|
||||
// module.qbittorrent.download.completed — a torrent finished downloading (progress reached 1).
|
||||
// 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.
|
||||
//
|
||||
// The torrent list is polled and diffed by hash, primed silently on the first look (like plex's and
|
||||
// sonarr's index.ts) so a restart does not re-announce everything already present. Completion is a
|
||||
// progress crossing from below 1 to exactly 1 — a torrent added already-complete is announced only
|
||||
// as added, never as freshly completed, since nothing was downloaded.
|
||||
|
||||
import { emit } from "@novox/mesh-sdk/events";
|
||||
import { QbittorrentClient } from "./client.js";
|
||||
|
||||
const qb = QbittorrentClient.fromEnv();
|
||||
|
||||
const progressByHash = new Map<string, number>();
|
||||
let primed = false;
|
||||
|
||||
async function pollTorrents(): Promise<void> {
|
||||
const torrents = await qb.getTorrents();
|
||||
const now = new Map(torrents.map((t) => [t.hash, t]));
|
||||
|
||||
if (primed) {
|
||||
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 });
|
||||
} else if (before < 1 && t.progress >= 1) {
|
||||
await emit("module.qbittorrent.download.completed", { name: t.name, category: t.category, sizeBytes: t.sizeBytes });
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
progressByHash.clear();
|
||||
for (const [hash, t] of now) progressByHash.set(hash, t.progress);
|
||||
primed = true;
|
||||
}
|
||||
|
||||
const tick = (fn: () => Promise<void>, everyMs: number): void => {
|
||||
const run = (): void => void fn().catch((err) => console.error(`[qbittorrent] ${err}`));
|
||||
setInterval(run, everyMs);
|
||||
run();
|
||||
};
|
||||
tick(pollTorrents, 20_000);
|
||||
|
||||
console.log("[qbittorrent] watching torrents, emitting adds and completions");
|
||||
@@ -4,6 +4,14 @@
|
||||
"capabilities": [
|
||||
"container-runtime"
|
||||
],
|
||||
"emits": [
|
||||
"module.qbittorrent.download.added",
|
||||
"module.qbittorrent.download.completed"
|
||||
],
|
||||
"consumes": [],
|
||||
"own-secrets": {
|
||||
"broker": "/var/lib/qbittorrent/broker"
|
||||
},
|
||||
"listens": [
|
||||
{
|
||||
"port": 8080,
|
||||
|
||||
@@ -0,0 +1,14 @@
|
||||
{
|
||||
"name": "@novox/module-qbittorrent",
|
||||
"version": "0.1.0",
|
||||
"description": "qbittorrent — BitTorrent download client. Its API client, tools and events live here (novox/hq ADR 0044).",
|
||||
"type": "module",
|
||||
"private": true,
|
||||
"dependencies": {
|
||||
"@novox/mesh-sdk": "^0.1.0"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/node": "^22.0.0",
|
||||
"typescript": "^5.6.0"
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,98 @@
|
||||
// qbittorrent's tools — living in the module (novox/hq ADR 0044), importing qbittorrent's own
|
||||
// client. They return structured data; the mesh serves them through the sdk's tool harness.
|
||||
|
||||
import { registerModuleTools, type ToolDefinition } from "@novox/mesh-sdk/tools";
|
||||
import { QbittorrentClient } from "../client.js";
|
||||
|
||||
export function getQbittorrentTools(qb: QbittorrentClient): ToolDefinition[] {
|
||||
return [
|
||||
{
|
||||
name: "qbittorrent_status",
|
||||
description: "qBittorrent status: version, global transfer rates, and how many torrents are active.",
|
||||
input: {},
|
||||
run: async () => {
|
||||
const [version, transfer, torrents] = await Promise.all([
|
||||
qb.getVersion(),
|
||||
qb.getTransferInfo(),
|
||||
qb.getTorrents(),
|
||||
]);
|
||||
const downloading = torrents.filter((t) => t.progress < 1).length;
|
||||
return { version, transfer, torrents: torrents.length, downloading, seeding: torrents.length - downloading };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "qbittorrent_torrents",
|
||||
description: "List torrents — name, state, progress and speed. Optional filter narrows the set.",
|
||||
input: {
|
||||
filter: { type: "string", description: "one of all|downloading|seeding|completed|paused|active|inactive|stalled" },
|
||||
},
|
||||
run: async (args) => {
|
||||
const items = await qb.getTorrents(args.filter ? String(args.filter) : undefined);
|
||||
return { count: items.length, torrents: items };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "qbittorrent_add",
|
||||
description: "Add a torrent by magnet link or .torrent URL, optionally into a category.",
|
||||
input: {
|
||||
url: { type: "string", description: "magnet link or http(s) URL to a .torrent" },
|
||||
category: { type: "string", description: "category name (determines save directory)" },
|
||||
savepath: { type: "string", description: "explicit save path (overrides the category default)" },
|
||||
paused: { type: "boolean", description: "add in paused state (default false)" },
|
||||
},
|
||||
run: async (args) => {
|
||||
await qb.add(
|
||||
String(args.url),
|
||||
args.category ? String(args.category) : "",
|
||||
args.savepath ? String(args.savepath) : "",
|
||||
args.paused === true || args.paused === "true",
|
||||
);
|
||||
return { added: String(args.url), category: args.category ? String(args.category) : null };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "qbittorrent_pause",
|
||||
description: "Pause torrents — a pipe-separated hash list, or 'all' (the default).",
|
||||
input: { hashes: { type: "string", description: "pipe-separated torrent hashes, or 'all' (default)" } },
|
||||
run: async (args) => {
|
||||
const hashes = args.hashes ? String(args.hashes) : "all";
|
||||
await qb.pause(hashes);
|
||||
return { paused: hashes };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "qbittorrent_resume",
|
||||
description: "Resume torrents — a pipe-separated hash list, or 'all' (the default).",
|
||||
input: { hashes: { type: "string", description: "pipe-separated torrent hashes, or 'all' (default)" } },
|
||||
run: async (args) => {
|
||||
const hashes = args.hashes ? String(args.hashes) : "all";
|
||||
await qb.resume(hashes);
|
||||
return { resumed: hashes };
|
||||
},
|
||||
},
|
||||
{
|
||||
name: "qbittorrent_delete",
|
||||
description: "Remove torrents by hash, optionally deleting their files on disk.",
|
||||
input: {
|
||||
hashes: { type: "string", description: "pipe-separated torrent hashes, or 'all'" },
|
||||
deleteFiles: { type: "boolean", description: "also delete downloaded files (default false)" },
|
||||
},
|
||||
run: async (args) => {
|
||||
const hashes = String(args.hashes);
|
||||
const deleteFiles = args.deleteFiles === true || args.deleteFiles === "true";
|
||||
await qb.delete(hashes, deleteFiles);
|
||||
return { deleted: hashes, deleteFiles };
|
||||
},
|
||||
},
|
||||
];
|
||||
}
|
||||
|
||||
// The tools exist only when qBittorrent is configured; without a URL and password, qbittorrent
|
||||
// contributes none rather than failing the whole runtime.
|
||||
registerModuleTools("qbittorrent", (env) => {
|
||||
try {
|
||||
return getQbittorrentTools(QbittorrentClient.fromEnv(env));
|
||||
} catch {
|
||||
return [];
|
||||
}
|
||||
});
|
||||
@@ -0,0 +1,12 @@
|
||||
{
|
||||
"compilerOptions": {
|
||||
"target": "ES2022",
|
||||
"module": "NodeNext",
|
||||
"moduleResolution": "NodeNext",
|
||||
"strict": true,
|
||||
"esModuleInterop": true,
|
||||
"skipLibCheck": true,
|
||||
"noEmit": true
|
||||
},
|
||||
"include": ["client.ts", "index.ts", "tools/index.ts"]
|
||||
}
|
||||
Reference in New Issue
Block a user