diff --git a/README.md b/README.md index 46c918d..944e09c 100644 --- a/README.md +++ b/README.md @@ -55,7 +55,7 @@ torhunt files range-aware HTTP server for media streaming torhunt attach persistent tmux session for remote SSH usage ``` -Append `--daemon` to run `watch`, `serve`, or `files` as background processes. Run `torhunt --help` for all commands and flags. +Append `--daemon` to run `watch`, `serve`, or `files` as background processes. `torhunt serve` also ships a built-in **web remote**: open `http://127.0.0.1:9161/` in any browser (phone included) to search all indexers, add magnets or info hashes, watch progress live, pause/resume, and manage seeding. Run with `--token` when exposing the port beyond loopback. Run `torhunt --help` for all commands and flags. ## Privacy & security diff --git a/scripts/postbuild.cjs b/scripts/postbuild.cjs index e28f88d..571eb0d 100644 --- a/scripts/postbuild.cjs +++ b/scripts/postbuild.cjs @@ -11,6 +11,12 @@ copyFileSync(src, dest); // The WebRTC fallback stub must ship beside cli.cjs, which resolves it via // __dirname when the node-datachannel binary is unavailable. copyFileSync(resolve(root, 'scripts/webrtc-stub.mjs'), resolve(root, 'dist/webrtc-stub.mjs')); +// The web remote's HTML shell ships beside the bundle; src/daemon/webui.ts +// loads it relative to the bundled entry at runtime. +copyFileSync( + resolve(root, 'src/daemon/assets/ui.html'), + resolve(root, 'dist/ui.html'), +); // On Windows chmod is effectively a no-op, and npm re-applies bin permissions on install anyway, so a failure // here shouldn't fail the build, but warn rather than swallow the error. @@ -20,4 +26,4 @@ try { console.warn('postbuild: could not set executable bit on dist/cli.cjs:', err.message); } -console.log('postbuild: wrote dist/cli.cjs and dist/webrtc-stub.mjs'); +console.log('postbuild: wrote dist/cli.cjs, dist/webrtc-stub.mjs and dist/ui.html'); diff --git a/src/daemon/assets/ui.html b/src/daemon/assets/ui.html new file mode 100644 index 0000000..806d1b4 --- /dev/null +++ b/src/daemon/assets/ui.html @@ -0,0 +1,852 @@ + + + + + +torhunt · remote + + + + + +
+
+ +
torhunt
+ remote +
+
offline
+
+
+ +
+ + +
+
Add by magnet or info hash
+
+ + +
+
Downloads start immediately and resume on their own if the daemon restarts.
+
+ +
+
–Downloading
+
–Total speed
+
–Seeding
+
+ + + +
+
Downloads
+
+
+ + + + +
+ + + +
+ + + + \ No newline at end of file diff --git a/src/daemon/serve.test.ts b/src/daemon/serve.test.ts index 5e0537a..339b338 100644 --- a/src/daemon/serve.test.ts +++ b/src/daemon/serve.test.ts @@ -1,8 +1,11 @@ import { describe, it, expect, vi, beforeEach, afterEach } from "vitest"; import os from "node:os"; import path from "node:path"; +import http from "node:http"; +import { AddressInfo } from "node:net"; +import { EventEmitter } from "node:events"; import { promises as fs } from "node:fs"; -import { handleApi, isAuthorized, extractMagnet, parseControl, applyControl } from "./serve"; +import { handleApi, isAuthorized, extractMagnet, parseControl, applyControl, createServeHandler } from "./serve"; import type { Runtime } from "./runtime"; const HASH = "abcdef0123456789abcdef0123456789abcdef01"; @@ -67,6 +70,14 @@ describe("handleApi", () => { expect(res.body.ok).toBe(true); }); + it("includes the configured TUI theme on /health", async () => { + const res = await handleApi(runtime, "tok", "GET", "/health", undefined, ""); + expect(res.status).toBe(200); + const theme = res.body.theme as { id: string; name: string; colors: Record }; + expect(typeof theme.id).toBe("string"); + expect(typeof theme.colors.accent).toBe("string"); + }); + it("401s a protected route without a token", async () => { const res = await handleApi(runtime, "tok", "POST", "/add", undefined, `{"magnet":"${MAGNET}"}`); expect(res.status).toBe(401); @@ -127,6 +138,186 @@ describe("handleApi", () => { expect(res.body).toMatchObject({ ok: true, action: "pause" }); expect(pause).toHaveBeenCalledWith(HASH); }); + + it("lists history on GET /history", async () => { + const completedAt = Date.now(); + runtime.queue = { + getItems: () => [], + getSeeds: () => [], + getHistory: () => [{ id: HASH, name: "Done", sizeBytes: 1234, completedAt }], + } as unknown as Runtime["queue"]; + const res = await handleApi(runtime, null, "GET", "/history", undefined, ""); + expect(res.status).toBe(200); + expect(res.body).toEqual({ + history: [{ id: HASH, name: "Done", sizeBytes: 1234, completedAt }], + }); + }); + + it("400s /search without a query", async () => { + const res = await handleApi(runtime, null, "GET", "/search", undefined, ""); + expect(res.status).toBe(400); + }); + + it("runs an injected search and validates the category", async () => { + const search = vi.fn().mockResolvedValue({ results: [], failed: [] }); + runtime.queue = { getItems: () => [], getSeeds: () => [] } as unknown as Runtime["queue"]; + + const ok = await handleApi( + runtime, null, "GET", "/search", undefined, "", + new URLSearchParams("q=ubuntu&cat=Movies"), search, + ); + expect(ok.status).toBe(200); + expect(search).toHaveBeenCalledWith("ubuntu", "Movies"); + + const bogus = await handleApi( + runtime, null, "GET", "/search", undefined, "", + new URLSearchParams("q=ubuntu&cat=Bogus"), search, + ); + expect(bogus.status).toBe(200); + expect(search).toHaveBeenLastCalledWith("ubuntu", null); + }); +}); + +describe("createServeHandler (web remote routes)", () => { + let server: http.Server; + + function fakeQueue(overrides: Record = {}): Runtime["queue"] { + const emitter = new EventEmitter(); + return Object.assign(emitter, { + getItems: () => [], + getSeeds: () => [], + getHistory: () => [], + has: () => false, + add: vi.fn(), + ...overrides, + }) as unknown as Runtime["queue"]; + } + + function start( + token: string | null, + queue: Runtime["queue"], + searchFn?: Parameters[3], + ): Promise { + const runtime = { queue, downloadDir: "unused" } as unknown as Runtime; + server = http.createServer(createServeHandler(runtime, token, () => {}, searchFn)); + return new Promise((resolve) => { + server.listen(0, "127.0.0.1", () => + resolve(`http://127.0.0.1:${(server.address() as AddressInfo).port}`), + ); + }); + } + + afterEach(async () => { + if (!server) return; + server.closeAllConnections?.(); + await new Promise((resolve) => server.close(() => resolve())); + server = undefined as unknown as http.Server; + }); + + it("serves the web remote shell on /", async () => { + const base = await start(null, fakeQueue()); + const res = await fetch(`${base}/`); + expect(res.status).toBe(200); + expect(res.headers.get("content-type")).toContain("text/html"); + const html = await res.text(); + expect(html).toContain(" { + const base = await start(null, fakeQueue()); + const res = await fetch(`${base}/ui`); + expect(res.status).toBe(200); + expect(await res.text()).toContain(" { + const base = await start(null, fakeQueue()); + const controller = new AbortController(); + const res = await fetch(`${base}/events`, { signal: controller.signal }); + expect(res.status).toBe(200); + expect(res.headers.get("content-type")).toContain("text/event-stream"); + const reader = res.body!.getReader()!; + const { value } = await reader.read(); + const text = new TextDecoder().decode(value); + expect(text).toContain("retry:"); + expect(text).toContain('"downloads"'); + controller.abort(); + }); + + it("401s /events with a wrong token", async () => { + const base = await start("tok", fakeQueue()); + const res = await fetch(`${base}/events?token=nope`); + expect(res.status).toBe(401); + }); + + it("accepts the correct token via query for /events", async () => { + const base = await start("tok", fakeQueue()); + const controller = new AbortController(); + const res = await fetch(`${base}/events?token=tok`, { signal: controller.signal }); + expect(res.status).toBe(200); + controller.abort(); + }); + + it("rejects cross-site POSTs by origin", async () => { + const base = await start(null, fakeQueue()); + const res = await fetch(`${base}/add`, { + method: "POST", + headers: { Origin: "http://evil.example", "Content-Type": "application/json" }, + body: JSON.stringify({ magnet: MAGNET }), + }); + expect(res.status).toBe(403); + expect(((await res.json()) as { error: string }).error).toContain("cross-origin"); + }); + + it("lets same-origin POSTs through to the API", async () => { + const add = vi.fn(); + const base = await start(null, fakeQueue({ add })); + const res = await fetch(`${base}/add`, { + method: "POST", + headers: { Origin: base, "Content-Type": "application/json" }, + body: JSON.stringify({ magnet: MAGNET }), + }); + expect(res.status).toBe(200); + expect(add).toHaveBeenCalled(); + }); + + it("keeps plain curl POSTs working (no Origin header)", async () => { + const add = vi.fn(); + const base = await start(null, fakeQueue({ add })); + const res = await fetch(`${base}/add`, { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ magnet: MAGNET }), + }); + expect(res.status).toBe(200); + expect(add).toHaveBeenCalled(); + }); + + it("exposes history over HTTP for the Completed tab", async () => { + const completedAt = Date.now(); + const base = await start( + null, + fakeQueue({ getHistory: () => [{ id: HASH, name: "Done", sizeBytes: 99, completedAt }] }), + ); + const res = await fetch(`${base}/history`); + expect(res.status).toBe(200); + expect(((await res.json()) as { history: unknown[] }).history).toHaveLength(1); + }); + + it("wires /search end to end with the query string", async () => { + const searchFn = vi.fn().mockResolvedValue({ + results: [{ infoHash: HASH, name: "Result", sizeBytes: 1, seeders: 2, leechers: 0, source: "yts", magnet: MAGNET }], + failed: ["EZTV"], + }); + const base = await start(null, fakeQueue(), searchFn); + const res = await fetch(`${base}/search?q=film&cat=Movies`); + expect(res.status).toBe(200); + expect(searchFn).toHaveBeenCalledWith("film", "Movies"); + const body = (await res.json()) as { results: unknown[]; failed: string[] }; + expect(body.results).toHaveLength(1); + expect(body.failed).toEqual(["EZTV"]); + }); }); describe("parseControl", () => { diff --git a/src/daemon/serve.ts b/src/daemon/serve.ts index 6e7c828..bf357a0 100644 --- a/src/daemon/serve.ts +++ b/src/daemon/serve.ts @@ -11,6 +11,10 @@ import http from "node:http"; import { startRuntime, addInput, type Runtime } from "./runtime"; import { startSeedReaper } from "./seed-reaper"; import { LOOPBACK_HOSTS, isAuthorized, hostHeaderOk } from "./auth"; +import { loadUiHtml, openEventStream, originAllowed } from "./webui"; +import { searchTorrents } from "./websearch"; +import { loadConfig } from "../config/config"; +import { THEMES } from "../ui/theme"; import { VERSION } from "../version"; export { isAuthorized } from "./auth"; @@ -136,6 +140,11 @@ function statusPayload(runtime: Runtime): Record { progress: it.progress, peers: it.peers, speed: it.speed, + // Extra context for the web remote; additive so existing API consumers + // that only read the original fields keep working untouched. + totalBytes: it.totalBytes, + downloadedBytes: it.downloadedBytes, + eta: it.eta, })); const seeds = runtime.queue.getSeeds().map((s) => ({ id: s.id, @@ -147,7 +156,21 @@ function statusPayload(runtime: Runtime): Record { return { downloads, seeds }; } +// Completed-download archive for the web remote's Completed tab. Read-only and +// capped by whatever the queue keeps (HISTORY_MAX), so the payload stays small. +function historyPayload(runtime: Runtime): Record { + const history = runtime.queue.getHistory().map((h) => ({ + id: h.id, + name: h.name, + sizeBytes: h.sizeBytes, + completedAt: h.completedAt, + })); + return { history }; +} + // Pure request router — no node:http types, so it's trivially testable. +// `query` carries the URL's search params (only /search reads them today) and +// `search` is injectable so tests never touch the network. export async function handleApi( runtime: Runtime, token: string | null, @@ -155,9 +178,23 @@ export async function handleApi( urlPath: string, authHeader: string | undefined, bodyText: string, + query?: URLSearchParams, + search: typeof searchTorrents = searchTorrents, ): Promise { if (method === "GET" && urlPath === "/health") { - return { status: 200, body: { ok: true, version: VERSION } }; + // The configured TUI theme rides along so the web remote can mirror the + // terminal's look. Read fresh from disk on every call: a theme switched in + // the TUI reaches an already-running daemon here without a restart. + const cfg = await loadConfig(); + const theme = THEMES.find((t) => t.id === cfg.theme) ?? THEMES[0]!; + return { + status: 200, + body: { + ok: true, + version: VERSION, + theme: { id: theme.id, name: theme.name, colors: theme.colors }, + }, + }; } if (!isAuthorized(token, authHeader)) { return { status: 401, body: { error: "unauthorized" } }; @@ -165,6 +202,26 @@ export async function handleApi( if (method === "GET" && (urlPath === "/downloads" || urlPath === "/status")) { return { status: 200, body: statusPayload(runtime) }; } + if (method === "GET" && urlPath === "/history") { + return { status: 200, body: historyPayload(runtime) }; + } + if (method === "GET" && urlPath === "/search") { + const q = (query?.get("q") ?? "").trim(); + if (!q) return { status: 400, body: { error: "missing q" } }; + const rawCat = query?.get("cat") ?? ""; + const category = + rawCat === "Games" || rawCat === "Movies" || rawCat === "TV" || rawCat === "Anime" + ? rawCat + : null; + try { + const outcome = await search(q, category); + return { status: 200, body: { results: outcome.results, failed: outcome.failed } }; + } catch { + // searchTorrents already absorbs per-source failures; reaching here means + // something systemic went wrong (e.g. the deadline itself threw). + return { status: 504, body: { error: "search failed" } }; + } + } if (method === "POST" && urlPath === "/add") { const magnet = extractMagnet(bodyText); if (!magnet) return { status: 400, body: { error: "missing magnet or info hash" } }; @@ -220,37 +277,64 @@ function log(message: string): void { console.log(`[torhunt serve] ${new Date().toISOString()} ${message}`); } -export async function runServe(options: ServeOptions = {}): Promise { - const port = options.port ?? DEFAULT_API_PORT; - const host = options.host ?? "127.0.0.1"; - const token = options.token && options.token.trim() ? options.token.trim() : null; - - // Fail soft, not open: never expose a public interface without a token. - if (!LOOPBACK_HOSTS.has(host) && !token) { - console.error( - `error: refusing to bind ${host} without a token. Pass --token ` + - `(or set TORHUNT_API_TOKEN), or bind 127.0.0.1.`, - ); - process.exit(1); - return; - } - - const runtime = await startRuntime(options.downloadDir); - - if (options.seedTimeMs && options.seedTimeMs > 0) { - startSeedReaper(runtime.queue, options.seedTimeMs, { deleteFiles: options.deleteFiles, log }); - } - - const server = http.createServer((req, res) => { +// Build the request handler for the headless add API. Split out from runServe +// so tests can spin a real server against a fake runtime. +export function createServeHandler( + runtime: Runtime, + token: string | null, + logFn: (message: string) => void = log, + searchFn: typeof searchTorrents = searchTorrents, +): (req: http.IncomingMessage, res: http.ServerResponse) => void { + return (req, res) => { void (async () => { const method = req.method ?? "GET"; - const urlPath = (req.url ?? "/").split("?")[0]!; + const url = new URL(req.url ?? "/", "http://localhost"); + const urlPath = url.pathname; // Tokenless means loopback-bound; require a loopback Host so a hostile // webpage can't reach us through DNS rebinding. if (!token && !hostHeaderOk(req.headers.host)) { res.writeHead(403, { "Content-Type": "application/json" }); res.end(JSON.stringify({ error: "forbidden host" })); - log(`${method} ${urlPath} -> 403 (host)`); + logFn(`${method} ${urlPath} -> 403 (host)`); + return; + } + // Web remote shell. Like /health it carries no secrets — data stays + // behind auth below — so any browser that passes the Host check gets it. + if (method === "GET" && (urlPath === "/" || urlPath === "/ui")) { + const html = loadUiHtml(); + if (html === null) { + res.writeHead(404, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ error: "web ui not bundled" })); + return; + } + res.writeHead(200, { "Content-Type": "text/html; charset=utf-8", "Cache-Control": "no-store" }); + res.end(html); + return; + } + // Live snapshot stream. EventSource can't set headers, so the token is + // also accepted as a query parameter here (never on JSON routes). + if (method === "GET" && urlPath === "/events") { + const queryToken = url.searchParams.get("token"); + const auth = + req.headers.authorization ?? (queryToken ? `Bearer ${queryToken}` : undefined); + if (!isAuthorized(token, auth)) { + res.writeHead(401, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ error: "unauthorized" })); + logFn(`GET /events -> 401`); + return; + } + openEventStream(res, runtime.queue, () => statusPayload(runtime)); + logFn("event stream attached"); + return; + } + // Browsers attach Origin to every POST they send; cross-site ones don't + // match our Host. curl and scripts send no Origin and pass untouched. + // This keeps the JSON API un-forgeable from hostile web pages even in + // tokenless mode (where parsers alone already reject form encodings). + if (method !== "GET" && method !== "HEAD" && !originAllowed(req.headers.origin, req.headers.host)) { + res.writeHead(403, { "Content-Type": "application/json" }); + res.end(JSON.stringify({ error: "cross-origin request rejected" })); + logFn(`${method} ${urlPath} -> 403 (origin)`); return; } const body = method === "POST" ? await readBody(req) : { text: "", tooLarge: false }; @@ -258,13 +342,22 @@ export async function runServe(options: ServeOptions = {}): Promise { res.writeHead(413, { "Content-Type": "application/json", Connection: "close" }); res.end(JSON.stringify({ error: "body too large" })); res.once("finish", () => req.destroy()); - log(`${method} ${urlPath} -> 413`); + logFn(`${method} ${urlPath} -> 413`); return; } const bodyText = body.text; let out: ApiResponse; try { - out = await handleApi(runtime, token, method, urlPath, req.headers.authorization, bodyText); + out = await handleApi( + runtime, + token, + method, + urlPath, + req.headers.authorization, + bodyText, + url.searchParams, + searchFn, + ); } catch { out = { status: 500, body: { error: "internal error" } }; } @@ -272,10 +365,34 @@ export async function runServe(options: ServeOptions = {}): Promise { res.writeHead(out.status, { "Content-Type": "application/json" }); res.end(payload); if (method !== "GET" || urlPath !== "/health") { - log(`${method} ${urlPath} -> ${out.status}`); + logFn(`${method} ${urlPath} -> ${out.status}`); } })(); - }); + }; +} + +export async function runServe(options: ServeOptions = {}): Promise { + const port = options.port ?? DEFAULT_API_PORT; + const host = options.host ?? "127.0.0.1"; + const token = options.token && options.token.trim() ? options.token.trim() : null; + + // Fail soft, not open: never expose a public interface without a token. + if (!LOOPBACK_HOSTS.has(host) && !token) { + console.error( + `error: refusing to bind ${host} without a token. Pass --token ` + + `(or set TORHUNT_API_TOKEN), or bind 127.0.0.1.`, + ); + process.exit(1); + return; + } + + const runtime = await startRuntime(options.downloadDir); + + if (options.seedTimeMs && options.seedTimeMs > 0) { + startSeedReaper(runtime.queue, options.seedTimeMs, { deleteFiles: options.deleteFiles, log }); + } + + const server = http.createServer(createServeHandler(runtime, token)); await new Promise((resolve) => { server.listen(port, host, () => { diff --git a/src/daemon/websearch.test.ts b/src/daemon/websearch.test.ts new file mode 100644 index 0000000..30ed5a8 --- /dev/null +++ b/src/daemon/websearch.test.ts @@ -0,0 +1,111 @@ +import { describe, it, expect, vi } from "vitest"; +import { searchTorrents, SEARCH_RESULT_CAP } from "./websearch"; +import type { Source, SourceGroup, TorrentResult } from "../sources/types"; + +const HASH_A = "a".repeat(40); +const HASH_B = "b".repeat(40); + +function result(over: Partial = {}): TorrentResult { + const infoHash = over.infoHash ?? HASH_A; + return { + infoHash, + name: "Torrent", + sizeBytes: 1024, + seeders: 5, + leechers: 1, + source: "yts", + magnet: `magnet:?xt=urn:btih:${infoHash}`, + ...over, + }; +} + +function makeSource( + id: Source["id"], + label: string, + groups: SourceGroup[], + behavior: (() => Promise) | TorrentResult[], +): Source { + return { + id, + label, + groups, + homepage: "https://example.invalid", + reportsHealth: true, + search: () => (typeof behavior === "function" ? behavior() : Promise.resolve(behavior)), + }; +} + +describe("searchTorrents", () => { + it("merges results from all sources, healthiest first", async () => { + const sources = [ + makeSource("yts", "YTS", ["Movies"], [result({ infoHash: HASH_A, seeders: 3 })]), + makeSource("eztv", "EZTV", ["TV"], [result({ infoHash: HASH_B, seeders: 30 })]), + ]; + const out = await searchTorrents("query", null, { sources }); + expect(out.results.map((r) => r.infoHash)).toEqual([HASH_B, HASH_A]); + expect(out.failed).toEqual([]); + }); + + it("dedupes the same torrent across indexers, keeping the healthiest row", async () => { + const sources = [ + makeSource("yts", "YTS", ["Movies"], [result({ seeders: 2 })]), + makeSource("tpb-movies", "TPB", ["Movies"], [result({ seeders: 9 })]), + ]; + const out = await searchTorrents("query", null, { sources }); + expect(out.results).toHaveLength(1); + expect(out.results[0]!.seeders).toBe(9); + }); + + it("degrades failing sources to the failed list instead of throwing", async () => { + const sources = [ + makeSource("yts", "YTS", ["Movies"], [result({})]), + makeSource("nyaa", "Nyaa", ["Anime"], async () => { + throw new Error("down"); + }), + makeSource("eztv", "EZTV", ["TV"], async () => { + throw new Error("also down"); + }), + makeSource("subsplease", "SubsPlease", ["Anime"], async () => { + throw new Error("down too"); + }), + ]; + const out = await searchTorrents("query", null, { sources }); + expect(out.results).toHaveLength(1); + expect(out.failed).toEqual(["Nyaa", "EZTV", "SubsPlease"]); + }); + + it("only queries sources matching the requested category", async () => { + const animeSearch = vi.fn(async () => [] as TorrentResult[]); + const movieSearch = vi.fn(async () => [] as TorrentResult[]); + const anime = makeSource("nyaa", "Nyaa", ["Anime"], []); + anime.search = animeSearch; + const movies = makeSource("yts", "YTS", ["Movies"], []); + movies.search = movieSearch; + + await searchTorrents("query", "Anime", { sources: [anime, movies] }); + + expect(animeSearch).toHaveBeenCalledTimes(1); + expect(movieSearch).not.toHaveBeenCalled(); + }); + + it("caps the merged result list", async () => { + const many = Array.from({ length: SEARCH_RESULT_CAP + 20 }, (_, i) => + result({ infoHash: i.toString(16).padStart(40, "0"), seeders: SEARCH_RESULT_CAP + 20 - i }), + ); + const out = await searchTorrents("query", null, { + sources: [makeSource("yts", "YTS", ["Movies"], many)], + }); + expect(out.results).toHaveLength(SEARCH_RESULT_CAP); + }); + + it("propagates the abort deadline to each source", async () => { + const seen: (AbortSignal | undefined)[] = []; + const src = makeSource("yts", "YTS", ["Movies"], []); + src.search = (_q, opts) => { + seen.push(opts?.signal); + return Promise.resolve([]); + }; + await searchTorrents("q", null, { sources: [src], timeoutMs: 1000 }); + expect(seen[0]).toBeInstanceOf(AbortSignal); + }); +}); diff --git a/src/daemon/websearch.ts b/src/daemon/websearch.ts new file mode 100644 index 0000000..49688db --- /dev/null +++ b/src/daemon/websearch.ts @@ -0,0 +1,74 @@ +// Headless search bridge: run the same source adapters the TUI's +// useConcurrentSearch hook runs, but resolved as a single HTTP response instead +// of a streaming render. Sources race in parallel under one deadline; a source +// failing or timing out degrades to an entry in `failed` rather than failing +// the whole request, mirroring how the TUI shows per-source status and keeps +// the rest of the results. +import { SOURCES } from "../sources/registry"; +import type { Source, SourceGroup, TorrentResult } from "../sources/types"; + +export const SEARCH_RESULT_CAP = 60; + +// The TUI gives multi-step scrapers 15s before timing them out gracefully; a +// search that outlives that is dead weight for an HTTP caller too. +const SOURCE_TIMEOUT_MS = 15_000; + +export interface SearchDeps { + sources?: readonly Source[]; + timeoutMs?: number; +} + +export interface SearchOutcome { + results: TorrentResult[]; + failed: string[]; +} + +// Same ordering the TUI defaults to: healthiest first, newest as tiebreak. +function defaultOrder(list: TorrentResult[]): TorrentResult[] { + return list.sort((a, b) => { + if (b.seeders !== a.seeders) return b.seeders - a.seeders; + return (b.added ?? 0) - (a.added ?? 0); + }); +} + +// Same infoHash dedupe the results view applies: several indexers often return +// the same torrent, and the healthiest row wins. +function dedupe(list: TorrentResult[]): TorrentResult[] { + const byHash = new Map(); + for (const r of list) { + const existing = byHash.get(r.infoHash); + if (!existing || r.seeders > existing.seeders) byHash.set(r.infoHash, r); + } + return [...byHash.values()]; +} + +export async function searchTorrents( + query: string, + category: SourceGroup | null, + deps: SearchDeps = {}, +): Promise { + const pool = deps.sources ?? SOURCES; + const picked = category + ? pool.filter((s) => s.groups?.includes(category)) + : pool; + const signal = AbortSignal.timeout(deps.timeoutMs ?? SOURCE_TIMEOUT_MS); + const settled = await Promise.all( + picked.map(async (source) => { + try { + return { + ok: true as const, + label: source.label, + results: await source.search(query, { signal }), + }; + } catch { + return { ok: false as const, label: source.label, results: [] as TorrentResult[] }; + } + }), + ); + const failed = [...new Set(settled.filter((s) => !s.ok).map((s) => s.label))]; + const results = defaultOrder(dedupe(settled.flatMap((s) => s.results))).slice( + 0, + SEARCH_RESULT_CAP, + ); + return { results, failed }; +} diff --git a/src/daemon/webui.test.ts b/src/daemon/webui.test.ts new file mode 100644 index 0000000..283bd63 --- /dev/null +++ b/src/daemon/webui.test.ts @@ -0,0 +1,105 @@ +import { describe, it, expect } from "vitest"; +import { EventEmitter } from "node:events"; +import { + loadUiHtml, + openEventStream, + originAllowed, + type EventStreamWriter, +} from "./webui"; + +describe("loadUiHtml", () => { + it("finds the web remote shell in the source tree", () => { + const html = loadUiHtml(); + expect(html).not.toBeNull(); + expect(html).toContain(" { + expect(loadUiHtml()).toBe(loadUiHtml()); + }); +}); + +describe("originAllowed", () => { + it("allows requests without an Origin (curl, scripts, health checks)", () => { + expect(originAllowed(undefined, "localhost:9161")).toBe(true); + }); + + it("allows same-origin browser posts regardless of case or trailing path", () => { + expect(originAllowed("http://localhost:9161", "localhost:9161")).toBe(true); + expect(originAllowed("http://LOCALHOST:9161/add", "localhost:9161")).toBe(true); + }); + + it("rejects cross-site origins", () => { + expect(originAllowed("http://evil.example", "localhost:9161")).toBe(false); + expect(originAllowed("http://localhost:9999", "localhost:9161")).toBe(false); + }); + + it("rejects malformed origins and requests with no host", () => { + expect(originAllowed("not a url", "localhost:9161")).toBe(false); + expect(originAllowed("http://localhost:9161", undefined)).toBe(false); + }); +}); + +describe("openEventStream", () => { + function fakeRes(): { + head: { status: number; headers: Record } | null; + writes: string[]; + on(event: string, listener: () => void): unknown; + writeHead(status: number, headers: Record): void; + write(chunk: string): void; + emitClose(): void; + } { + const listeners = new Map void>(); + return { + head: null, + writes: [], + on(event, listener) { + listeners.set(event, listener); + return this; + }, + writeHead(status, headers) { + this.head = { status, headers }; + }, + write(chunk) { + this.writes.push(chunk); + }, + emitClose() { + listeners.get("close")?.(); + }, + }; + } + + it("writes sse headers, an initial snapshot, and pushes on queue updates", () => { + const res = fakeRes(); + const queue = new EventEmitter(); + openEventStream(res as unknown as EventStreamWriter, queue, () => ({ n: 1 })); + expect(res.head?.status).toBe(200); + expect(res.head?.headers["Content-Type"]).toContain("text/event-stream"); + expect(res.head?.headers["Cache-Control"]).toBe("no-store"); + expect(res.writes[0]).toContain("retry:"); + expect(res.writes[1]).toBe('data: {"n":1}\n\n'); + queue.emit("update"); + expect(res.writes[2]).toBe('data: {"n":1}\n\n'); + }); + + it("detaches from the queue when the client disconnects", () => { + const res = fakeRes(); + const queue = new EventEmitter(); + openEventStream(res as unknown as EventStreamWriter, queue, () => ({})); + res.emitClose(); + const before = res.writes.length; + queue.emit("update"); + expect(res.writes.length).toBe(before); + }); + + it("survives a snapshot function that throws", () => { + const res = fakeRes(); + const queue = new EventEmitter(); + openEventStream(res as unknown as EventStreamWriter, queue, () => { + throw new Error("boom"); + }); + queue.emit("update"); + expect(res.writes.every((w) => !w.startsWith("data: {"))).toBe(true); + }); +}); diff --git a/src/daemon/webui.ts b/src/daemon/webui.ts new file mode 100644 index 0000000..9bb6d0a --- /dev/null +++ b/src/daemon/webui.ts @@ -0,0 +1,107 @@ +// Web remote: the static front-end served by `torhunt serve`, plus the two +// pieces of server plumbing it needs that the plain JSON API doesn't provide — +// an asset locator for the bundled HTML shell and a server-sent-event stream +// that pushes queue snapshots to the browser. Everything here is structural and +// dependency-free so it stays trivially testable without node:http. +// +// Security posture (mirrors serve.ts): the HTML shell carries no secrets, so it +// is served like /health — after the tokenless Host check but without bearer +// auth. Every data route (/events included) still requires authorization, and +// cross-site POSTs are rejected by origin comparison before any handler runs. + +import { readFileSync } from "node:fs"; + +// The shell ships beside the bundle in dist/ (postbuild copies it there) and +// lives under src/daemon/assets/ in a checkout. Try both layouts; cache the +// first hit so a hot loop never re-stats the disk. +let cachedHtml: string | null | undefined; + +export function loadUiHtml(): string | null { + if (cachedHtml !== undefined) return cachedHtml; + const candidates = [ + new URL("./ui.html", import.meta.url), // bundled: dist/ui.html + new URL("./assets/ui.html", import.meta.url), // source tree + ]; + for (const url of candidates) { + try { + cachedHtml = readFileSync(url, "utf8"); + return cachedHtml; + } catch { + // try the next layout + } + } + cachedHtml = null; + return cachedHtml; +} + +// Minimal structural types so tests can pass plain objects instead of real +// http.ServerResponse / EventEmitter instances. +export interface EventStreamWriter { + writeHead(status: number, headers: Record): unknown; + write(chunk: string): unknown; + on(event: string, listener: () => void): unknown; +} + +export interface QueueEventSource { + on(event: string, listener: () => void): unknown; + off(event: string, listener: () => void): unknown; +} + +export const SSE_HEARTBEAT_MS = 15_000; + +// Open an SSE stream on `res` and push a fresh snapshot on every queue update, +// plus a keepalive comment so intermediaries don't reap an idle connection. +// The stream cleans up after itself when the client disconnects. +export function openEventStream( + res: EventStreamWriter, + queue: QueueEventSource, + snapshot: () => unknown, +): void { + res.writeHead(200, { + "Content-Type": "text/event-stream; charset=utf-8", + "Cache-Control": "no-store", + Connection: "keep-alive", + }); + // Ask browsers to retry quickly after a daemon restart. + res.write("retry: 3000\n\n"); + + const push = (): void => { + try { + res.write(`data: ${JSON.stringify(snapshot())}\n\n`); + } catch { + // A torn write means the socket died; the close listener detaches us. + } + }; + push(); + + queue.on("update", push); + const heartbeat = setInterval(() => { + try { + res.write(": keepalive\n\n"); + } catch { + // same as above — close will clean up + } + }, SSE_HEARTBEAT_MS); + heartbeat.unref?.(); + + const detach = (): void => { + queue.off("update", push); + clearInterval(heartbeat); + }; + res.on("close", detach); +} + +// Cross-site request guard for state-changing calls. Browsers attach an Origin +// header to every POST they send — same-origin ones match the request's Host, +// forged ones don't. curl and scripts send no Origin at all and pass through. +export function originAllowed(origin: string | undefined, hostHeader: string | undefined): boolean { + if (!origin) return true; + let originHost: string; + try { + originHost = new URL(origin).host.toLowerCase(); + } catch { + return false; // malformed Origin is never trusted + } + if (!hostHeader) return false; + return originHost === hostHeader.trim().toLowerCase(); +}