Keep one long-lived yt-dlp worker process instead of spawning per call
This commit is contained in:
112
server/ytdlp-pool.js
Normal file
112
server/ytdlp-pool.js
Normal file
@@ -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 <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; } };
|
||||
}
|
||||
Reference in New Issue
Block a user