113 lines
4.5 KiB
JavaScript
113 lines
4.5 KiB
JavaScript
/* 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 <args>` 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; } };
|
|
}
|