- MEDIA_DIR moves to /mnt/data/ytplayer (bind-mounted at /app/bulk); DB, lyrics, notes and uploads stay on the ytplayer-data volume - MEDIA_VOLUME_MARKER: while the marker on the drive is missing the cache writes, serves and deletes nothing, so playback streams and an unmounted drive never fills the disk under its mountpoint or loses its index on a reboot - Budget 250 GiB, keep 20 GiB free on the drive - /api/channel falls back to the uploads playlist and the channel home page when a channel has no Videos tab (auto-generated '- Topic' music channels)
736 lines
34 KiB
JavaScript
736 lines
34 KiB
JavaScript
/* ============================================================================
|
|
* media-cache.js — server-side video cache fed by background jobs
|
|
*
|
|
* Every video that is played or saved gets ONE validated copy on disk:
|
|
* <dir>/<videoId>.<gen>.mp4 ≤720p H.264 (8-bit 4:2:0) + AAC, faststart
|
|
* <dir>/<videoId>.<gen>.m4a audio-only sidecar (stream copy) for audio mode
|
|
*
|
|
* - Jobs are owned by the server, not by an HTTP request: a client closing
|
|
* its tab never kills a download, the copy just lands for next time.
|
|
* - Nothing is promoted without passing validateMedia() — a truncated or
|
|
* undecodable file must never be served "forever".
|
|
* - Copies never expire by time. They leave only through LRU eviction at the
|
|
* byte budget, or through redownload() (the client's "Broken" button).
|
|
* - `gen` bumps on every replacement and is part of the filename, so a URL
|
|
* carrying ?g=<gen> always maps to the same bytes; superseded files are kept
|
|
* for a grace period so in-flight Range playback never splices two encodes.
|
|
* - After a copy is ready, a separate low-priority lane re-encodes it with
|
|
* x264 CRF and keeps the result only when it is meaningfully smaller.
|
|
* ========================================================================== */
|
|
|
|
import { spawn } from 'node:child_process';
|
|
import {
|
|
copyFileSync, existsSync, linkSync, mkdirSync, readdirSync, renameSync,
|
|
statSync, statfsSync, unlinkSync,
|
|
} from 'node:fs';
|
|
import { join } from 'node:path';
|
|
import { sha256File } from './hash.js';
|
|
|
|
export const HIGH = 0; // explicit save / Broken re-download
|
|
export const LOW = 1; // auto-cache after a play
|
|
|
|
// Compression-lane revision. `optimized` stores the revision a copy was last
|
|
// evaluated with; bumping this re-queues every cached copy once (rev 1 was
|
|
// x264, which never beat YouTube's own encode; rev 2 is HEVC).
|
|
export const OPT_REV = 2;
|
|
|
|
const ID_RE = /^[A-Za-z0-9_-]{6,64}$/;
|
|
export const isMediaId = (id) => typeof id === 'string' && ID_RE.test(id);
|
|
|
|
const MIN_BYTES = 100 * 1024;
|
|
const OLD_GEN_GRACE_MS = 30 * 60_000; // keep superseded files for in-flight playback
|
|
const EVICT_PROTECT_MS = 10 * 60_000; // never evict something played this recently
|
|
const TOUCH_DEBOUNCE_MS = 60_000;
|
|
const EST_BYTES_PER_SEC = 250 * 1024; // ~2 Mbps: pre-compression 720p + 128k AAC
|
|
const MIN_SAVING = 0.15; // keep a re-encode only if ≥15% smaller
|
|
const SAMPLE_SECONDS = 10; // per sample window (two windows)
|
|
const SAMPLE_MIN_DURATION = 60; // shorter videos are cheap enough to encode whole
|
|
// x265 tuning from the codec benchmark: stronger adaptive quantisation and
|
|
// psy settings, no SAO blur, lighter deblock, longer lookahead.
|
|
const X265_TUNING = 'aq-mode=3:aq-strength=0.8:psy-rd=1.5:psy-rdoq=1.0:no-sao=1:deblock=-1,-1:'
|
|
+ 'bframes=8:rc-lookahead=60:ref=6:no-strong-intra-smoothing=1';
|
|
const DAY_MS = 24 * 3600_000;
|
|
|
|
// A condition that means "don't cache this right now, just stream it" — not a
|
|
// broken download. code: SKIPPED (live, too long, disk/budget) | BACKOFF.
|
|
export class MediaSkip extends Error {
|
|
constructor(message, code = 'SKIPPED') { super(message); this.code = code; }
|
|
}
|
|
|
|
const tail = (s, n = 400) => String(s || '').trim().slice(-n);
|
|
const sizeOf = (p) => { try { return statSync(p).size; } catch { return 0; } };
|
|
const fmtMB = (b) => (b / 1048576).toFixed(1) + ' MB';
|
|
|
|
// Spawn a process and collect its output. Never rejects.
|
|
function run(cmd, args) {
|
|
return new Promise((resolve) => {
|
|
let child;
|
|
try { child = spawn(cmd, args, { stdio: ['ignore', 'pipe', 'pipe'] }); }
|
|
catch (e) { resolve({ code: -1, stdout: '', stderr: e.message }); return; }
|
|
let out = '';
|
|
let err = '';
|
|
child.stdout.setEncoding('utf8');
|
|
child.stderr.setEncoding('utf8');
|
|
child.stdout.on('data', (d) => { out += d; });
|
|
child.stderr.on('data', (d) => { err = (err + d).slice(-8000); });
|
|
child.on('error', (e) => resolve({ code: -1, stdout: out, stderr: e.message }));
|
|
child.on('close', (code) => resolve({ code, stdout: out, stderr: err }));
|
|
});
|
|
}
|
|
|
|
// ----------------------------------------------------------------------------
|
|
// Validation gate. Throws with a reason; resolves to probe facts on success.
|
|
// 1. file exists and is not tiny
|
|
// 2. exactly one 8-bit 4:2:0 video stream in an allowed codec + an AAC
|
|
// audio stream. Downloads must be H.264 (AV1/VP9 play on almost nothing
|
|
// Apple — see SAVE_MUX_FORMAT); the compression lane may produce HEVC,
|
|
// which must carry the `hvc1` tag or Safari refuses it.
|
|
// 3. container AND per-stream durations cover the expected length
|
|
// 4. full demux pass reports no errors (catches corrupt/truncated boxes)
|
|
// 5. the first and last seconds actually decode
|
|
// ----------------------------------------------------------------------------
|
|
export async function validateMedia(path, expected, { ffmpeg = 'ffmpeg', ffprobe = 'ffprobe', codecs = ['h264', 'hevc'] } = {}) {
|
|
const size = sizeOf(path);
|
|
if (!size) throw new Error('validation: output file missing');
|
|
if (size < MIN_BYTES) throw new Error(`validation: file too small (${size} bytes)`);
|
|
|
|
const pr = await run(ffprobe, [
|
|
'-v', 'error',
|
|
'-show_entries', 'format=duration:stream=codec_type,codec_name,codec_tag_string,height,pix_fmt,duration:stream_disposition=attached_pic',
|
|
'-of', 'json', path,
|
|
]);
|
|
if (pr.code !== 0) throw new Error('validation: ffprobe failed: ' + tail(pr.stderr));
|
|
let j;
|
|
try { j = JSON.parse(pr.stdout); } catch { throw new Error('validation: unreadable ffprobe output'); }
|
|
const streams = Array.isArray(j.streams) ? j.streams : [];
|
|
const video = streams.filter((s) => s.codec_type === 'video' && !(s.disposition && s.disposition.attached_pic));
|
|
const audio = streams.filter((s) => s.codec_type === 'audio');
|
|
if (video.length !== 1) throw new Error(`validation: expected 1 video stream, found ${video.length}`);
|
|
const v = video[0];
|
|
if (!codecs.includes(v.codec_name)) throw new Error(`validation: video codec ${v.codec_name} is not ${codecs.join('/')}`);
|
|
if (v.codec_name === 'hevc' && v.codec_tag_string !== 'hvc1') {
|
|
throw new Error(`validation: HEVC tagged ${v.codec_tag_string} (Safari needs hvc1)`);
|
|
}
|
|
if (v.pix_fmt && !/^yuvj?420p$/.test(v.pix_fmt)) throw new Error(`validation: pixel format ${v.pix_fmt} is not 8-bit 4:2:0`);
|
|
const a = audio.find((s) => s.codec_name === 'aac');
|
|
if (!a) throw new Error('validation: no AAC audio stream');
|
|
|
|
const duration = Number(j.format && j.format.duration) || 0;
|
|
if (duration <= 0) throw new Error('validation: zero duration');
|
|
if (expected > 0) {
|
|
const min = expected - Math.max(3, expected * 0.03);
|
|
for (const [label, d] of [['file', duration], ['video', Number(v.duration)], ['audio', Number(a.duration)]]) {
|
|
if (Number.isFinite(d) && d > 0 && d < min) {
|
|
throw new Error(`validation: ${label} truncated (${d.toFixed(1)}s of ${expected}s)`);
|
|
}
|
|
}
|
|
}
|
|
|
|
const demux = await run(ffmpeg, ['-v', 'error', '-nostdin', '-i', path, '-map', '0', '-c', 'copy', '-f', 'null', '-']);
|
|
if (demux.code !== 0 || demux.stderr.trim()) throw new Error('validation: demux errors: ' + tail(demux.stderr));
|
|
const head = await run(ffmpeg, ['-v', 'error', '-nostdin', '-i', path, '-t', '3', '-f', 'null', '-']);
|
|
if (head.code !== 0 || head.stderr.trim()) throw new Error('validation: head decode failed: ' + tail(head.stderr));
|
|
const end = await run(ffmpeg, ['-v', 'error', '-nostdin', '-sseof', '-5', '-i', path, '-f', 'null', '-']);
|
|
if (end.code !== 0 || end.stderr.trim()) throw new Error('validation: tail decode failed: ' + tail(end.stderr));
|
|
|
|
return { size, height: v.height || 0, vcodec: v.codec_name, acodec: a.codec_name, duration };
|
|
}
|
|
|
|
function metaFromInfo(info, duration) {
|
|
const s = (...keys) => {
|
|
for (const k of keys) if (typeof info[k] === 'string' && info[k].trim()) return info[k].trim();
|
|
return '';
|
|
};
|
|
return {
|
|
title: s('title') || '(untitled)',
|
|
channel: s('channel', 'uploader'),
|
|
channelId: s('channel_id', 'uploader_id'),
|
|
channelUrl: s('channel_url', 'uploader_url'),
|
|
duration: typeof info.duration === 'number' ? info.duration : Math.round(duration || 0),
|
|
};
|
|
}
|
|
|
|
// ----------------------------------------------------------------------------
|
|
// createMediaCache(opts)
|
|
// dir cache directory (tmp work dir lives in <dir>/.tmp, same fs)
|
|
// db { getMedia, upsertMedia, deleteMedia, listMedia, listMediaLru, touchMedia, mediaStats }
|
|
// download async (videoId, outPath) => path — yt-dlp into outPath
|
|
// getInfo async (videoId) => yt-dlp -J info — duration, live flags, meta
|
|
// ----------------------------------------------------------------------------
|
|
export function createMediaCache({
|
|
dir, db, download, getInfo,
|
|
ffmpeg = 'ffmpeg', ffprobe = 'ffprobe', nice = 'nice',
|
|
maxBytes = 10 * 1024 ** 3,
|
|
minFreeBytes = 5 * 1024 ** 3,
|
|
autoMaxSeconds = 3 * 3600,
|
|
saveMaxSeconds = 3 * 3600,
|
|
transcode = { enabled: true, codec: 'hevc', crf: 28, preset: 'medium', threads: 2, maxSeconds: 3600 },
|
|
freeBytes = () => { const s = statfsSync(dir); return s.bavail * s.bsize; },
|
|
now = () => Date.now(),
|
|
log = console,
|
|
hashFile = sha256File, // server-computed SHA-256 = the file's P2P content id
|
|
onReady = null, // ({ id, gen, path, sha256, size, … }) after a validated copy lands
|
|
backfillDelayMs = 30_000, // hash pre-existing copies this long after init()
|
|
// Optional marker file that only exists on the real cache volume (e.g. a USB
|
|
// drive bind-mounted into the container). While it is missing the volume is
|
|
// treated as OFFLINE: nothing is written (so an unmounted drive never fills
|
|
// the disk underneath its mountpoint), nothing is served (playback streams
|
|
// instead), and nothing is deleted (a drive that comes back keeps its copies).
|
|
volumeMarker = null,
|
|
} = {}) {
|
|
const available = () => {
|
|
if (!volumeMarker) return true;
|
|
try { return existsSync(volumeMarker); } catch { return false; }
|
|
};
|
|
const tmpDir = join(dir, '.tmp');
|
|
const fileFor = (id, gen, kind = 'mp4') => join(dir, `${id}.${gen}.${kind}`);
|
|
|
|
const jobs = new Map(); // videoId -> fetch job (queued or running)
|
|
const fetchQueue = [];
|
|
let fetchActive = null; // job currently downloading/validating
|
|
let seq = 0;
|
|
const optQueue = [];
|
|
const optPending = new Set();
|
|
let optActive = null; // videoId currently being re-encoded
|
|
const lastTouch = new Map();
|
|
|
|
function sweepTmp(prefix) {
|
|
let names = [];
|
|
try { names = readdirSync(tmpDir); } catch { return; }
|
|
for (const n of names) {
|
|
if (!prefix || n.startsWith(prefix)) { try { unlinkSync(join(tmpDir, n)); } catch { /* gone */ } }
|
|
}
|
|
}
|
|
|
|
// Delete this video's files, optionally keeping one generation. Superseded
|
|
// generations can be deleted after a delay so playback in flight finishes.
|
|
function removeFiles(id, keepGen = null, delayMs = 0) {
|
|
let names = [];
|
|
try { names = readdirSync(dir); } catch { return; }
|
|
for (const n of names) {
|
|
const m = /^(.+)\.(\d+)\.(mp4|m4a)$/.exec(n);
|
|
if (!m || m[1] !== id || (keepGen !== null && Number(m[2]) === keepGen)) continue;
|
|
const p = join(dir, n);
|
|
const rm = () => { try { unlinkSync(p); } catch { /* gone */ } };
|
|
if (delayMs > 0) { const t = setTimeout(rm, delayMs); if (t && t.unref) t.unref(); }
|
|
else rm();
|
|
}
|
|
}
|
|
|
|
function readyFilesExist(row) {
|
|
return available() && existsSync(fileFor(row.video_id, row.gen)) && existsSync(fileFor(row.video_id, row.gen, 'm4a'));
|
|
}
|
|
|
|
function safeFree() {
|
|
try { return freeBytes(); } catch { return null; }
|
|
}
|
|
|
|
async function evict(id, why) {
|
|
await db.deleteMedia(id);
|
|
removeFiles(id);
|
|
log.info?.(`[media] evicted ${id} (${why})`);
|
|
}
|
|
|
|
// Make room for ~est bytes under the budget (LRU), then check the disk guard.
|
|
async function makeRoom(id, est) {
|
|
if (!available()) throw new MediaSkip('cache volume offline — streaming instead');
|
|
let { bytes } = await db.mediaStats();
|
|
if (bytes + est > maxBytes) {
|
|
const cutoff = now() - EVICT_PROTECT_MS;
|
|
// Retention order (plan 010) when the server provides it: copies that
|
|
// are neither top nor recent go first. Otherwise plain LRU.
|
|
const order = db.listMediaEvictionOrder ? await db.listMediaEvictionOrder() : await db.listMediaLru();
|
|
for (const r of order) {
|
|
if (bytes + est <= maxBytes) break;
|
|
if (r.video_id === id || r.last_access > cutoff || jobs.has(r.video_id) || optActive === r.video_id) continue;
|
|
await evict(r.video_id, 'budget');
|
|
bytes -= r.size || 0;
|
|
}
|
|
}
|
|
if (bytes + est > maxBytes) throw new MediaSkip('cache budget is full of recently played videos');
|
|
const free = safeFree();
|
|
if (free !== null && free - est < minFreeBytes) {
|
|
throw new MediaSkip(`disk guard: only ${fmtMB(free)} free on the cache volume`);
|
|
}
|
|
}
|
|
|
|
// Tell the P2P layer a validated file with a server-computed hash exists.
|
|
// Never throws and never delays the caller.
|
|
function notifyReady(id, gen, sha256, probe, metaJson) {
|
|
if (!onReady || !sha256) return;
|
|
let meta = {};
|
|
try { meta = JSON.parse(metaJson || '{}'); } catch { /* keep {} */ }
|
|
Promise.resolve()
|
|
.then(() => onReady({
|
|
id, gen, path: fileFor(id, gen), sha256, size: probe.size, height: probe.height,
|
|
vcodec: probe.vcodec, acodec: probe.acodec, duration: probe.duration, meta,
|
|
}))
|
|
.catch((e) => log.warn?.(`[media] onReady ${id} failed: ${e.message}`));
|
|
}
|
|
|
|
async function setStatus(job, status) {
|
|
job.status = status;
|
|
await db.upsertMedia(job.id, { status, updated_at: now() });
|
|
}
|
|
|
|
function bump(job, priority, auto) {
|
|
if (priority < job.priority) job.priority = priority;
|
|
if (!auto) job.auto = false; // an explicit request lifts the auto length limit
|
|
}
|
|
|
|
function startJob(id, priority, auto, row) {
|
|
let resolve, reject;
|
|
const promise = new Promise((res, rej) => { resolve = res; reject = rej; });
|
|
promise.catch(() => {}); // auto jobs have nobody awaiting them
|
|
const job = { id, priority, auto, seq: seq++, status: 'queued', promise, resolve, reject };
|
|
job.persisted = db.upsertMedia(id, {
|
|
status: 'queued', priority, auto: auto ? 1 : 0,
|
|
gen: (row && row.gen) || 0,
|
|
created_at: (row && row.created_at) || now(),
|
|
updated_at: now(),
|
|
}).catch((e) => log.warn?.(`[media] ${id} queue write failed: ${e.message}`));
|
|
jobs.set(id, job);
|
|
fetchQueue.push(job);
|
|
pumpFetch();
|
|
return promise;
|
|
}
|
|
|
|
function pumpFetch() {
|
|
if (fetchActive || !fetchQueue.length) return;
|
|
fetchQueue.sort((a, b) => a.priority - b.priority || a.seq - b.seq);
|
|
const job = fetchQueue.shift();
|
|
fetchActive = job;
|
|
runFetch(job).then(job.resolve, job.reject).finally(() => {
|
|
jobs.delete(job.id);
|
|
fetchActive = null;
|
|
pumpFetch();
|
|
});
|
|
}
|
|
|
|
async function runFetch(job) {
|
|
const { id } = job;
|
|
const prefix = `${id}-${now()}`;
|
|
const base = join(tmpDir, prefix);
|
|
await job.persisted;
|
|
try {
|
|
await setStatus(job, 'downloading');
|
|
const info = await getInfo(id);
|
|
const live = info.is_live || ['is_live', 'post_live', 'is_upcoming'].includes(info.live_status);
|
|
if (live) throw new MediaSkip('live streams are not cached');
|
|
const expected = typeof info.duration === 'number' ? info.duration : 0;
|
|
const limit = job.auto ? autoMaxSeconds : saveMaxSeconds;
|
|
if (expected > limit) {
|
|
throw new MediaSkip(`too long to cache (${Math.round(expected / 60)} min, limit ${Math.round(limit / 60)} min)`);
|
|
}
|
|
await makeRoom(id, Math.max(expected, 60) * EST_BYTES_PER_SEC);
|
|
|
|
const raw = await download(id, `${base}.dl.mp4`);
|
|
await setStatus(job, 'validating');
|
|
|
|
// Normalise: first video + first audio, moov atom up front so the copy
|
|
// starts instantly under Range requests. Stream copy — no quality loss.
|
|
const norm = `${base}.norm.mp4`;
|
|
const rm = await run(ffmpeg, ['-v', 'error', '-nostdin', '-y', '-i', raw,
|
|
'-map', '0:v:0', '-map', '0:a:0', '-c', 'copy', '-movflags', '+faststart', '-f', 'mp4', norm]);
|
|
if (rm.code !== 0) throw new Error('remux failed: ' + tail(rm.stderr));
|
|
const probe = await validateMedia(norm, expected, { ffmpeg, ffprobe, codecs: ['h264'] });
|
|
|
|
const m4a = `${base}.norm.m4a`;
|
|
const ex = await run(ffmpeg, ['-v', 'error', '-nostdin', '-y', '-i', norm,
|
|
'-map', '0:a:0', '-c', 'copy', '-movflags', '+faststart', '-f', 'mp4', m4a]);
|
|
if (ex.code !== 0 || sizeOf(m4a) < 1024) throw new Error('audio sidecar failed: ' + tail(ex.stderr));
|
|
|
|
const sha256 = await hashFile(norm);
|
|
const prev = await db.getMedia(id);
|
|
const gen = ((prev && prev.gen) || 0) + 1;
|
|
renameSync(norm, fileFor(id, gen));
|
|
renameSync(m4a, fileFor(id, gen, 'm4a'));
|
|
const t = now();
|
|
const fields = {
|
|
status: 'ready', gen,
|
|
size: probe.size + sizeOf(fileFor(id, gen, 'm4a')),
|
|
height: probe.height, vcodec: probe.vcodec, acodec: probe.acodec, duration: probe.duration,
|
|
optimized: 0, meta: JSON.stringify(metaFromInfo(info, probe.duration)),
|
|
attempts: 0, error: null, retry_at: 0, updated_at: t, last_access: t, sha256,
|
|
};
|
|
await db.upsertMedia(id, fields);
|
|
removeFiles(id, gen);
|
|
job.status = 'ready';
|
|
log.info?.(`[media] cached ${id} ${probe.height}p ${fmtMB(fields.size)}`);
|
|
enqueueOptimize(id);
|
|
notifyReady(id, gen, sha256, probe, fields.meta);
|
|
return { ...(prev || {}), video_id: id, ...fields };
|
|
} catch (err) {
|
|
job.status = 'failed';
|
|
if (err instanceof MediaSkip) {
|
|
await db.deleteMedia(id).catch(() => {});
|
|
log.info?.(`[media] skipped ${id}: ${err.message}`);
|
|
} else {
|
|
const prev = await db.getMedia(id).catch(() => null);
|
|
const attempts = ((prev && prev.attempts) || 0) + 1;
|
|
await db.upsertMedia(id, {
|
|
status: 'failed', attempts, error: String(err.message || err).slice(0, 500),
|
|
retry_at: now() + Math.min(DAY_MS, 15 * 60_000 * 2 ** (attempts - 1)),
|
|
updated_at: now(),
|
|
}).catch(() => {});
|
|
log.warn?.(`[media] ${id} failed (attempt ${attempts}): ${err.message}`);
|
|
}
|
|
throw err;
|
|
} finally {
|
|
sweepTmp(prefix);
|
|
}
|
|
}
|
|
|
|
// ---- Compression lane ----------------------------------------------------
|
|
// Re-encodes a ready copy with CRF (quality-targeted) and keeps the result
|
|
// only when it is ≥15% smaller. Measured 2026-09-13: x264 can't beat
|
|
// YouTube's own H.264 (it needed MORE bitrate for the same VMAF), while
|
|
// x265 8-bit medium CRF 28 cut live worship video ~44% (VMAF ~93) at ~2x
|
|
// realtime on the homelab's 2 niced threads. HEVC plays on iPhone Safari,
|
|
// Android and Chrome/Edge with hardware decode; the server only serves an
|
|
// HEVC copy to clients that ask for it (?hevc=1), everyone else streams.
|
|
// 8-bit on purpose — the widest hardware-decode support. 60 fps is capped
|
|
// to 30; audio is copied untouched.
|
|
function enqueueOptimize(id) {
|
|
if (!transcode || !transcode.enabled || optPending.has(id)) return;
|
|
optPending.add(id);
|
|
optQueue.push(id);
|
|
pumpOpt();
|
|
}
|
|
|
|
function pumpOpt() {
|
|
if (optActive || !optQueue.length) return;
|
|
const id = optQueue.shift();
|
|
optActive = id;
|
|
runOptimize(id)
|
|
.catch((err) => log.warn?.(`[media] optimize ${id} failed: ${err.message}`))
|
|
.finally(() => { optPending.delete(id); optActive = null; pumpOpt(); });
|
|
}
|
|
|
|
async function runOptimize(id) {
|
|
const row = await db.getMedia(id);
|
|
if (!row || row.status !== 'ready' || (row.optimized || 0) >= OPT_REV) return;
|
|
const src = fileFor(id, row.gen);
|
|
if (!existsSync(src)) return;
|
|
const markDone = async () => {
|
|
const cur = await db.getMedia(id);
|
|
if (cur && cur.status === 'ready' && cur.gen === row.gen) await db.upsertMedia(id, { optimized: OPT_REV });
|
|
};
|
|
// Very long videos would tie the lane up for hours on 2 niced threads.
|
|
if (transcode.maxSeconds && row.duration > transcode.maxSeconds) { await markDone(); return; }
|
|
const hevc = transcode.codec !== 'h264';
|
|
const prefix = `${id}-${now()}-${hevc ? 'x265' : 'x264'}`;
|
|
const out = join(tmpDir, `${prefix}.mp4`);
|
|
const threads = String(transcode.threads);
|
|
const video = hevc
|
|
? ['-c:v', 'libx265', '-preset', String(transcode.preset), '-crf', String(transcode.crf),
|
|
'-pix_fmt', 'yuv420p', '-tag:v', 'hvc1',
|
|
'-x265-params', `log-level=error:pools=${threads}:${X265_TUNING}`]
|
|
: ['-c:v', 'libx264', '-preset', String(transcode.preset), '-crf', String(transcode.crf),
|
|
'-profile:v', 'high', '-level:v', '4.0', '-pix_fmt', 'yuv420p', '-threads', threads];
|
|
try {
|
|
// Sample first. Re-encoding YouTube's already-compressed H.264 usually
|
|
// needs MORE bits than YouTube spent (measured on prod: 0 of 5 videos
|
|
// shrank at CRF 28; one grew 25%), so spend ~20 s of video finding out
|
|
// before ~2x realtime of CPU on the whole thing.
|
|
if (row.duration >= SAMPLE_MIN_DURATION) {
|
|
const ratio = await sampleRatio(src, row.duration, video, prefix);
|
|
const mp4 = sizeOf(src);
|
|
const videoBytes = Math.max(0, mp4 - sizeOf(fileFor(id, row.gen, 'm4a')));
|
|
const predicted = mp4 - videoBytes * (1 - ratio);
|
|
if (predicted > mp4 * (1 - MIN_SAVING)) {
|
|
await markDone();
|
|
log.info?.(`[media] kept original ${id}: sample predicts ${Math.round(predicted / mp4 * 100)}% of ${fmtMB(mp4)}`);
|
|
return;
|
|
}
|
|
}
|
|
const args = ['-v', 'error', '-nostdin', '-y', '-i', src,
|
|
'-map', '0:v:0', '-map', '0:a:0',
|
|
'-vf', 'scale=-2:min(720\\,ih)',
|
|
...video, '-fpsmax', '30',
|
|
'-c:a', 'copy', '-movflags', '+faststart', '-f', 'mp4', out];
|
|
const r = nice ? await run(nice, ['-n', '19', ffmpeg, ...args]) : await run(ffmpeg, args);
|
|
if (r.code !== 0) throw new Error(`${hevc ? 'x265' : 'x264'} encode failed: ` + tail(r.stderr));
|
|
const probe = await validateMedia(out, row.duration, { ffmpeg, ffprobe, codecs: [hevc ? 'hevc' : 'h264'] });
|
|
|
|
const cur = await db.getMedia(id);
|
|
if (!cur || cur.status !== 'ready' || cur.gen !== row.gen) return; // replaced or evicted meanwhile
|
|
const oldSize = sizeOf(src);
|
|
if (probe.size > oldSize * (1 - MIN_SAVING)) {
|
|
await db.upsertMedia(id, { optimized: OPT_REV });
|
|
log.info?.(`[media] kept original ${id}: ${probe.vcodec} ${fmtMB(probe.size)} vs ${fmtMB(oldSize)}`);
|
|
return;
|
|
}
|
|
const sha256 = await hashFile(out);
|
|
const gen = row.gen + 1;
|
|
renameSync(out, fileFor(id, gen));
|
|
const oldAudio = fileFor(id, row.gen, 'm4a');
|
|
const newAudio = fileFor(id, gen, 'm4a');
|
|
try { linkSync(oldAudio, newAudio); } catch { copyFileSync(oldAudio, newAudio); }
|
|
await db.upsertMedia(id, {
|
|
gen, size: probe.size + sizeOf(newAudio), height: probe.height, vcodec: probe.vcodec,
|
|
duration: probe.duration, optimized: OPT_REV, updated_at: now(), sha256,
|
|
});
|
|
notifyReady(id, gen, sha256, probe, cur.meta);
|
|
removeFiles(id, gen, OLD_GEN_GRACE_MS);
|
|
log.info?.(`[media] optimized ${id}: ${fmtMB(oldSize)} → ${fmtMB(probe.size)} (${probe.vcodec})`);
|
|
} catch (err) {
|
|
await markDone().catch(() => {}); // don't retry a source the encoder can't improve
|
|
throw err;
|
|
} finally {
|
|
sweepTmp(prefix);
|
|
}
|
|
}
|
|
|
|
// Encode two short windows (at 25% and 60% of the video) and compare them
|
|
// with the source's own video bytes over exactly the same timestamps.
|
|
// Returns encoded/source (below 1 = smaller).
|
|
async function sampleRatio(src, duration, videoArgs, prefix) {
|
|
let srcBytes = 0;
|
|
let encBytes = 0;
|
|
const starts = [0.25, 0.6].map((f) => Math.floor(duration * f));
|
|
for (const [i, start] of starts.entries()) {
|
|
const end = start + SAMPLE_SECONDS;
|
|
const pr = await run(ffprobe, ['-v', 'error', '-select_streams', 'v:0',
|
|
'-read_intervals', `${Math.max(0, start - 10)}%${end + 1}`,
|
|
'-show_entries', 'packet=pts_time,size', '-of', 'csv=p=0', src]);
|
|
if (pr.code !== 0) throw new Error('sample probe failed: ' + tail(pr.stderr));
|
|
for (const line of pr.stdout.split('\n')) {
|
|
const [pts, size] = line.split(',');
|
|
const t = parseFloat(pts);
|
|
if (t >= start && t < end) srcBytes += parseInt(size, 10) || 0;
|
|
}
|
|
const out = join(tmpDir, `${prefix}.sample${i}.mp4`);
|
|
const args = ['-v', 'error', '-nostdin', '-y', '-ss', String(start), '-t', String(SAMPLE_SECONDS), '-i', src,
|
|
'-map', '0:v:0', '-an', '-vf', 'scale=-2:min(720\\,ih)', ...videoArgs, '-fpsmax', '30', '-f', 'mp4', out];
|
|
const r = nice ? await run(nice, ['-n', '19', ffmpeg, ...args]) : await run(ffmpeg, args);
|
|
if (r.code !== 0) throw new Error('sample encode failed: ' + tail(r.stderr));
|
|
encBytes += sizeOf(out);
|
|
}
|
|
return srcBytes > 0 ? encBytes / srcBytes : 1;
|
|
}
|
|
|
|
// ---- Public API ----------------------------------------------------------
|
|
|
|
// Resolve to the ready row, fetching it first if needed. Rejects with a
|
|
// MediaSkip (stream instead) or the download/validation error.
|
|
async function ensureCached(id, { priority = LOW, auto = false, force = false } = {}) {
|
|
if (!isMediaId(id)) throw new MediaSkip('invalid video id');
|
|
const existing = jobs.get(id);
|
|
if (existing) { bump(existing, priority, auto); return existing.promise; }
|
|
const row = await db.getMedia(id);
|
|
const again = jobs.get(id);
|
|
if (again) { bump(again, priority, auto); return again.promise; }
|
|
if (row && row.status === 'ready') {
|
|
if (readyFilesExist(row)) return row;
|
|
if (!available()) throw new MediaSkip('cache volume offline — streaming instead');
|
|
await db.deleteMedia(id);
|
|
removeFiles(id);
|
|
} else if (row && row.status === 'failed' && !force && now() < row.retry_at) {
|
|
throw new MediaSkip(`cache fetch failed recently (${row.error || 'unknown error'})`, 'BACKOFF');
|
|
}
|
|
if (jobs.get(id)) { bump(jobs.get(id), priority, auto); return jobs.get(id).promise; }
|
|
return startJob(id, priority, auto, row);
|
|
}
|
|
|
|
// The "Broken" button: drop the copy and fetch it again at high priority.
|
|
// Promote a file that came from somewhere else (P2P intake / rehydrate from
|
|
// a device) as this video's server copy. The caller has already hashed it
|
|
// and run validateMedia(). The bytes are kept EXACTLY (no remux — the cid
|
|
// must stay true); only the audio sidecar is derived. Never replaces an
|
|
// existing ready copy or races a running fetch.
|
|
async function adoptFile(id, src, { sha256, probe, meta = {} }) {
|
|
if (!isMediaId(id)) throw new Error('adopt: bad video id');
|
|
const cur = await db.getMedia(id);
|
|
if (cur && cur.status === 'ready' && readyFilesExist(cur)) return { adopted: false, reason: 'already cached' };
|
|
if (jobs.has(id)) return { adopted: false, reason: 'fetch in progress' };
|
|
await makeRoom(id, probe.size);
|
|
const prefix = `${id}-${now()}-adopt`;
|
|
const m4a = join(tmpDir, `${prefix}.m4a`);
|
|
try {
|
|
const ex = await run(ffmpeg, ['-v', 'error', '-nostdin', '-y', '-i', src,
|
|
'-map', '0:a:0', '-c', 'copy', '-movflags', '+faststart', '-f', 'mp4', m4a]);
|
|
if (ex.code !== 0 || sizeOf(m4a) < 1024) throw new Error('audio sidecar failed: ' + tail(ex.stderr));
|
|
const gen = ((cur && cur.gen) || 0) + 1;
|
|
try { renameSync(src, fileFor(id, gen)); } catch { copyFileSync(src, fileFor(id, gen)); unlinkSync(src); }
|
|
renameSync(m4a, fileFor(id, gen, 'm4a'));
|
|
const t = now();
|
|
await db.upsertMedia(id, {
|
|
status: 'ready', gen, size: probe.size + sizeOf(fileFor(id, gen, 'm4a')),
|
|
height: probe.height, vcodec: probe.vcodec, acodec: probe.acodec, duration: probe.duration,
|
|
optimized: OPT_REV, // re-encoding would change the bytes other devices hold
|
|
meta: JSON.stringify(meta), priority: HIGH, auto: 0, attempts: 0, error: null, retry_at: 0,
|
|
created_at: (cur && cur.created_at) || t, updated_at: t, last_access: t, sha256,
|
|
});
|
|
removeFiles(id, gen);
|
|
log.info?.(`[media] adopted ${id} ${probe.height}p ${fmtMB(probe.size)} from P2P`);
|
|
return { adopted: true, gen };
|
|
} finally {
|
|
sweepTmp(prefix);
|
|
}
|
|
}
|
|
|
|
async function redownload(id) {
|
|
if (!isMediaId(id)) throw new MediaSkip('invalid video id');
|
|
if (!available()) throw new MediaSkip('cache volume offline — streaming instead');
|
|
const existing = jobs.get(id);
|
|
if (existing) { bump(existing, HIGH, false); return { status: existing.status }; }
|
|
const row = await db.getMedia(id);
|
|
if (jobs.get(id)) return { status: jobs.get(id).status };
|
|
removeFiles(id);
|
|
await db.upsertMedia(id, {
|
|
status: 'queued', gen: ((row && row.gen) || 0) + 1, size: 0, optimized: 0,
|
|
attempts: 0, error: null, retry_at: 0, updated_at: now(),
|
|
created_at: (row && row.created_at) || now(),
|
|
});
|
|
log.info?.(`[media] ${id} reported broken — re-downloading`);
|
|
startJob(id, HIGH, false, { ...(row || {}), gen: ((row && row.gen) || 0) + 1 });
|
|
return { status: 'queued' };
|
|
}
|
|
|
|
// Ready row whose files are on disk, else null. Hot path for /api/streams.
|
|
async function getReady(id) {
|
|
if (!isMediaId(id)) return null;
|
|
const row = await db.getMedia(id);
|
|
return row && row.status === 'ready' && readyFilesExist(row) ? row : null;
|
|
}
|
|
|
|
// Path for /api/media. With a gen, ONLY that generation is served (never
|
|
// different bytes under the same URL); without one, the current copy.
|
|
async function filePath(id, gen, kind = 'mp4') {
|
|
if (!isMediaId(id) || !available()) return null;
|
|
let g = gen !== undefined && gen !== null && gen !== '' ? Number(gen) : null;
|
|
if (g !== null && !Number.isInteger(g)) return null;
|
|
if (g === null) {
|
|
const row = await db.getMedia(id);
|
|
if (!row || row.status !== 'ready') return null;
|
|
g = row.gen;
|
|
}
|
|
const p = fileFor(id, g, kind);
|
|
return existsSync(p) ? p : null;
|
|
}
|
|
|
|
// A client fell back from our copy (?nocache=1). Its decode failure may be
|
|
// device-specific, so re-run the server-side gate instead of trusting it:
|
|
// a copy that fails is dropped and refetched; one that passes stays.
|
|
const lastVerify = new Map();
|
|
async function verify(id) {
|
|
if (!isMediaId(id)) return null;
|
|
const t = now();
|
|
if (t - (lastVerify.get(id) || 0) < 10 * 60_000) return null;
|
|
lastVerify.set(id, t);
|
|
if (lastVerify.size > 5000) lastVerify.clear();
|
|
const row = await getReady(id);
|
|
if (!row || jobs.has(id)) return null;
|
|
try {
|
|
await validateMedia(fileFor(id, row.gen), row.duration, { ffmpeg, ffprobe });
|
|
return true;
|
|
} catch (err) {
|
|
const cur = await db.getMedia(id);
|
|
if (!cur || cur.gen !== row.gen) return null; // replaced meanwhile
|
|
log.warn?.(`[media] ${id} copy failed re-validation (${err.message}) — re-downloading`);
|
|
await redownload(id);
|
|
return false;
|
|
}
|
|
}
|
|
|
|
function touch(id) {
|
|
const t = now();
|
|
if (t - (lastTouch.get(id) || 0) < TOUCH_DEBOUNCE_MS) return;
|
|
lastTouch.set(id, t);
|
|
if (lastTouch.size > 5000) lastTouch.clear();
|
|
db.touchMedia(id, t).catch(() => {});
|
|
}
|
|
|
|
async function status(id) {
|
|
if (!isMediaId(id)) return { status: 'none' };
|
|
const job = jobs.get(id);
|
|
const row = await db.getMedia(id);
|
|
return {
|
|
status: job ? job.status : row ? row.status : 'none',
|
|
size: row ? row.size : 0,
|
|
height: row ? row.height : 0,
|
|
vcodec: row ? row.vcodec : null,
|
|
optimized: !!(row && (row.optimized || 0) >= OPT_REV),
|
|
optimizing: optActive === id,
|
|
error: row && row.status === 'failed' ? row.error : null,
|
|
retryAt: row && row.status === 'failed' ? row.retry_at : 0,
|
|
};
|
|
}
|
|
|
|
async function stats() {
|
|
return {
|
|
...(await db.mediaStats()),
|
|
maxBytes,
|
|
freeBytes: safeFree(),
|
|
queued: fetchQueue.length,
|
|
active: fetchActive ? fetchActive.id : null,
|
|
optimizing: optActive,
|
|
optimizeQueued: optQueue.length,
|
|
};
|
|
}
|
|
|
|
// Boot: sweep partial work, drop files/rows that disagree, resume jobs.
|
|
async function init() {
|
|
if (!available()) {
|
|
// Reconciling now would drop every row (files "missing") and then, once
|
|
// the drive is back, delete every file as an orphan. Leave it all alone.
|
|
log.warn?.(`[media] cache volume offline (no ${volumeMarker}) — caching paused, playback streams`);
|
|
return;
|
|
}
|
|
mkdirSync(tmpDir, { recursive: true });
|
|
sweepTmp();
|
|
const rows = await db.listMedia();
|
|
const byId = new Map(rows.map((r) => [r.video_id, r]));
|
|
for (const name of readdirSync(dir)) {
|
|
if (name === '.tmp') continue;
|
|
const m = /^(.+)\.(\d+)\.(mp4|m4a)$/.exec(name);
|
|
const r = m && byId.get(m[1]);
|
|
if (!m || !r || r.status !== 'ready' || r.gen !== Number(m[2])) {
|
|
try { unlinkSync(join(dir, name)); } catch { /* gone */ }
|
|
}
|
|
}
|
|
let resumed = 0;
|
|
for (const r of rows) {
|
|
if (r.status === 'ready') {
|
|
if (!readyFilesExist(r)) { await db.deleteMedia(r.video_id); removeFiles(r.video_id); continue; }
|
|
if ((r.optimized || 0) < OPT_REV) enqueueOptimize(r.video_id);
|
|
} else if (r.status === 'queued' || r.status === 'downloading' || r.status === 'validating') {
|
|
startJob(r.video_id, r.priority === HIGH ? HIGH : LOW, !!r.auto, r);
|
|
resumed++;
|
|
}
|
|
}
|
|
const s = await db.mediaStats();
|
|
log.info?.(`[media] ${s.count} cached (${fmtMB(s.bytes)}), ${resumed} job(s) resumed, dir ${dir}`);
|
|
if (backfillDelayMs >= 0) {
|
|
const t = setTimeout(() => backfillHashes().catch(() => {}), backfillDelayMs);
|
|
t.unref?.();
|
|
}
|
|
}
|
|
|
|
// Copies cached before hashing existed get their sha256 in the background,
|
|
// one at a time, so their content ids can be admitted too.
|
|
async function backfillHashes() {
|
|
let n = 0;
|
|
for (const r of await db.listMedia()) {
|
|
if (r.status !== 'ready' || r.sha256) continue;
|
|
const p = fileFor(r.video_id, r.gen);
|
|
if (!existsSync(p)) continue;
|
|
try {
|
|
const sha256 = await hashFile(p);
|
|
const cur = await db.getMedia(r.video_id);
|
|
if (!cur || cur.status !== 'ready' || cur.gen !== r.gen) continue;
|
|
await db.upsertMedia(r.video_id, { sha256 });
|
|
notifyReady(r.video_id, r.gen, sha256,
|
|
{ size: sizeOf(p), height: r.height, vcodec: r.vcodec, acodec: r.acodec, duration: r.duration }, r.meta);
|
|
n++;
|
|
} catch (e) {
|
|
log.warn?.(`[media] hash backfill ${r.video_id}: ${e.message}`);
|
|
}
|
|
}
|
|
return n;
|
|
}
|
|
|
|
return { init, ensureCached, redownload, verify, getReady, filePath, touch, status, stats, isMediaId, backfillHashes, adoptFile, available };
|
|
}
|