/* ============================================================================
* media-cache.js — server-side video cache fed by background jobs
*
* Every video that is played or saved gets ONE validated copy on disk:
*
/..mp4 ≤720p H.264 (8-bit 4:2:0) + AAC, faststart
* /..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= 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 /.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',
gen: row ? row.gen : null,
sha256: row && row.sha256 ? row.sha256 : null,
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 };
}