Coalesce save downloads, kill them on disconnect, refuse live and over-long videos

This commit is contained in:
Jonathan Sykes
2026-08-23 23:50:10 +08:00
parent ea9e372c1e
commit 04c419d151

View File

@@ -113,9 +113,16 @@ const CHANNEL_LIMIT = 60;
// concurrent request — including in-flight /api/download proxy streams, which // concurrent request — including in-flight /api/download proxy streams, which
// Bun then kills at its idle timeout ("fetch failed" mid-download on clients). // Bun then kills at its idle timeout ("fetch failed" mid-download on clients).
// Rejects on non-zero exit. // Rejects on non-zero exit.
function runYtdlp(args) { function runYtdlp(args, { signal } = {}) {
return new Promise((resolve, reject) => { return new Promise((resolve, reject) => {
const child = spawn(YTDLP, args, { stdio: ['ignore', 'pipe', 'pipe'] }); const child = spawn(YTDLP, args, { stdio: ['ignore', 'pipe', 'pipe'] });
// Kill the download when the requesting client goes away — otherwise an
// aborted/retried save leaves yt-dlp running to completion (8 copies of
// one video were found pulling in parallel after the client retried).
if (signal) {
if (signal.aborted) child.kill('SIGTERM');
else signal.addEventListener('abort', () => child.kill('SIGTERM'), { once: true });
}
let out = ''; let out = '';
let err = ''; let err = '';
child.stdout.setEncoding('utf8'); child.stdout.setEncoding('utf8');
@@ -673,35 +680,73 @@ app.get('/api/play', async (c) => {
// with --get-url and re-fetching it server-side with hand-rolled browser // with --get-url and re-fetching it server-side with hand-rolled browser
// headers intermittently got 403s from the YouTube CDN when the User-Agent // headers intermittently got 403s from the YouTube CDN when the User-Agent
// didn't match the extraction client. // didn't match the extraction client.
async function ytdlpDownloadResponse(videoId, fp, formatArgs) { // Live streams have no end: yt-dlp/ffmpeg would pull the HLS manifest forever
// (one such "save" ran for an hour and ate 14 GB). Refuse them up front.
const MAX_SAVE_SECONDS = Number(process.env.MAX_SAVE_SECONDS || 4 * 3600);
async function assertNotLive(videoId) {
const { info } = await resolveStreams(videoId);
if (info.is_live || info.live_status === 'is_live' || info.live_status === 'post_live') {
throw new Error('live streams cannot be saved');
}
if (typeof info.duration === 'number' && info.duration > MAX_SAVE_SECONDS) {
throw new Error(`video is too long to save (${Math.round(info.duration / 3600)} h, limit ${MAX_SAVE_SECONDS / 3600} h)`);
}
}
// Same video + same format args requested while a download is already in
// flight (the OPFS worker retries, then the main-thread fallback retries
// again) share ONE yt-dlp run instead of spawning another each time.
const inflightDownloads = new Map(); // key -> { promise, waiters, tmpBase }
async function ytdlpDownloadResponse(videoId, fp, formatArgs, signal) {
await assertNotLive(videoId);
const key = videoId + '|' + formatArgs.join(' ');
let entry = inflightDownloads.get(key);
if (!entry) {
const tmpBase = `ytp-dl-${videoId}-${Date.now()}`; const tmpBase = `ytp-dl-${videoId}-${Date.now()}`;
const tmp = `${tmpdir()}/${tmpBase}.mp4`; const tmp = `${tmpdir()}/${tmpBase}.mp4`;
let size, fd; const ctl = new AbortController();
try { entry = { waiters: 0, tmpBase, ctl };
await runYtdlp([ entry.promise = runYtdlp([
`https://www.youtube.com/watch?v=${videoId}`, `https://www.youtube.com/watch?v=${videoId}`,
'--no-warnings', '--no-playlist', '--no-warnings', '--no-playlist',
...formatArgs, ...formatArgs,
'--limit-rate', DOWNLOAD_RATE, '--limit-rate', DOWNLOAD_RATE,
'-o', tmp, '-o', tmp,
]); ], { signal: ctl.signal }).then(() => ({ tmp, size: statSync(tmp).size }));
size = statSync(tmp).size; inflightDownloads.set(key, entry);
// Open the fd BEFORE the finally unlinks: on Linux the data stays }
entry.waiters++;
// Only abort the shared yt-dlp when EVERY waiter has gone away.
let gone = false;
const leave = () => { if (gone) return; gone = true; if (--entry.waiters <= 0) { entry.ctl.abort(); } };
if (signal) signal.addEventListener('abort', leave, { once: true });
let size, fd;
try {
const res = await entry.promise;
size = res.size;
// Open the fd BEFORE the sweep unlinks: on Linux the data stays
// readable until the fd closes, so the temp file cleans itself up even // readable until the fd closes, so the temp file cleans itself up even
// if the client disconnects mid-transfer. // if the client disconnects mid-transfer.
fd = openSync(tmp, 'r'); fd = openSync(res.tmp, 'r');
} finally { } finally {
// Sweep everything yt-dlp may have left under this request's unique if (signal) signal.removeEventListener('abort', leave);
if (!gone) { gone = true; entry.waiters--; }
if (entry.waiters <= 0 && inflightDownloads.get(key) === entry) {
inflightDownloads.delete(key);
// Sweep everything yt-dlp may have left under this run's unique
// prefix: the output itself, .part partials, and .fNNN single-format // prefix: the output itself, .part partials, and .fNNN single-format
// intermediates (left when ffmpeg is missing — yt-dlp then downloads // intermediates (left when ffmpeg is missing — yt-dlp then downloads
// the streams separately, exits 0 without merging, and statSync above // the streams separately, exits 0 without merging, and statSync above
// throws on the absent merged file). // throws on the absent merged file).
for (const name of readdirSync(tmpdir())) { for (const name of readdirSync(tmpdir())) {
if (name.startsWith(tmpBase)) { if (name.startsWith(entry.tmpBase)) {
try { unlinkSync(`${tmpdir()}/${name}`); } catch { /* already gone */ } try { unlinkSync(`${tmpdir()}/${name}`); } catch { /* already gone */ }
} }
} }
} }
}
const stream = createReadStream('', { fd }); const stream = createReadStream('', { fd });
if (fp) recordVideoAccess(fp, { id: videoId }).catch(() => {}); if (fp) recordVideoAccess(fp, { id: videoId }).catch(() => {});
@@ -723,7 +768,8 @@ async function ytdlpDownloadResponse(videoId, fp, formatArgs) {
// concatenate them into one continuous mp4 — the user's custom cut. The result // concatenate them into one continuous mp4 — the user's custom cut. The result
// is streamed to the browser exactly like a normal save, so OPFS stores it // is streamed to the browser exactly like a normal save, so OPFS stores it
// under the caller-chosen custom id. Every temp file is swept afterwards. // under the caller-chosen custom id. Every temp file is swept afterwards.
async function ytdlpEditedDownloadResponse(videoId, fp, keep) { async function ytdlpEditedDownloadResponse(videoId, fp, keep, signal) {
await assertNotLive(videoId);
const tmpBase = `ytp-edit-${videoId}-${Date.now()}`; const tmpBase = `ytp-edit-${videoId}-${Date.now()}`;
const srcTmp = `${tmpdir()}/${tmpBase}.src.mp4`; const srcTmp = `${tmpdir()}/${tmpBase}.src.mp4`;
const outTmp = `${tmpdir()}/${tmpBase}.out.mp4`; const outTmp = `${tmpdir()}/${tmpBase}.out.mp4`;
@@ -737,7 +783,7 @@ async function ytdlpEditedDownloadResponse(videoId, fp, keep) {
'--merge-output-format', 'mp4', '--merge-output-format', 'mp4',
'--limit-rate', DOWNLOAD_RATE, '--limit-rate', DOWNLOAD_RATE,
'-o', srcTmp, '-o', srcTmp,
]); ], { signal });
// 2) Trim + concat the keep segments into the final custom video. // 2) Trim + concat the keep segments into the final custom video.
await runFfmpeg([ await runFfmpeg([
'-y', '-hide_banner', '-loglevel', 'error', '-y', '-hide_banner', '-loglevel', 'error',
@@ -787,7 +833,7 @@ app.get('/api/download/:videoId', async (c) => {
const keep = parseKeepParam(c.req.query('keep') || ''); const keep = parseKeepParam(c.req.query('keep') || '');
if (!keep.length) return c.json({ ok: false, error: 'missing or invalid keep segments' }, 400); if (!keep.length) return c.json({ ok: false, error: 'missing or invalid keep segments' }, 400);
try { try {
return await ytdlpEditedDownloadResponse(videoId, fp, keep); return await ytdlpEditedDownloadResponse(videoId, fp, keep, c.req.raw.signal);
} catch (err) { } catch (err) {
return c.json({ ok: false, error: err.message }, 500); return c.json({ ok: false, error: err.message }, 500);
} }
@@ -802,7 +848,7 @@ app.get('/api/download/:videoId', async (c) => {
return await ytdlpDownloadResponse(videoId, fp, [ return await ytdlpDownloadResponse(videoId, fp, [
'-f', 'bv*[height<=720][ext=mp4]+ba[ext=m4a]/bv*[height<=720]+ba/b[ext=mp4]/b', '-f', 'bv*[height<=720][ext=mp4]+ba[ext=m4a]/bv*[height<=720]+ba/b[ext=mp4]/b',
'--merge-output-format', 'mp4', '--merge-output-format', 'mp4',
]); ], c.req.raw.signal);
} catch (err) { } catch (err) {
console.warn(`[ytplayer] mux download failed for ${videoId}, falling back to progressive:`, err.message); console.warn(`[ytplayer] mux download failed for ${videoId}, falling back to progressive:`, err.message);
} }
@@ -818,7 +864,7 @@ app.get('/api/download/:videoId', async (c) => {
'-f', 'bestvideo[ext=mp4][acodec!=none]/bestvideo[acodec!=none]/best[ext=mp4][acodec!=none]/' '-f', 'bestvideo[ext=mp4][acodec!=none]/bestvideo[acodec!=none]/best[ext=mp4][acodec!=none]/'
+ 'bv*[height<=720][ext=mp4]+ba[ext=m4a]/bv*[height<=720]+ba/b', + 'bv*[height<=720][ext=mp4]+ba[ext=m4a]/bv*[height<=720]+ba/b',
'--merge-output-format', 'mp4', '--merge-output-format', 'mp4',
]); ], c.req.raw.signal);
} catch (err) { } catch (err) {
return c.json({ ok: false, error: err.message }, 500); return c.json({ ok: false, error: err.message }, 500);
} }