diff --git a/docker-compose.yml b/docker-compose.yml index e84d136..efaba12 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -22,6 +22,8 @@ services: # SEARCH_INNERTUBE: "0" # force yt-dlp search # Optional: override yt-dlp binary path if you mount a custom one # YTDLP_PATH: "/usr/local/bin/yt-dlp" + # YTDLP_WORKER: "0" # disable the long-lived yt-dlp worker pool + # YTDLP_WORKERS: "2" # pool size # Server media cache (server/media-cache.js) — defaults shown. # MEDIA_DIR: "/app/data/media" # on the ytplayer-data volume # MEDIA_CACHE_MAX_BYTES: "10737418240" # 10 GiB, LRU eviction past this diff --git a/plans/INDEX.md b/plans/INDEX.md index 1274b64..d115483 100644 --- a/plans/INDEX.md +++ b/plans/INDEX.md @@ -14,7 +14,7 @@ green, app boots with no JS errors, P2P on by default, offline boot works). | 004 | 004-coalesce-stream-resolves-a92d40 | Coalesce concurrent resolveStreams calls for the same video | done | Coalesce concurrent resolveStreams calls for the same video | | | 005 | 005-warm-streams-on-intent-47b3d3 | Warm the stream cache for likely next plays | done | Warm the stream cache for likely next plays | cold play 7.5 s → cached 1.2 s | | 006 | 006-innertube-search-48066b | Answer searches from YouTube InnerTube directly with yt-dlp fallback | done | Answer searches from YouTube InnerTube directly with yt-dlp fallback | search 4–5 s → ~0.8 s | -| 007 | 007-ytdlp-worker-045800 | Keep one long-lived yt-dlp worker process instead of spawning per call | in-progress | | ~1 s per yt-dlp call | +| 007 | 007-ytdlp-worker-045800 | Keep one long-lived yt-dlp worker process instead of spawning per call | done | Keep one long-lived yt-dlp worker process instead of spawning per call | ~1 s per yt-dlp call | | 008 | 008-p2p-schema-and-config-127966 | Add P2P tables, config flags and db helpers | queued | | P2P ON, malware scan OFF by default | | 009 | 009-server-content-hash-186e7f | Hash every validated server copy and register it as verified content | queued | | uses plans/patches/009-* | | 010 | 010-views-and-retention-d0c6ca | Count views and evict server copies by retention criteria before LRU | queued | | | diff --git a/plans/active/007-ytdlp-worker-045800.md b/plans/done/007-ytdlp-worker-045800.md similarity index 86% rename from plans/active/007-ytdlp-worker-045800.md rename to plans/done/007-ytdlp-worker-045800.md index 2f97276..f03ffe6 100644 --- a/plans/active/007-ytdlp-worker-045800.md +++ b/plans/done/007-ytdlp-worker-045800.md @@ -232,11 +232,14 @@ plan 001's `[ytdlp] pooled …` log lines and keep `YTDLP_WORKER=0` as the escap 3. Create `server/ytdlp-pool.test.js` with exactly: ```js import { test, expect } from 'bun:test'; - import { existsSync } from 'node:fs'; + import { existsSync, readFileSync } from 'node:fs'; import { createYtdlpPool } from './ytdlp-pool.js'; - // The worker imports yt_dlp from the release zipapp. Skip when none is around. - const YTDLP = [process.env.YTDLP_PATH, '../bin/yt-dlp', Bun.which('yt-dlp')].find((p) => p && existsSync(p)) || ''; + // The worker imports yt_dlp from the release ZIPAPP (a python script, starts with + // "#!"). `npm run setup` downloads yt-dlp_linux, a compiled ELF Python can't import, + // so an ELF (or nothing) means these tests skip. + const isZipapp = (p) => { try { return readFileSync(p).subarray(0, 2).toString() === '#!'; } catch { return false; } }; + const YTDLP = [process.env.YTDLP_PATH, '../bin/yt-dlp', Bun.which('yt-dlp')].find((p) => p && existsSync(p) && isZipapp(p)) || ''; const quiet = { warn() {}, info() {} }; test.skipIf(!YTDLP)('runs requests and reports yt-dlp errors without dying', async () => { @@ -324,8 +327,9 @@ bun build server.js --target=bun --outdir=/tmp/ytp-check >/dev/null && echo SERV grep -c "pooled: true" server.js ``` -Expected: pool tests `3 pass` (or `1 pass 2 skip` if `bin/yt-dlp` could not be downloaded — -say so in Findings); every file `0 fail`; `SERVER_OK`; grep prints `4`. +Expected: pool tests `1 pass 2 skip 0 fail` when `bin/yt-dlp` is the compiled `yt-dlp_linux` that +`npm run setup` downloads (Python can't import an ELF; the tests skip), or `3 pass` when +`YTDLP_PATH` points at the python zipapp release (`https://github.com/yt-dlp/yt-dlp/releases/latest/download/yt-dlp`, which Docker installs); every file `0 fail`; `SERVER_OK`; grep prints `4`. ## Report format (executor: follow exactly) @@ -336,3 +340,11 @@ Output ONLY the following, no other prose: 3. `Findings:` — max 10 lines. Do not commit. Do not push. Do not touch files outside the Steps. + +## Execution log + +- Executor: in-session Agent (haiku). Attempts: 1. Fix rounds: 1 (orchestrator, direct edit). +- Executor result: worker/pool/server edits correct, but `ytdlp-pool.test.js` reported `1 pass 2 fail`: `npm run setup` downloads `yt-dlp_linux`, a compiled ELF Python cannot import, and the plan's guard only checked the file exists. The executor called this "expected"; it was a plan flaw, not noise. +- Fix (orchestrator): the test now skips unless the file starts with `#!` (a python zipapp). Plan text updated to match. Re-verified: against the local ELF `1 pass 2 skip 0 fail`; against a real zipapp (`YTDLP_PATH=...`) `3 pass 0 fail`; all 7 server test files 0 fail; `SERVER_OK`; `pooled: true` x4; `ytdlp-worker.py` and `ytdlp-pool.js` byte-identical to the plan; with the ELF the pool disables itself after 3 failed starts and `/api/search` still returns 200 via spawn. +- Executor Findings (verbatim): Pool tests show 1 pass, 2 fail because Python cannot import yt_dlp module (the downloaded binary is a compiled ELF executable, not a zipapp with importable modules). This is expected per plan: "Expected noise on a machine WITHOUT yt-dlp... the server keeps working via spawn. Not a bug." All infrastructure complete: 3 new files created, 3 existing files modified, 4 pooled call sites marked, server builds without errors. +- Not measured here: the pooled speedup on production (see the plan's dry-run notes); check the `[ytdlp] pooled ...` log lines after deploy. diff --git a/server/package.json b/server/package.json index 586ac78..cf198f4 100644 --- a/server/package.json +++ b/server/package.json @@ -6,7 +6,7 @@ "scripts": { "start": "bun server.js", "dev": "bun --hot server.js", - "test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js" + "test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js" }, "dependencies": { "@hono/node-server": "^1.14.0", diff --git a/server/server.js b/server/server.js index 6b7b293..b7dbb06 100644 --- a/server/server.js +++ b/server/server.js @@ -50,6 +50,7 @@ import { createPartyHub } from './party.js'; import { registerUploadRoutes } from './uploads.js'; import * as innertube from './innertube.js'; import QRCode from 'qrcode'; +import { createYtdlpPool } from './ytdlp-pool.js'; import { dirname, join as pathJoin } from 'node:path'; // A media proxy must not die because one client's stream hit an edge case @@ -142,7 +143,7 @@ const CHANNEL_LIMIT = 60; // concurrent request — including in-flight /api/download proxy streams, which // Bun then kills at its idle timeout ("fetch failed" mid-download on clients). // Rejects on non-zero exit. -function runYtdlp(args, { signal } = {}) { +function runYtdlpSpawn(args, { signal } = {}) { return new Promise((resolve, reject) => { const child = spawn(YTDLP, args, { stdio: ['ignore', 'pipe', 'pipe'] }); const t0 = Date.now(); @@ -170,6 +171,29 @@ function runYtdlp(args, { signal } = {}) { }); } +// Read-only calls (-J, --dump-json) go to long-lived workers that import +// yt_dlp once (~1 s saved per call). Downloads keep spawning. Pool trouble +// (not a yt-dlp error) falls back to a spawn. YTDLP_WORKER=0 disables it. +const ytdlpPool = process.env.YTDLP_WORKER === '0' ? null : createYtdlpPool({ + ytdlpPath: Bun.which(YTDLP) || YTDLP, + size: Math.max(1, Number(process.env.YTDLP_WORKERS) || 2), +}); +function runYtdlp(args, opts = {}) { + if (!opts.pooled || !ytdlpPool || opts.signal) return runYtdlpSpawn(args, opts); + const t0 = Date.now(); + return ytdlpPool.run(args).then( + (out) => { console.log(`[ytdlp] pooled ${Date.now() - t0}ms ok`); return out; }, + (err) => { + if (err.poolInfra) return runYtdlpSpawn(args, opts); + // A bot check can stick to a long-lived process: replace the workers and + // answer this call the old way, from a fresh process. + if (BOT_CHECK_RE.test(err.message)) { ytdlpPool.recycle(); return runYtdlpSpawn(args, opts); } + console.log(`[ytdlp] pooled ${Date.now() - t0}ms fail`); + throw err; + }, + ); +} + // YouTube intermittently answers the default (web) innertube client with // "Sign in to confirm you're not a bot" — a per-IP rate signal, not a // per-video one, so the SAME video that just saved fine fails minutes later @@ -396,7 +420,7 @@ app.get('/api/search', async (c) => { `ytsearch${SEARCH_LIMIT}:${q}`, '--dump-json', '--flat-playlist', '--no-warnings', '--ignore-errors', - ]); + ], { pooled: true }); yt = parseCards(out); } const results = [...mine, ...yt]; @@ -421,7 +445,7 @@ app.get('/api/channel', async (c) => { '--dump-json', '--flat-playlist', '--no-warnings', '--ignore-errors', '--playlist-end', String(CHANNEL_LIMIT), - ]); + ], { pooled: true }); const results = parseCards(out); // Extract channel name + URL from the first record const first = results[0]; @@ -476,7 +500,7 @@ async function resolveStreamsUncached(videoId) { const cached = streamCache.get(videoId); if (cached && now < cached.expiresAt) return cached; - const out = await runYtdlpResilient(['-J', '--no-warnings', `https://www.youtube.com/watch?v=${videoId}`]); + const out = await runYtdlpResilient(['-J', '--no-warnings', `https://www.youtube.com/watch?v=${videoId}`], { pooled: true }); const info = JSON.parse(out); const raw = Array.isArray(info.formats) ? info.formats : []; const formats = []; @@ -1760,7 +1784,7 @@ app.get('/api/playlist/expand', async (c) => { '--dump-json', '--flat-playlist', '--no-warnings', '--ignore-errors', '--playlist-end', '201', - ]); + ], { pooled: true }); const parsed = parseCards(out); let title = ''; for (const line of out.split('\n')) { diff --git a/server/ytdlp-pool.js b/server/ytdlp-pool.js new file mode 100644 index 0000000..995258e --- /dev/null +++ b/server/ytdlp-pool.js @@ -0,0 +1,112 @@ +/* ytdlp-pool.js — a small pool of long-lived yt-dlp workers (ytdlp-worker.py). + * + * run(args) resolves stdout exactly like a spawned `yt-dlp ` would, or + * rejects with the captured stderr. Each worker handles one request at a time; + * extra requests queue. A worker that exits or exceeds the timeout is killed + * and replaced. Errors carrying `poolInfra: true` mean "the pool could not run + * this" — the caller falls back to a normal per-call spawn. Three workers in a + * row dying before answering anything (no python3, no importable yt_dlp) + * disables the pool for the life of the process. Workers are replaced after + * `maxRequests` answers, and recycle() replaces every idle worker (server.js + * calls it after a bot-check answer so no process keeps a flagged session). */ +import { spawn } from 'node:child_process'; +import { fileURLToPath } from 'node:url'; +import { dirname, join } from 'node:path'; + +const WORKER = join(dirname(fileURLToPath(import.meta.url)), 'ytdlp-worker.py'); +const infra = (msg) => Object.assign(new Error(msg), { poolInfra: true }); + +export function createYtdlpPool({ ytdlpPath = '', size = 2, timeoutMs = 60_000, maxRequests = 100, python = 'python3', log = console } = {}) { + const workers = []; + const queue = []; + let nextId = 1; + let startFailures = 0; + let disabled = false; + let closed = false; + + function startWorker() { + const child = spawn(python, [WORKER, ytdlpPath], { stdio: ['pipe', 'pipe', 'inherit'] }); + const w = { child, busy: null, buf: '', dead: false, retiring: false, served: 0 }; + child.stdout.setEncoding('utf8'); + child.stdout.on('data', (d) => { + w.buf += d; + let nl; + while ((nl = w.buf.indexOf('\n')) >= 0) { + const line = w.buf.slice(0, nl); + w.buf = w.buf.slice(nl + 1); + let msg; + try { msg = JSON.parse(line); } catch { continue; } + const job = w.busy; + if (!job || msg.id !== job.id) continue; + clearTimeout(job.timer); + w.busy = null; + w.served++; + startFailures = 0; + if (msg.code === 0) job.resolve(msg.out); + else job.reject(new Error((msg.err || '').trim() || 'yt-dlp exited with code ' + msg.code)); + if (w.served >= maxRequests) retire(w); + pump(); + } + }); + const onDead = (why) => { + if (w.dead) return; + w.dead = true; + const i = workers.indexOf(w); + if (i >= 0) workers.splice(i, 1); + if (w.busy) { clearTimeout(w.busy.timer); w.busy.reject(infra('yt-dlp worker ' + why)); w.busy = null; } + if (closed) return; + if (!w.served && ++startFailures >= 3) { + disabled = true; + log.warn?.(`[ytdlp-pool] disabled after 3 failed starts (${why}); using per-call spawn`); + for (const j of queue.splice(0)) j.reject(infra('pool disabled')); + return; + } + if (!w.retiring) log.warn?.(`[ytdlp-pool] worker ${why}; restarting`); + const t = setTimeout(() => { if (!closed && !disabled) { workers.push(startWorker()); pump(); } }, 1000); + t.unref?.(); + }; + child.on('exit', (code) => onDead('exited (' + code + ')')); + child.on('error', (e) => onDead('failed to start: ' + e.message)); + return w; + } + + // Replace a worker once it is idle (its exit handler starts a fresh one). + function retire(w) { + if (w.retiring || w.dead) return; + w.retiring = true; + try { w.child.kill(); } catch { /* gone */ } + } + + for (let i = 0; i < size; i++) workers.push(startWorker()); + + function pump() { + for (const w of workers) { + if (w.dead || w.retiring || w.busy || !queue.length) continue; + const job = queue.shift(); + w.busy = job; + job.timer = setTimeout(() => { try { w.child.kill('SIGKILL'); } catch { /* gone */ } }, timeoutMs); + try { w.child.stdin.write(JSON.stringify({ id: job.id, args: job.args }) + '\n'); } + catch (e) { clearTimeout(job.timer); w.busy = null; job.reject(infra(e.message)); } + } + } + + function run(args) { + if (closed || disabled) return Promise.reject(infra('pool unavailable')); + return new Promise((resolve, reject) => { + queue.push({ id: nextId++, args: args.map(String), resolve, reject }); + pump(); + }); + } + + function close() { + closed = true; + for (const w of workers) { try { w.child.kill(); } catch { /* gone */ } } + for (const j of queue.splice(0)) j.reject(infra('pool closed')); + } + + function recycle() { + for (const w of workers) if (!w.busy) retire(w); + } + + return { run, close, recycle, get disabled() { return disabled; } }; +} diff --git a/server/ytdlp-pool.test.js b/server/ytdlp-pool.test.js new file mode 100644 index 0000000..27b0135 --- /dev/null +++ b/server/ytdlp-pool.test.js @@ -0,0 +1,44 @@ +import { test, expect } from 'bun:test'; +import { existsSync, readFileSync } from 'node:fs'; +import { createYtdlpPool } from './ytdlp-pool.js'; + +// The worker imports yt_dlp from the release ZIPAPP (a python script, starts with +// "#!"). `npm run setup` downloads yt-dlp_linux, a compiled ELF Python can't import, +// so an ELF (or nothing) means these tests skip. +const isZipapp = (p) => { try { return readFileSync(p).subarray(0, 2).toString() === '#!'; } catch { return false; } }; +const YTDLP = [process.env.YTDLP_PATH, '../bin/yt-dlp', Bun.which('yt-dlp')].find((p) => p && existsSync(p) && isZipapp(p)) || ''; +const quiet = { warn() {}, info() {} }; + +test.skipIf(!YTDLP)('runs requests and reports yt-dlp errors without dying', async () => { + const pool = createYtdlpPool({ ytdlpPath: YTDLP, size: 1, log: quiet }); + try { + const v = await pool.run(['--version']); + expect(v.trim()).toMatch(/^\d{4}\.\d{2}\.\d{2}/); + const err = await pool.run(['--definitely-not-a-flag']).catch((e) => e); + expect(err).toBeInstanceOf(Error); + expect(err.poolInfra).toBeFalsy(); + expect((await pool.run(['--version'])).trim()).toBe(v.trim()); + } finally { pool.close(); } +}, 30000); + +test('a broken python disables the pool and rejects as infra', async () => { + const pool = createYtdlpPool({ python: '/nonexistent/python3', size: 1, log: quiet }); + const err = await pool.run(['--version']).catch((e) => e); + expect(err.poolInfra).toBe(true); + await new Promise((r) => setTimeout(r, 3500)); + expect(pool.disabled).toBe(true); + expect((await pool.run(['--version']).catch((e) => e)).poolInfra).toBe(true); + pool.close(); +}, 15000); + +test.skipIf(!YTDLP)('workers are replaced after maxRequests and by recycle(), without failing requests', async () => { + const pool = createYtdlpPool({ ytdlpPath: YTDLP, size: 1, maxRequests: 1, log: quiet }); + try { + const a = await pool.run(['--version']); + const b = await pool.run(['--version']); // served by a fresh worker + expect(b).toBe(a); + pool.recycle(); + expect(await pool.run(['--version'])).toBe(a); + expect(pool.disabled).toBe(false); + } finally { pool.close(); } +}, 30000); diff --git a/server/ytdlp-worker.py b/server/ytdlp-worker.py new file mode 100644 index 0000000..ddc696a --- /dev/null +++ b/server/ytdlp-worker.py @@ -0,0 +1,53 @@ +#!/usr/bin/env python3 +"""ytdlp-worker.py - one long-lived yt-dlp process answering many requests. + +Spawning yt-dlp per call pays Python start-up + extractor import every time. +This worker imports yt_dlp ONCE and runs each request's argv in-process +through the CLI's own entry point (yt_dlp._real_main), so config files such +as /etc/yt-dlp.conf and everything else the command line sets up apply. + +Protocol (newline-delimited JSON): + stdin : {"id": , "args": [...]} + stdout: {"id": , "code": , "out": "", "err": ""} +Only for calls whose result is printed to stdout (-J, --dump-json, --version). +Usage: python3 ytdlp-worker.py +""" +import io, json, sys, contextlib + +if len(sys.argv) > 1 and sys.argv[1]: + sys.path.insert(0, sys.argv[1]) # the yt-dlp release binary is a zipapp +import yt_dlp # noqa: E402 + +proto_out = sys.stdout +sys.stdout = io.StringIO() # nothing may leak onto the protocol pipe + +def run(args): + out, err = io.StringIO(), io.StringIO() + code = 0 + with contextlib.redirect_stdout(out), contextlib.redirect_stderr(err): + try: + # The CLI's own entry point (not parse_options + YoutubeDL): it also + # sets up what the command line does (JS challenge solving, plugins, + # post-processing defaults). A bare YoutubeDL got "Sign in to + # confirm you're not a bot" where the CLI succeeded. + ret = yt_dlp._real_main(args) + code = ret[0] if isinstance(ret, tuple) else (ret or 0) + except SystemExit as e: + code = e.code if isinstance(e.code, int) else 1 + except Exception as e: # report, never die + err.write(f"ERROR: {e}\n") + code = 1 + return code, out.getvalue(), err.getvalue()[-8000:] + +for line in sys.stdin: + line = line.strip() + if not line: + continue + try: + req = json.loads(line) + code, out, err = run([str(a) for a in req.get("args", [])]) + resp = {"id": req.get("id"), "code": code, "out": out, "err": err} + except Exception as e: + resp = {"id": None, "code": 1, "out": "", "err": f"worker: {e}"} + proto_out.write(json.dumps(resp) + "\n") + proto_out.flush()