2222 lines
97 KiB
JavaScript
2222 lines
97 KiB
JavaScript
/* ============================================================================
|
|
* server.js — YT Player PWA backend
|
|
*
|
|
* Runtime: Bun (https://bun.sh)
|
|
* Framework: Hono v4
|
|
* DB: libsql (concurrent SQLite fork, embedded file mode)
|
|
*
|
|
* Endpoints:
|
|
* GET /api/search?q=<query> yt-dlp search → slim card array
|
|
* GET /api/channel?c=<channel> yt-dlp channel uploads → slim card array
|
|
* GET /api/streams?v=<videoId> yt-dlp stream info → {meta, audioUrl, qualities} (proxied URLs)
|
|
* GET /api/play?v=<id>&f=<fmt> same-origin playback proxy (Range-aware) → media bytes
|
|
* GET /api/download/:videoId server-cached copy (fetched by a background job) → binary
|
|
* GET /api/media/:id?g=<gen>[&a=1] Range-aware server-cached mp4 (or m4a audio sidecar)
|
|
* GET /api/media/:id/peaks loudness envelope (400 buckets, 0..100) for the waveform seek bar
|
|
* GET /api/media/:id/gif?t=&d=&w= short looping GIF of a moment (timestamp sharing)
|
|
* GET /api/media/:id/clip?start=&end=&fmt=mp3|m4r soundbite / iPhone ringtone
|
|
* GET /api/media/:id/status { status, size, height, optimized, error }
|
|
* POST /api/media/:id/redownload "Broken" button: drop the copy, fetch it again
|
|
* GET /api/media/stats cache totals, budget, free disk, queue
|
|
* GET /api/version { version }
|
|
* POST /api/user/sync upsert user playlists + last-seen version
|
|
* GET /api/user/data?fp=<fp> retrieve stored playlists + history
|
|
* POST /api/playlist/share share a single playlist → { ok, code }
|
|
* GET /api/playlist/shared?code=<code> retrieve shared playlist → { ok, code, playlist, createdAt }
|
|
* GET /api/playlist/expand?url=<url> expand YouTube playlist → { ok, title, entries, truncated? }
|
|
* GET /* serve frontend/public static files
|
|
*
|
|
* JSON shapes mirror the Tauri (Rust) bridge exactly so the existing app.js
|
|
* UI code works without modification in WEB mode.
|
|
* ========================================================================== */
|
|
|
|
import { Hono } from 'hono';
|
|
import { serveStatic } from 'hono/bun';
|
|
import { logger } from 'hono/logger';
|
|
import { spawn } from 'node:child_process';
|
|
import { createServer } from 'node:http';
|
|
import { readFileSync, readdirSync, existsSync, statSync, openSync, unlinkSync, createReadStream, mkdirSync } from 'node:fs';
|
|
import { Readable } from 'node:stream';
|
|
import { tmpdir } from 'node:os';
|
|
import { createHash } from 'node:crypto';
|
|
import { brotliCompressSync, constants as zlibConstants } from 'node:zlib';
|
|
import { initDb, upsertUser, recordVideoAccess, getUserData, createProfile, getProfile, saveProfile, createSharedPlaylist, getSharedPlaylist, queueInboxPlaylist, listInbox, deleteInboxItem, countInbox,
|
|
getMedia, upsertMedia, deleteMedia, listMedia, listMediaLru, touchMedia, mediaStats } from './db.js';
|
|
import { createMediaCache, HIGH, LOW, validateMedia } from './media-cache.js';
|
|
import * as notesDb from './db.js';
|
|
import { registerNoteRoutes, parseLrc, sanitizeLyrics } from './notes.js';
|
|
import { createRemoteHub } from './remote.js';
|
|
import { createPartyHub } from './party.js';
|
|
import { registerUploadRoutes } from './uploads.js';
|
|
import { admitFile } from './p2p-admit.js';
|
|
import { P2P } from './p2p-config.js';
|
|
import * as p2pDb from './p2p-db.js';
|
|
import { registerP2pRoutes } from './p2p-routes.js';
|
|
import { registerIntakeRoutes } from './p2p-intake.js';
|
|
import { createP2pHub, holdersPayload, createRehydrator } from './p2p-hub.js';
|
|
import { sha256Range } from './hash.js';
|
|
import * as innertube from './innertube.js';
|
|
import QRCode from 'qrcode';
|
|
import { createYtdlpPool } from './ytdlp-pool.js';
|
|
import { dirname, join as pathJoin } from 'node:path';
|
|
|
|
// A media proxy must not die because one client's stream hit an edge case
|
|
// (see /api/play cancel()): log and keep serving instead of crash-looping.
|
|
process.on('uncaughtException', (err) => console.error('[ytplayer] uncaught exception:', err));
|
|
process.on('unhandledRejection', (err) => console.error('[ytplayer] unhandled rejection:', err));
|
|
|
|
const PORT = parseInt(process.env.PORT || '3000', 10);
|
|
const APP_VERSION = process.env.APP_VERSION || '1.0.0';
|
|
const YTDLP = process.env.YTDLP_PATH || 'yt-dlp';
|
|
const FFMPEG = process.env.FFMPEG_PATH || 'ffmpeg';
|
|
const FFPROBE = process.env.FFPROBE_PATH || 'ffprobe';
|
|
// Cap for server-side SAVE downloads. yt-dlp with -N 4 pinned the homelab's
|
|
// whole downlink (~7.6 MB/s measured), and since /api/play fetches its own
|
|
// googlevideo slices over the same link, one long save starved every
|
|
// concurrent playback (stalled at ~27 s, 11 KB/s). Leave headroom.
|
|
const DOWNLOAD_RATE = process.env.DOWNLOAD_RATE || '2M';
|
|
|
|
// Saves run ONE AT A TIME. Two concurrent rate-capped saves plus playback
|
|
// still filled the homelab's ~7.5 MB/s downlink and playback starved, so
|
|
// additional saves wait their turn (the client just sees a longer save).
|
|
let saveChain = Promise.resolve();
|
|
function withSaveSlot(fn) {
|
|
const run = saveChain.then(fn, fn);
|
|
saveChain = run.catch(() => {});
|
|
return run;
|
|
}
|
|
|
|
// ----------------------------------------------------------------------------
|
|
// BUILD_TAG — must be DETERMINISTIC across restarts of identical code.
|
|
//
|
|
// Previously this was `Date.now().toString(36)`, which changes every time the
|
|
// process starts even if nothing was deployed (crash-loop, healthcheck
|
|
// restart, container reschedule). The frontend's checkBuildTag() polls
|
|
// /api/version and re-shows the "Update available" modal the instant the tag
|
|
// drifts — so a restarting-but-unchanged server kept re-announcing an update
|
|
// that never actually happened, and clicking "Refresh UI" (which itself
|
|
// reloads the page and re-polls) never made the prompt go away for good.
|
|
//
|
|
// Fix: hash the actual served frontend files. Identical code → identical
|
|
// hash → identical tag, no matter how many times the process restarts. A
|
|
// real deploy (changed files) still produces a new tag as intended.
|
|
// process.env.BUILD_TAG still wins if a CI pipeline already injects a git
|
|
// SHA — that's an even better source of truth than a content hash.
|
|
// ----------------------------------------------------------------------------
|
|
function computeBuildTag() {
|
|
try {
|
|
// Hash EVERY served frontend file (recursively, in sorted order), not a
|
|
// hand-picked subset — a change to any shell file (e.g. sw-update.js or
|
|
// opfs.js) must produce a new tag, or clients keep their old SW cache
|
|
// and never receive the change.
|
|
const hash = createHash('sha256');
|
|
const walk = (dir) => {
|
|
for (const name of readdirSync(dir).sort()) {
|
|
const path = `${dir}/${name}`;
|
|
if (statSync(path).isDirectory()) walk(path);
|
|
else { hash.update(path); hash.update(readFileSync(path)); }
|
|
}
|
|
};
|
|
walk('./public');
|
|
return hash.digest('hex').slice(0, 12);
|
|
} catch {
|
|
// Frontend files not readable (e.g. unit tests run outside ./public) —
|
|
// fall back to a fixed tag rather than Date.now(), so it still never
|
|
// drifts spuriously between restarts.
|
|
return 'dev-build';
|
|
}
|
|
}
|
|
|
|
const BUILD_TAG = process.env.BUILD_TAG || computeBuildTag();
|
|
|
|
// BUILD_TIME — human-readable "when was this image built". Written by the
|
|
// Dockerfile at image build time (never at container start, so restarts
|
|
// don't drift it). Kept OUTSIDE ./public so it can't perturb BUILD_TAG.
|
|
const BUILD_TIME = process.env.BUILD_TIME || (() => {
|
|
try { return readFileSync('./build-time.txt', 'utf8').trim(); }
|
|
catch { return null; }
|
|
})();
|
|
|
|
const SEARCH_LIMIT = 25;
|
|
const CHANNEL_LIMIT = 60;
|
|
|
|
// ============================================================================
|
|
// yt-dlp helpers
|
|
// ============================================================================
|
|
|
|
// Run yt-dlp asynchronously and resolve stdout as a string.
|
|
// MUST stay async (spawn, not spawnSync): a sync child process blocks Bun's
|
|
// event loop for the full yt-dlp runtime (~2-3s per call), which stalls every
|
|
// concurrent request — including in-flight /api/download proxy streams, which
|
|
// Bun then kills at its idle timeout ("fetch failed" mid-download on clients).
|
|
// Rejects on non-zero exit.
|
|
function runYtdlpSpawn(args, { signal } = {}) {
|
|
return new Promise((resolve, reject) => {
|
|
const child = spawn(YTDLP, args, { stdio: ['ignore', 'pipe', 'pipe'] });
|
|
const t0 = Date.now();
|
|
const kind = String(args.find((a) => /^ytsearch|^https?:/.test(String(a))) || args[0] || '')
|
|
.replace(/^ytsearch\d*:.*/, 'search').replace(/^https?:\/\/[^/]+\/watch.*/, 'video').slice(0, 40);
|
|
// 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 err = '';
|
|
child.stdout.setEncoding('utf8');
|
|
child.stderr.setEncoding('utf8');
|
|
child.stdout.on('data', (d) => { out += d; });
|
|
child.stderr.on('data', (d) => { err += d; });
|
|
child.on('error', (e) => reject(new Error('yt-dlp not found: ' + e.message)));
|
|
child.on('close', (code) => {
|
|
console.log(`[ytdlp] ${kind} ${Date.now() - t0}ms ${code === 0 ? 'ok' : 'fail'}`);
|
|
if (code !== 0) reject(new Error(err.trim() || 'yt-dlp exited with code ' + code));
|
|
else resolve(out);
|
|
});
|
|
});
|
|
}
|
|
|
|
// Read-only calls (-J, --dump-json) go to long-lived workers that import
|
|
// yt_dlp once (~1 s saved per call). Downloads keep spawning. Pool trouble
|
|
// (not a yt-dlp error) falls back to a spawn. YTDLP_WORKER=0 disables it.
|
|
const ytdlpPool = process.env.YTDLP_WORKER === '0' ? null : createYtdlpPool({
|
|
ytdlpPath: Bun.which(YTDLP) || YTDLP,
|
|
size: Math.max(1, Number(process.env.YTDLP_WORKERS) || 2),
|
|
});
|
|
function runYtdlp(args, opts = {}) {
|
|
if (!opts.pooled || !ytdlpPool || opts.signal) return runYtdlpSpawn(args, opts);
|
|
const t0 = Date.now();
|
|
return ytdlpPool.run(args).then(
|
|
(out) => { console.log(`[ytdlp] pooled ${Date.now() - t0}ms ok`); return out; },
|
|
(err) => {
|
|
if (err.poolInfra) return runYtdlpSpawn(args, opts);
|
|
// A bot check can stick to a long-lived process: replace the workers and
|
|
// answer this call the old way, from a fresh process.
|
|
if (BOT_CHECK_RE.test(err.message)) { ytdlpPool.recycle(); return runYtdlpSpawn(args, opts); }
|
|
console.log(`[ytdlp] pooled ${Date.now() - t0}ms fail`);
|
|
throw err;
|
|
},
|
|
);
|
|
}
|
|
|
|
// YouTube intermittently answers the default (web) innertube client with
|
|
// "Sign in to confirm you're not a bot" — a per-IP rate signal, not a
|
|
// per-video one, so the SAME video that just saved fine fails minutes later
|
|
// and the user sees "sign in required". Other player clients are not gated by
|
|
// that check from a datacentre IP, so retry the whole yt-dlp call against each
|
|
// in turn instead of demanding cookies. Order is quality-first: web_embedded
|
|
// still exposes the adaptive DASH ladder (399+251), while tv_simply /
|
|
// android_vr / mweb typically only offer progressive itag 18 (360p) — a 360p
|
|
// save beats a failed save.
|
|
const BOT_CHECK_RE = /Sign in to confirm|not a bot|confirm you.{0,3}re not a bot/i;
|
|
const FALLBACK_CLIENTS = (process.env.YTDLP_FALLBACK_CLIENTS
|
|
|| 'web_embedded,tv_simply,android_vr,mweb').split(',').map((s) => s.trim()).filter(Boolean);
|
|
// Optional cookies jar (Netscape format) for the rare case every client is
|
|
// gated. Mounted read-only; absent by default and never required.
|
|
const YTDLP_COOKIES = process.env.YTDLP_COOKIES || '';
|
|
|
|
function withCookies(args) {
|
|
if (!YTDLP_COOKIES || !existsSync(YTDLP_COOKIES)) return args;
|
|
return ['--cookies', YTDLP_COOKIES, ...args];
|
|
}
|
|
|
|
// runYtdlp + bot-check fallback. Every YouTube-facing call goes through this.
|
|
async function runYtdlpResilient(args, opts = {}) {
|
|
const hasClientArg = args.some((a) => String(a).includes('player_client='));
|
|
try {
|
|
return await runYtdlp(withCookies(args), opts);
|
|
} catch (err) {
|
|
if (hasClientArg || !BOT_CHECK_RE.test(err.message)) throw err;
|
|
if (opts.signal?.aborted) throw err;
|
|
let last = err;
|
|
for (const client of FALLBACK_CLIENTS) {
|
|
if (opts.signal?.aborted) throw last;
|
|
try {
|
|
const out = await runYtdlp(
|
|
withCookies(['--extractor-args', `youtube:player_client=${client}`, ...args]),
|
|
opts,
|
|
);
|
|
console.warn(`[ytplayer] bot check on default client, succeeded via player_client=${client}`);
|
|
return out;
|
|
} catch (e) {
|
|
last = e;
|
|
// A client that simply lacks the requested format is not a bot check;
|
|
// keep walking the list either way, but surface the last real error.
|
|
if (!BOT_CHECK_RE.test(e.message) && !/format is not available/i.test(e.message)) throw e;
|
|
}
|
|
}
|
|
throw last;
|
|
}
|
|
}
|
|
|
|
// Run ffmpeg the same way — async spawn so a multi-minute trim/concat never
|
|
// blocks Bun's event loop. Rejects on non-zero exit with ffmpeg's stderr tail.
|
|
function runFfmpeg(args) {
|
|
return new Promise((resolve, reject) => {
|
|
const child = spawn(FFMPEG, args, { stdio: ['ignore', 'ignore', 'pipe'] });
|
|
let err = '';
|
|
child.stderr.setEncoding('utf8');
|
|
// ffmpeg is extremely chatty on stderr; keep only the tail so an error
|
|
// message stays useful without buffering the whole progress log.
|
|
child.stderr.on('data', (d) => { err = (err + d).slice(-4000); });
|
|
child.on('error', (e) => reject(new Error('ffmpeg not found: ' + e.message)));
|
|
child.on('close', (code) => {
|
|
if (code !== 0) reject(new Error(err.trim() || 'ffmpeg exited with code ' + code));
|
|
else resolve();
|
|
});
|
|
});
|
|
}
|
|
|
|
// Parse the compact "s-e,s-e" keep-segment string (see frontend/video-edit.js)
|
|
// into an array of {start,end} second ranges. Skips malformed / non-increasing
|
|
// tokens; returns [] on empty or all-garbage input. Kept in lockstep with the
|
|
// frontend parseKeepParam so both ends agree on the wire format.
|
|
function parseKeepParam(str) {
|
|
if (typeof str !== 'string') return [];
|
|
const out = [];
|
|
for (const tok of str.split(',')) {
|
|
const t = tok.trim();
|
|
if (!t) continue;
|
|
const m = t.match(/^(\d+(?:\.\d+)?)-(\d+(?:\.\d+)?)$/);
|
|
if (!m) continue;
|
|
const a = parseFloat(m[1]);
|
|
const b = parseFloat(m[2]);
|
|
if (!isFinite(a) || !isFinite(b) || b <= a) continue;
|
|
out.push({ start: a, end: b });
|
|
}
|
|
return out;
|
|
}
|
|
|
|
// Build an ffmpeg filter_complex that trims `src` to the keep segments and
|
|
// concatenates them back into a single continuous stream. Re-encodes (the cut
|
|
// points rarely fall on keyframes, so stream-copy would glitch), producing one
|
|
// clean mp4. Returns the ffmpeg argv (input already appended by the caller).
|
|
function buildTrimArgs(keep) {
|
|
const parts = [];
|
|
keep.forEach((k, i) => {
|
|
parts.push(
|
|
`[0:v]trim=start=${k.start}:end=${k.end},setpts=PTS-STARTPTS[v${i}]`,
|
|
`[0:a]atrim=start=${k.start}:end=${k.end},asetpts=PTS-STARTPTS[a${i}]`,
|
|
);
|
|
});
|
|
const concatInputs = keep.map((_, i) => `[v${i}][a${i}]`).join('');
|
|
const filter = parts.join(';') + ';' +
|
|
`${concatInputs}concat=n=${keep.length}:v=1:a=1[outv][outa]`;
|
|
return [
|
|
'-filter_complex', filter,
|
|
'-map', '[outv]', '-map', '[outa]',
|
|
'-c:v', 'libx264', '-preset', 'veryfast', '-crf', '20',
|
|
'-c:a', 'aac', '-b:a', '160k',
|
|
'-movflags', '+faststart',
|
|
];
|
|
}
|
|
|
|
// Helpers to pick the right field from a yt-dlp JSON record
|
|
function pick(obj, ...keys) {
|
|
for (const k of keys) {
|
|
const v = obj[k];
|
|
if (v && typeof v === 'string' && v.trim()) return v.trim();
|
|
}
|
|
return '';
|
|
}
|
|
|
|
function pickChannel(obj) { return pick(obj, 'channel', 'uploader'); }
|
|
function pickChannelUrl(obj) { return pick(obj, 'channel_url', 'uploader_url'); }
|
|
function pickChannelId(obj) { return pick(obj, 'channel_id', 'uploader_id'); }
|
|
|
|
// Normalise a flat-playlist yt-dlp record into the slim UI card shape
|
|
function slimEntry(j) {
|
|
const id = pick(j, 'id');
|
|
if (!id) return null;
|
|
return {
|
|
id,
|
|
title: pick(j, 'title') || '(untitled)',
|
|
channel: pickChannel(j),
|
|
channelId: pickChannelId(j),
|
|
channelUrl: pickChannelUrl(j),
|
|
duration: typeof j.duration === 'number' ? j.duration : 0,
|
|
thumbnail: `https://i.ytimg.com/vi/${id}/mqdefault.jpg`,
|
|
};
|
|
}
|
|
|
|
// Parse multi-line JSON output from yt-dlp --dump-json --flat-playlist
|
|
function parseCards(output) {
|
|
const results = [];
|
|
for (const line of output.split('\n')) {
|
|
const t = line.trim();
|
|
if (!t) continue;
|
|
try {
|
|
const j = JSON.parse(t);
|
|
const card = slimEntry(j);
|
|
if (card) results.push(card);
|
|
} catch { /* skip malformed lines */ }
|
|
}
|
|
return results;
|
|
}
|
|
|
|
// Resolve a channel identifier to a /videos URL yt-dlp can fetch
|
|
function channelToUrl(c) {
|
|
let base = c.trim();
|
|
if (!base.startsWith('http')) {
|
|
base = base.startsWith('@') ? `https://www.youtube.com/${base}`
|
|
: base.startsWith('UC') ? `https://www.youtube.com/channel/${base}`
|
|
: `https://www.youtube.com/@${base}`;
|
|
}
|
|
base = base.replace(/\/$/, '');
|
|
return base.endsWith('/videos') ? base : base + '/videos';
|
|
}
|
|
|
|
// ============================================================================
|
|
// App setup
|
|
// ============================================================================
|
|
|
|
const app = new Hono();
|
|
app.use('*', logger());
|
|
|
|
// ============================================================================
|
|
// API routes
|
|
// ============================================================================
|
|
|
|
// GET /api/version
|
|
// Returns version string + a build tag that changes on every server restart/deploy.
|
|
// Clients poll this to detect when a new build is live and prompt a reload.
|
|
app.get('/api/version', (c) =>
|
|
c.json(
|
|
{ version: APP_VERSION, buildTag: BUILD_TAG, buildTime: BUILD_TIME },
|
|
200,
|
|
{ 'Cache-Control': 'no-store, no-cache, must-revalidate' }
|
|
)
|
|
);
|
|
|
|
// Same query repeats a lot (users re-searching, several devices, the
|
|
// landing-page chips) and yt-dlp's own search takes several seconds — cache
|
|
// the combined (library + YouTube) result for a short TTL. Same shape as
|
|
// resolveStreams()' streamCache below, just keyed by normalized query text
|
|
// instead of videoId. Only a full success is cached, so a transient yt-dlp
|
|
// failure still gets retried on the next request.
|
|
const SEARCH_CACHE_MAX = 100;
|
|
const SEARCH_CACHE_TTL_MS = 3 * 60_000;
|
|
const searchCache = new Map(); // lowercased q -> { results, expiresAt }
|
|
|
|
// GET /api/search?q=<query>
|
|
app.get('/api/search', async (c) => {
|
|
const q = (c.req.query('q') || '').trim();
|
|
if (!q) return c.json({ ok: false, error: 'empty query' }, 400);
|
|
|
|
const cacheKey = q.toLowerCase();
|
|
const cached = searchCache.get(cacheKey);
|
|
if (cached && Date.now() < cached.expiresAt) return c.json({ ok: true, results: cached.results });
|
|
|
|
// This server's own library first — and it still answers when YouTube
|
|
// (yt-dlp) is unreachable or rate-limited.
|
|
let mine = [];
|
|
try { mine = (await notesDb.listUploads({ q, limit: 20 })).map(uploads.card); } catch { /* library optional */ }
|
|
try {
|
|
let yt = [];
|
|
if (process.env.SEARCH_INNERTUBE !== '0') {
|
|
try {
|
|
yt = await innertube.search(q);
|
|
} catch (e) {
|
|
console.warn(`[search] innertube failed, using yt-dlp: ${e.message}`);
|
|
}
|
|
}
|
|
if (!yt.length) {
|
|
const out = await runYtdlpResilient([
|
|
`ytsearch${SEARCH_LIMIT}:${q}`,
|
|
'--dump-json', '--flat-playlist',
|
|
'--no-warnings', '--ignore-errors',
|
|
], { pooled: true });
|
|
yt = parseCards(out);
|
|
}
|
|
const results = [...mine, ...yt];
|
|
if (searchCache.size >= SEARCH_CACHE_MAX) searchCache.delete(searchCache.keys().next().value);
|
|
searchCache.set(cacheKey, { results, expiresAt: Date.now() + SEARCH_CACHE_TTL_MS });
|
|
return c.json({ ok: true, results });
|
|
} catch (err) {
|
|
if (mine.length) return c.json({ ok: true, results: mine, youtubeError: err.message });
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// GET /api/channel?c=<channel>
|
|
app.get('/api/channel', async (c) => {
|
|
const chan = (c.req.query('c') || '').trim();
|
|
if (!chan) return c.json({ ok: false, error: 'missing channel' }, 400);
|
|
|
|
try {
|
|
const url = channelToUrl(chan);
|
|
const out = await runYtdlpResilient([
|
|
url,
|
|
'--dump-json', '--flat-playlist',
|
|
'--no-warnings', '--ignore-errors',
|
|
'--playlist-end', String(CHANNEL_LIMIT),
|
|
], { pooled: true });
|
|
const results = parseCards(out);
|
|
// Extract channel name + URL from the first record
|
|
const first = results[0];
|
|
let channelName = '', channelUrl = '';
|
|
for (const line of out.split('\n')) {
|
|
const t = line.trim();
|
|
if (!t) continue;
|
|
try {
|
|
const j = JSON.parse(t);
|
|
channelName = channelName || pickChannel(j);
|
|
channelUrl = channelUrl || pickChannelUrl(j);
|
|
if (channelName && channelUrl) break;
|
|
} catch { /* skip */ }
|
|
}
|
|
return c.json({ ok: true, channel: channelName || (first?.channel || ''), channelUrl, results });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// ============================================================================
|
|
// Stream resolution + same-origin playback proxy
|
|
// ----------------------------------------------------------------------------
|
|
// googlevideo stream URLs are bound to the innertube client AND the IP that
|
|
// extracted them, and they expire (~6h via their `expire=` param). Handing
|
|
// them straight to the browser — which fetches from a DIFFERENT IP than this
|
|
// server — is what produced the intermittent 403s on playback (only the
|
|
// download path was ever fixed, in 70b5214). So playback now flows through
|
|
// /api/play: the browser hits our own origin, and THIS server fetches the
|
|
// googlevideo bytes (its IP matches the extractor) using yt-dlp's own
|
|
// http_headers, forwarding the browser's Range header so seeking still works.
|
|
// ============================================================================
|
|
|
|
// videoId -> { info, formats:[{formatId,height,vcodec,acodec,abr,ext,url,headers}], expiresAt }
|
|
const streamCache = new Map();
|
|
const STREAM_CACHE_MAX = 200;
|
|
// Max bytes per upstream googlevideo request (its adaptive streams are cut
|
|
// after ~10 MB; matches yt-dlp's http_chunk_size).
|
|
const PLAY_CHUNK = 10 * 1024 * 1024 - 1024;
|
|
|
|
// googlevideo URLs carry `expire=<unix-seconds>`; return that as an ms epoch.
|
|
function parseExpiry(url) {
|
|
const m = /[?&]expire=(\d+)/.exec(url || '');
|
|
return m ? Number(m[1]) * 1000 : 0;
|
|
}
|
|
|
|
// Resolve (and briefly cache) a video's playable formats via yt-dlp -J. The
|
|
// cache spares a fresh ~2-3s yt-dlp run on every Range request the media
|
|
// element fires; its TTL is bounded a minute inside the URLs' own expiry.
|
|
async function resolveStreamsUncached(videoId) {
|
|
const now = Date.now();
|
|
const cached = streamCache.get(videoId);
|
|
if (cached && now < cached.expiresAt) return cached;
|
|
|
|
const out = await runYtdlpResilient(['-J', '--no-warnings', `https://www.youtube.com/watch?v=${videoId}`], { pooled: true });
|
|
const info = JSON.parse(out);
|
|
const raw = Array.isArray(info.formats) ? info.formats : [];
|
|
const formats = [];
|
|
let soonest = Infinity;
|
|
for (const f of raw) {
|
|
if (!f.url) continue;
|
|
// Only keep direct https byte streams. Newer yt-dlp also surfaces HLS
|
|
// (m3u8/m3u8_native) and DASH-segment formats whose `url` is a manifest,
|
|
// not media bytes — proxying those hands the <video> element an m3u8 it
|
|
// can't play (Chrome) and defeats Range seeking (416). https formats are
|
|
// exactly the ones the dual-stream engine relied on before this change.
|
|
if (f.protocol && f.protocol !== 'https') continue;
|
|
const exp = parseExpiry(f.url);
|
|
if (exp) soonest = Math.min(soonest, exp);
|
|
formats.push({
|
|
formatId: String(f.format_id || ''),
|
|
height: f.height || 0,
|
|
vcodec: f.vcodec,
|
|
acodec: f.acodec,
|
|
abr: f.abr || 0,
|
|
ext: f.ext || '',
|
|
url: f.url,
|
|
headers: f.http_headers || {},
|
|
});
|
|
}
|
|
const ttlUntil = soonest === Infinity ? now + 30 * 60_000 : soonest - 60_000;
|
|
const expiresAt = Math.max(now + 60_000, Math.min(ttlUntil, now + 3 * 3600_000));
|
|
const entry = { info, formats, expiresAt };
|
|
if (streamCache.size >= STREAM_CACHE_MAX) streamCache.delete(streamCache.keys().next().value);
|
|
streamCache.set(videoId, entry);
|
|
return entry;
|
|
}
|
|
|
|
// One yt-dlp -J per video at a time: concurrent callers (warm-up + play,
|
|
// two devices, the media cache's getInfo) share the in-flight promise.
|
|
const inflightStreams = new Map(); // videoId -> Promise<entry>
|
|
function resolveStreams(videoId) {
|
|
const cached = streamCache.get(videoId);
|
|
if (cached && Date.now() < cached.expiresAt) return Promise.resolve(cached);
|
|
let p = inflightStreams.get(videoId);
|
|
if (!p) {
|
|
p = resolveStreamsUncached(videoId).finally(() => inflightStreams.delete(videoId));
|
|
inflightStreams.set(videoId, p);
|
|
}
|
|
return p;
|
|
}
|
|
|
|
const isVideoFmt = (f) => f.vcodec && f.vcodec !== 'none';
|
|
const isAudioFmt = (f) => f.acodec && f.acodec !== 'none';
|
|
|
|
// Highest-bitrate audio-only format (prefer mp4a/m4a).
|
|
function pickBestAudio(formats) {
|
|
let best = null, bestScore = -1;
|
|
for (const f of formats) {
|
|
if (isVideoFmt(f) || !isAudioFmt(f)) continue;
|
|
let score = f.abr || 0;
|
|
if (f.acodec && f.acodec.includes('mp4a')) score += 1000;
|
|
if (score > bestScore) { bestScore = score; best = f; }
|
|
}
|
|
return best;
|
|
}
|
|
|
|
// Pick the format /api/play should serve — by exact id first, then by the same
|
|
// height/progressive-first ordering /api/streams used to build the quality.
|
|
function pickFormat(formats, { formatId, wantAudio, wantHeight }) {
|
|
if (formatId) {
|
|
const byId = formats.find((f) => f.formatId === formatId);
|
|
if (byId) return byId;
|
|
}
|
|
if (wantAudio) return pickBestAudio(formats);
|
|
if (wantHeight) {
|
|
for (const wantProg of [false, true]) {
|
|
for (const f of formats) {
|
|
if (!isVideoFmt(f)) continue;
|
|
if (isAudioFmt(f) !== wantProg) continue;
|
|
if ((f.height || 0) === wantHeight) return f;
|
|
}
|
|
}
|
|
}
|
|
return null;
|
|
}
|
|
|
|
// GET /api/streams/warm?v=<id> — fire-and-forget: resolve streams into
|
|
// streamCache so the real /api/streams a moment later is instant. Bounded
|
|
// so a scrolling user can't queue dozens of yt-dlp processes.
|
|
const WARM_MAX = 2;
|
|
let warmActive = 0;
|
|
app.get('/api/streams/warm', async (c) => {
|
|
const id = (c.req.query('v') || '').trim();
|
|
if (!/^[A-Za-z0-9_-]{11}$/.test(id)) return c.body(null, 204);
|
|
if (warmActive >= WARM_MAX) return c.body(null, 204);
|
|
try { if (await media.getReady(id)) return c.body(null, 204); } catch { /* fall through */ }
|
|
warmActive++;
|
|
resolveStreams(id).catch(() => {}).finally(() => { warmActive--; });
|
|
return c.body(null, 204);
|
|
});
|
|
|
|
// One view per client per video per 30 min (a play and its save count once).
|
|
const viewSeen = new Map(); // `${who}|${id}` -> ms
|
|
function countView(c, videoId) {
|
|
const who = (c.req.header('x-forwarded-for') || '').split(',')[0].trim() || c.req.query('fp') || 'local';
|
|
const key = who + '|' + videoId;
|
|
const now = Date.now();
|
|
if (now - (viewSeen.get(key) || 0) < 30 * 60_000) return;
|
|
viewSeen.set(key, now);
|
|
if (viewSeen.size > 20000) viewSeen.clear();
|
|
p2pDb.addView(videoId, now).catch(() => {});
|
|
}
|
|
|
|
// GET /api/streams?v=<videoId> — meta + proxied audio/quality URLs.
|
|
app.get('/api/streams', async (c) => {
|
|
const videoId = (c.req.query('v') || '').replace(/[/\\:?<>|*"]/g, '').trim();
|
|
if (!videoId) return c.json({ ok: false, error: 'missing videoId' }, 400);
|
|
countView(c, videoId);
|
|
|
|
// An upload from the server's own library — no yt-dlp, no media cache.
|
|
if (isUpload(videoId)) {
|
|
const u = await notesDb.getUpload(videoId);
|
|
if (!u) return c.json({ ok: false, error: 'upload not found' }, 404);
|
|
return c.json({ ok: true, data: uploads.streamsPayload(u) });
|
|
}
|
|
|
|
// Server already holds a validated copy → answer from the DB alone, no
|
|
// yt-dlp round trip. ?nocache=1 (the client's fallback when a cached copy
|
|
// won't play on its device) forces the YouTube path below.
|
|
if (c.req.query('nocache') !== '1') {
|
|
try {
|
|
const row = await media.getReady(videoId);
|
|
if (servableTo(row, c.req.query('hevc') === '1')) {
|
|
media.touch(videoId);
|
|
return c.json({ ok: true, data: cachedStreamsPayload(videoId, row) });
|
|
}
|
|
} catch (err) {
|
|
console.warn(`[media] cache lookup failed for ${videoId}:`, err.message);
|
|
}
|
|
} else {
|
|
// A client couldn't play our copy — re-check it server-side (background).
|
|
media.verify(videoId).catch(() => {});
|
|
}
|
|
|
|
try {
|
|
const { info, formats } = await resolveStreams(videoId);
|
|
|
|
// Proxied audio URL (same-origin, re-resolved fresh on each play).
|
|
const bestAudio = pickBestAudio(formats);
|
|
const audioUrl = bestAudio
|
|
? `/api/play?v=${videoId}&audio=1&f=${encodeURIComponent(bestAudio.formatId)}`
|
|
: null;
|
|
|
|
// Quality list — one entry per height. ADAPTIVE (video-only, paired with
|
|
// the separate audioUrl for dual-stream playback) is preferred over
|
|
// progressive single-file: adaptive https streams proxy reliably, while
|
|
// progressive muxed formats (itag 18 and friends) are increasingly
|
|
// SABR-gated and flaky. Progressive is kept only when no adaptive stream
|
|
// exists at that height.
|
|
const qualities = [];
|
|
const seen = new Set();
|
|
for (const wantProg of [false, true]) {
|
|
for (const f of formats) {
|
|
if (!isVideoFmt(f)) continue;
|
|
if (isAudioFmt(f) !== wantProg) continue;
|
|
const h = f.height || 0;
|
|
if (h <= 0 || seen.has(h)) continue;
|
|
seen.add(h);
|
|
qualities.push({
|
|
label: h + 'p',
|
|
height: h,
|
|
hasAudio: wantProg,
|
|
url: `/api/play?v=${videoId}&h=${h}&f=${encodeURIComponent(f.formatId)}`,
|
|
ext: f.ext || '',
|
|
});
|
|
}
|
|
}
|
|
qualities.sort((a, b) => b.height - a.height);
|
|
|
|
// Nothing proxyable (live streams only expose HLS manifests) — say so
|
|
// instead of handing the player an empty list it silently stalls on.
|
|
if (!qualities.length && !audioUrl) {
|
|
const live = info.is_live || info.live_status === 'is_live';
|
|
return c.json({ ok: false, error: live ? 'Live streams are not supported' : 'No playable formats for this video' }, 422);
|
|
}
|
|
|
|
// Played but not cached yet: fetch a copy in the background so the next
|
|
// play is served from disk. Runs server-side; the client never waits.
|
|
media.ensureCached(videoId, { priority: LOW, auto: true }).catch(() => {});
|
|
|
|
return c.json({
|
|
ok: true,
|
|
data: {
|
|
meta: {
|
|
id: videoId,
|
|
title: pick(info, 'title') || '(untitled)',
|
|
channel: pickChannel(info),
|
|
channelId: pickChannelId(info),
|
|
channelUrl: pickChannelUrl(info),
|
|
duration: typeof info.duration === 'number' ? info.duration : 0,
|
|
thumbnail: `https://i.ytimg.com/vi/${videoId}/hqdefault.jpg`,
|
|
},
|
|
audioUrl,
|
|
qualities,
|
|
},
|
|
});
|
|
} catch (err) {
|
|
// The source failed. If a device holds a verified copy, ask it to send one
|
|
// to the server so the video comes back (P2P flow 8).
|
|
let restoring = false;
|
|
try { restoring = await p2pRehydrate(videoId); } catch { /* best effort */ }
|
|
if (restoring) {
|
|
return c.json({
|
|
ok: false, restoring: true,
|
|
error: 'This video is unavailable at the source — a device that has it is sending a copy to the server. Try again in a minute.',
|
|
}, 503);
|
|
}
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// Stream a single format straight out of yt-dlp's stdout. Some formats
|
|
// (notably progressive itag 18, and any SABR-gated stream) 403 on a plain GET
|
|
// even from the extractor IP — the bytes are only reachable through yt-dlp's
|
|
// full protocol. This can't honour a byte Range (yt-dlp writes start-to-end to
|
|
// the pipe), so it always answers 200; the media element still plays, it just
|
|
// can't seek past what it has buffered. Used only as the /api/play fallback.
|
|
//
|
|
// It PEEKS the first chunk before committing to a 200: if yt-dlp errors or
|
|
// exits without emitting any bytes (stale binary, ffmpeg segfault on merge,
|
|
// dead itag), it resolves to null so the caller can return a real error and
|
|
// the frontend advances to the next candidate instead of playing an empty 200.
|
|
function ytdlpPipeResponse(videoId, formatArgs, contentType) {
|
|
return new Promise((resolve) => {
|
|
const child = spawn(YTDLP, [
|
|
`https://www.youtube.com/watch?v=${videoId}`,
|
|
'--no-warnings', '--no-playlist',
|
|
...formatArgs,
|
|
'-o', '-',
|
|
], { stdio: ['ignore', 'pipe', 'ignore'] });
|
|
|
|
let settled = false;
|
|
const finish = (val) => { if (!settled) { settled = true; resolve(val); } };
|
|
|
|
child.stdout.once('data', (first) => {
|
|
// Re-emit the peeked chunk so the response stream is byte-complete.
|
|
child.stdout.unshift(first);
|
|
finish(new Response(Readable.toWeb(child.stdout), {
|
|
status: 200,
|
|
headers: {
|
|
'Content-Type': contentType,
|
|
'Accept-Ranges': 'none',
|
|
'Cache-Control': 'no-store',
|
|
},
|
|
}));
|
|
});
|
|
child.on('error', () => { try { child.kill(); } catch {} finish(null); });
|
|
child.on('close', () => finish(null)); // closed before any data → failure
|
|
});
|
|
}
|
|
|
|
// GET /api/play?v=<id>&f=<formatId>[&h=<height>][&audio=1]
|
|
// Same-origin streaming proxy for playback. Fetches the googlevideo bytes from
|
|
// THIS server (whose IP matches the extractor) with yt-dlp's own http_headers,
|
|
// forwarding the browser's Range header so seeking works. This is what makes
|
|
// the previously-403ing direct URLs play reliably. If the direct URL still
|
|
// 403s (SABR/itag-18 formats reachable only via yt-dlp's protocol), it falls
|
|
// back to piping yt-dlp itself.
|
|
app.get('/api/play', async (c) => {
|
|
const videoId = (c.req.query('v') || '').replace(/[/\\:?<>|*"]/g, '').trim();
|
|
if (!videoId) return c.json({ ok: false, error: 'missing videoId' }, 400);
|
|
const sel = {
|
|
formatId: (c.req.query('f') || '').trim(),
|
|
wantAudio: c.req.query('audio') === '1',
|
|
wantHeight: Number(c.req.query('h') || 0),
|
|
};
|
|
const range = c.req.header('range');
|
|
|
|
// Fetch the chosen format's bytes; on a stale-URL 403/410 drop the cache and
|
|
// re-resolve once before giving up.
|
|
async function upstreamFetch(forceFresh) {
|
|
if (forceFresh) streamCache.delete(videoId);
|
|
const { formats } = await resolveStreams(videoId);
|
|
const fmt = pickFormat(formats, sel);
|
|
if (!fmt) return { error: 'no matching format' };
|
|
const headers = { ...fmt.headers };
|
|
// Always ask upstream for a bounded slice (see PLAY_CHUNK below); the
|
|
// splicer extends it to the client's full requested range.
|
|
const rq = /bytes=(\d+)-(\d*)/.exec(range || '');
|
|
const from = rq ? Number(rq[1]) : 0;
|
|
const to = rq && rq[2] !== '' ? Math.min(Number(rq[2]), from + PLAY_CHUNK - 1) : from + PLAY_CHUNK - 1;
|
|
headers['Range'] = range && !rq ? range : `bytes=${from}-${to}`;
|
|
return { fmt, res: await fetch(fmt.url, { headers, redirect: 'follow' }) };
|
|
}
|
|
|
|
try {
|
|
let attempt = await upstreamFetch(false);
|
|
if (attempt.error) return c.json({ ok: false, error: attempt.error }, 404);
|
|
if (attempt.res.status === 403 || attempt.res.status === 410) {
|
|
const retry = await upstreamFetch(true);
|
|
if (!retry.error && retry.res) attempt = retry;
|
|
}
|
|
|
|
// Direct URL is genuinely ungettable (SABR / itag-18) — let yt-dlp fetch it.
|
|
if (attempt.res.status === 403 || attempt.res.status === 410) {
|
|
const fmtArgs = sel.formatId ? ['-f', sel.formatId]
|
|
: sel.wantAudio ? ['-f', 'bestaudio[ext=m4a]/bestaudio']
|
|
: sel.wantHeight ? ['-f', `best[height=${sel.wantHeight}]/bv*[height=${sel.wantHeight}]`]
|
|
: ['-f', 'best'];
|
|
const piped = await ytdlpPipeResponse(videoId, fmtArgs, sel.wantAudio ? 'audio/mp4' : 'video/mp4');
|
|
if (piped) return piped;
|
|
return c.json({ ok: false, error: 'stream unavailable (403)' }, 502);
|
|
}
|
|
|
|
const up = attempt.res;
|
|
const headers = new Headers();
|
|
headers.set('Content-Type', up.headers.get('content-type') || (sel.wantAudio ? 'audio/mp4' : 'video/mp4'));
|
|
headers.set('Accept-Ranges', 'bytes');
|
|
headers.set('Cache-Control', 'no-store');
|
|
|
|
// googlevideo silently closes adaptive-stream connections after ~10 MB
|
|
// (yt-dlp's own `http_chunk_size` exists for exactly this), so a single
|
|
// forwarded `bytes=0-` truncated at ~30 s of 720p and the <video> element
|
|
// stalled with no error. Re-issue upstream fetches in ≤10 MB Range slices
|
|
// and splice them into ONE response body that is byte-complete for the
|
|
// client's requested range.
|
|
const cr = /bytes (\d+)-(\d+)\/(\d+|\*)/.exec(up.headers.get('content-range') || '');
|
|
const total = cr && cr[3] !== '*' ? Number(cr[3])
|
|
: up.status === 200 ? Number(up.headers.get('content-length') || 0) : 0;
|
|
if (!total || up.status !== 206 && up.status !== 200) {
|
|
for (const k of ['content-length', 'content-range']) {
|
|
const v = up.headers.get(k);
|
|
if (v) headers.set(k, v);
|
|
}
|
|
return new Response(up.body, { status: up.status, headers });
|
|
}
|
|
|
|
const rq = /bytes=(\d*)-(\d*)/.exec(range || '');
|
|
let start = cr ? Number(cr[1]) : 0;
|
|
let end = total - 1;
|
|
if (rq && rq[2] !== '') end = Math.min(total - 1, Number(rq[2]));
|
|
if (rq && rq[1] === '' && rq[2] !== '') { start = Math.max(0, total - Number(rq[2])); end = total - 1; }
|
|
const length = end - start + 1;
|
|
|
|
const fmt = attempt.fmt;
|
|
// Client-disconnect handling: the active upstream body is LOCKED by its
|
|
// reader, so `body.cancel()` on it throws ("Cannot cancel a locked
|
|
// ReadableStream") — and an exception thrown from cancel() took the whole
|
|
// Bun process down on every seek / quality switch / aborted fetch. Cancel
|
|
// through the reader instead, and stop the slice loop via a flag.
|
|
let aborted = false;
|
|
let reader = null;
|
|
const body = new ReadableStream({
|
|
async start(ctrl) {
|
|
// Upstream body is a slice already; the slice boundaries below are
|
|
// relative to `start`, and the first slice reuses the in-flight fetch.
|
|
let pos = start;
|
|
let res = up;
|
|
try {
|
|
while (pos <= end && !aborted) {
|
|
const sliceEnd = Math.min(end, pos + PLAY_CHUNK - 1);
|
|
if (res === null) {
|
|
res = await fetch(fmt.url, { headers: { ...fmt.headers, Range: `bytes=${pos}-${sliceEnd}` }, redirect: 'follow' });
|
|
if (res.status !== 206 && res.status !== 200) throw new Error('upstream ' + res.status);
|
|
}
|
|
reader = res.body.getReader();
|
|
let got = 0;
|
|
while (got < sliceEnd - pos + 1 && !aborted) {
|
|
const { value, done } = await reader.read();
|
|
if (done) break;
|
|
const room = sliceEnd - pos + 1 - got;
|
|
const chunk = value.length > room ? value.subarray(0, room) : value;
|
|
ctrl.enqueue(chunk);
|
|
got += chunk.length;
|
|
}
|
|
try { await reader.cancel(); } catch {}
|
|
reader = null;
|
|
if (aborted) break;
|
|
if (got === 0) throw new Error('upstream returned no bytes');
|
|
pos += got;
|
|
res = null;
|
|
}
|
|
if (!aborted) ctrl.close();
|
|
} catch (err) {
|
|
if (!aborted) { try { ctrl.error(err); } catch {} }
|
|
}
|
|
},
|
|
cancel() {
|
|
aborted = true;
|
|
const r = reader;
|
|
if (r) r.cancel().catch(() => {});
|
|
},
|
|
});
|
|
|
|
headers.set('Content-Length', String(length));
|
|
if (range) {
|
|
headers.set('Content-Range', `bytes ${start}-${end}/${total}`);
|
|
return new Response(body, { status: 206, headers });
|
|
}
|
|
return new Response(body, { status: 200, headers });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: 'stream proxy failed: ' + err.message }, 502);
|
|
}
|
|
});
|
|
|
|
// Download via yt-dlp into a self-cleaning temp file, then stream it.
|
|
// yt-dlp MUST perform the HTTP fetch itself: googlevideo stream URLs are
|
|
// bound to the innertube client that extracted them, so resolving the URL
|
|
// 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
|
|
// didn't match the extraction client.
|
|
// 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 || 3 * 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)`);
|
|
}
|
|
}
|
|
|
|
// Format selectors for saves. `[ext=mp4]` alone is NOT enough: YouTube ships
|
|
// AV1 in mp4 containers too, and the bot-check fallback clients (web_embedded)
|
|
// rank the AV1 ladder as "bestvideo" — an AV1 save plays on almost nothing
|
|
// (iOS Safari has no AV1 decoder), so cached copies errored and fell back to
|
|
// streaming. Demand h264 (`vcodec^=avc1`) + m4a/aac first; anything exotic is
|
|
// a last resort for videos with no avc ladder at all.
|
|
const SAVE_MUX_FORMAT =
|
|
'bv*[height<=720][vcodec^=avc1]+ba[ext=m4a]/bv*[height<=720][vcodec^=avc1]+ba/'
|
|
+ 'b[ext=mp4][vcodec^=avc1]/bv*[height<=720]+ba/b[ext=mp4]/b';
|
|
const SAVE_DEFAULT_FORMAT =
|
|
'best[ext=mp4][vcodec^=avc1][acodec!=none]/'
|
|
+ 'bv*[height<=720][vcodec^=avc1]+ba[ext=m4a]/bv*[height<=720][vcodec^=avc1]+ba/'
|
|
+ 'best[ext=mp4][acodec!=none]/bv*[height<=720]+ba/b';
|
|
|
|
// 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 tmp = `${tmpdir()}/${tmpBase}.mp4`;
|
|
const ctl = new AbortController();
|
|
entry = { waiters: 0, tmpBase, ctl };
|
|
entry.promise = withSaveSlot(() => {
|
|
if (ctl.signal.aborted) throw new Error('save cancelled');
|
|
return runYtdlpResilient([
|
|
`https://www.youtube.com/watch?v=${videoId}`,
|
|
'--no-warnings', '--no-playlist',
|
|
...formatArgs,
|
|
'--limit-rate', DOWNLOAD_RATE,
|
|
'-o', tmp,
|
|
], { signal: ctl.signal });
|
|
}).then(() => ({ tmp, size: statSync(tmp).size }));
|
|
inflightDownloads.set(key, entry);
|
|
}
|
|
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
|
|
// if the client disconnects mid-transfer.
|
|
fd = openSync(res.tmp, 'r');
|
|
} finally {
|
|
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
|
|
// intermediates (left when ffmpeg is missing — yt-dlp then downloads
|
|
// the streams separately, exits 0 without merging, and statSync above
|
|
// throws on the absent merged file).
|
|
for (const name of readdirSync(tmpdir())) {
|
|
if (name.startsWith(entry.tmpBase)) {
|
|
try { unlinkSync(`${tmpdir()}/${name}`); } catch { /* already gone */ }
|
|
}
|
|
}
|
|
}
|
|
}
|
|
const stream = createReadStream('', { fd });
|
|
|
|
if (fp) recordVideoAccess(fp, { id: videoId }).catch(() => {});
|
|
|
|
return new Response(Readable.toWeb(stream), {
|
|
status: 200,
|
|
headers: {
|
|
'Content-Type': 'video/mp4',
|
|
'Content-Length': String(size),
|
|
'Content-Disposition': `attachment; filename="${videoId}.mp4"`,
|
|
'Cache-Control': 'no-store',
|
|
'Access-Control-Allow-Origin': '*',
|
|
},
|
|
});
|
|
}
|
|
|
|
// "Edit & download": fetch the source with yt-dlp (muxed up to 720p, same as
|
|
// the mux path), then run ffmpeg to KEEP only the requested segments and
|
|
// 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
|
|
// under the caller-chosen custom id. Every temp file is swept afterwards.
|
|
async function ytdlpEditedDownloadResponse(videoId, fp, keep, signal, cachedSrc = null) {
|
|
if (!cachedSrc) await assertNotLive(videoId);
|
|
const tmpBase = `ytp-edit-${videoId}-${Date.now()}`;
|
|
const srcTmp = cachedSrc || `${tmpdir()}/${tmpBase}.src.mp4`;
|
|
const outTmp = `${tmpdir()}/${tmpBase}.out.mp4`;
|
|
let size, fd;
|
|
try {
|
|
// 1) Grab the full source (video+audio merged) so ffmpeg has both streams
|
|
// — unless the server cache already holds it.
|
|
if (!cachedSrc) await withSaveSlot(() => runYtdlpResilient([
|
|
`https://www.youtube.com/watch?v=${videoId}`,
|
|
'--no-warnings', '--no-playlist',
|
|
'-f', SAVE_MUX_FORMAT,
|
|
'--merge-output-format', 'mp4',
|
|
'--limit-rate', DOWNLOAD_RATE,
|
|
'-o', srcTmp,
|
|
], { signal }));
|
|
// 2) Trim + concat the keep segments into the final custom video.
|
|
await runFfmpeg([
|
|
'-y', '-hide_banner', '-loglevel', 'error',
|
|
'-i', srcTmp,
|
|
...buildTrimArgs(keep),
|
|
outTmp,
|
|
]);
|
|
size = statSync(outTmp).size;
|
|
fd = openSync(outTmp, 'r');
|
|
} finally {
|
|
for (const name of readdirSync(tmpdir())) {
|
|
if (name.startsWith(tmpBase)) {
|
|
try { unlinkSync(`${tmpdir()}/${name}`); } catch { /* already gone */ }
|
|
}
|
|
}
|
|
}
|
|
const stream = createReadStream('', { fd });
|
|
|
|
if (fp) recordVideoAccess(fp, { id: videoId }).catch(() => {});
|
|
|
|
return new Response(Readable.toWeb(stream), {
|
|
status: 200,
|
|
headers: {
|
|
'Content-Type': 'video/mp4',
|
|
'Content-Length': String(size),
|
|
'Content-Disposition': `attachment; filename="${videoId}-edited.mp4"`,
|
|
'Cache-Control': 'no-store',
|
|
'Access-Control-Allow-Origin': '*',
|
|
},
|
|
});
|
|
}
|
|
|
|
// ============================================================================
|
|
// Server-side media cache (see media-cache.js). One validated ≤720p H.264 +
|
|
// AAC copy per played/saved video under $MEDIA_DIR, fetched by server-owned
|
|
// jobs that survive the client going away. Budget-bounded LRU; never expires
|
|
// by time. Downloads still go through withSaveSlot so they stay serialized
|
|
// with any legacy fallback save (the homelab downlink is the constraint).
|
|
// ============================================================================
|
|
const envNum = (k, d) => (process.env[k] !== undefined && process.env[k] !== '' ? Number(process.env[k]) : d);
|
|
const MEDIA_DIR = process.env.MEDIA_DIR || './data/media';
|
|
mkdirSync(MEDIA_DIR, { recursive: true });
|
|
|
|
const media = createMediaCache({
|
|
dir: MEDIA_DIR,
|
|
db: {
|
|
getMedia, upsertMedia, deleteMedia, listMedia, listMediaLru, touchMedia, mediaStats,
|
|
// Retention (docs/p2p-architecture.md flow 9): cold copies go before popular ones.
|
|
listMediaEvictionOrder: () => p2pDb.listMediaEvictionOrder({
|
|
now: Date.now(), keepMinViews: P2P.keepMinViews, keepDays: P2P.keepDays, keepRecentDays: P2P.keepRecentDays,
|
|
}),
|
|
},
|
|
getInfo: async (videoId) => (await resolveStreams(videoId)).info,
|
|
download: async (videoId, out) => {
|
|
await withSaveSlot(() => runYtdlpResilient([
|
|
`https://www.youtube.com/watch?v=${videoId}`,
|
|
'--no-warnings', '--no-playlist',
|
|
'-f', SAVE_MUX_FORMAT,
|
|
'--merge-output-format', 'mp4',
|
|
'--limit-rate', DOWNLOAD_RATE,
|
|
'-o', out,
|
|
]));
|
|
return out;
|
|
},
|
|
ffmpeg: FFMPEG,
|
|
ffprobe: process.env.FFPROBE_PATH || 'ffprobe',
|
|
maxBytes: envNum('MEDIA_CACHE_MAX_BYTES', 10 * 1024 ** 3),
|
|
minFreeBytes: envNum('MEDIA_MIN_FREE_BYTES', 5 * 1024 ** 3),
|
|
autoMaxSeconds: envNum('MEDIA_AUTO_MAX_SECONDS', 3 * 3600),
|
|
saveMaxSeconds: MAX_SAVE_SECONDS,
|
|
transcode: {
|
|
enabled: process.env.MEDIA_TRANSCODE !== '0',
|
|
codec: process.env.MEDIA_CODEC || 'hevc', // hevc | h264
|
|
crf: envNum('MEDIA_CRF', 28),
|
|
preset: process.env.MEDIA_PRESET || 'medium',
|
|
threads: envNum('MEDIA_THREADS', 2),
|
|
maxSeconds: envNum('MEDIA_OPT_MAX_SECONDS', 3600),
|
|
},
|
|
// P2P (docs/p2p-architecture.md): every validated copy's server-computed
|
|
// hash becomes verified content, after the optional malware scan.
|
|
onReady: (info) => admitFile(
|
|
{ ...info, cid: info.sha256, videoId: info.id, origin: 'server' },
|
|
{ cfg: P2P, upsertContent: p2pDb.upsertContent },
|
|
),
|
|
});
|
|
|
|
// A cached copy the compression lane turned into HEVC is only handed to
|
|
// clients that said they can decode it (?hevc=1 — the page checks
|
|
// canPlayType for hvc1). Everyone else keeps the pre-cache behaviour.
|
|
function servableTo(row, hevcOk) {
|
|
return !!row && (row.vcodec !== 'hevc' || hevcOk);
|
|
}
|
|
|
|
// /api/streams response for a server-cached video: one progressive quality
|
|
// (the mp4 carries its own audio) plus the m4a sidecar for audio-only mode.
|
|
// Same shape as the YouTube path; `serverCached` and `codec` are additive.
|
|
function cachedStreamsPayload(videoId, row) {
|
|
let meta = {};
|
|
try { meta = JSON.parse(row.meta || '{}'); } catch { /* corrupt — use defaults */ }
|
|
const url = `/api/media/${videoId}?g=${row.gen}`;
|
|
// data.cid is additive and web-only (like serverCached) — the Tauri bridge ignores it.
|
|
return {
|
|
meta: {
|
|
id: videoId,
|
|
title: meta.title || '(untitled)',
|
|
channel: meta.channel || '',
|
|
channelId: meta.channelId || '',
|
|
channelUrl: meta.channelUrl || '',
|
|
duration: meta.duration || Math.round(row.duration || 0),
|
|
thumbnail: `https://i.ytimg.com/vi/${videoId}/hqdefault.jpg`,
|
|
},
|
|
audioUrl: `${url}&a=1`,
|
|
qualities: [{
|
|
label: (row.height || 720) + 'p',
|
|
height: row.height || 720,
|
|
hasAudio: true,
|
|
url,
|
|
ext: 'mp4',
|
|
codec: row.vcodec || 'h264',
|
|
}],
|
|
serverCached: true,
|
|
cid: row.sha256 || null,
|
|
};
|
|
}
|
|
|
|
// Serve a file with byte-Range support (the <video> element seeks with it).
|
|
function rangeFileResponse(c, path, contentType, cacheControl) {
|
|
const file = Bun.file(path);
|
|
const total = file.size;
|
|
const headers = {
|
|
'Content-Type': contentType,
|
|
'Accept-Ranges': 'bytes',
|
|
'Cache-Control': cacheControl,
|
|
};
|
|
const range = c.req.header('range');
|
|
if (!range) return new Response(file, { status: 200, headers });
|
|
const m = /^bytes=(\d*)-(\d*)$/.exec(range.trim());
|
|
let start, end;
|
|
if (m && m[1] !== '') {
|
|
start = Number(m[1]);
|
|
end = m[2] !== '' ? Math.min(Number(m[2]), total - 1) : total - 1;
|
|
} else if (m && m[2] !== '') {
|
|
start = Math.max(0, total - Number(m[2]));
|
|
end = total - 1;
|
|
}
|
|
if (start === undefined || start > end || start >= total) {
|
|
return new Response(null, { status: 416, headers: { ...headers, 'Content-Range': `bytes */${total}` } });
|
|
}
|
|
return new Response(file.slice(start, end + 1), {
|
|
status: 206,
|
|
headers: { ...headers, 'Content-Range': `bytes ${start}-${end}/${total}` },
|
|
});
|
|
}
|
|
|
|
// GET /api/media/stats — registered before /api/media/:id so it isn't an id.
|
|
app.get('/api/media/stats', async (c) => {
|
|
try {
|
|
return c.json({ ok: true, ...(await media.stats()) }, 200, { 'Cache-Control': 'no-store' });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// GET /api/media/:id?g=<gen>[&a=1] — the cached mp4 (or its m4a sidecar).
|
|
// A URL with ?g= is immutable: a replaced copy gets a new gen, and an old gen
|
|
// is served only while its file still exists (never different bytes).
|
|
app.get('/api/media/:id', async (c) => {
|
|
const id = c.req.param('id');
|
|
const g = c.req.query('g');
|
|
const audio = c.req.query('a') === '1';
|
|
const path = await media.filePath(id, g, audio ? 'm4a' : 'mp4');
|
|
if (!path) return c.json({ ok: false, error: 'not cached' }, 404);
|
|
media.touch(id);
|
|
return rangeFileResponse(c, path, audio ? 'audio/mp4' : 'video/mp4',
|
|
g ? 'public, max-age=31536000, immutable' : 'no-cache');
|
|
});
|
|
|
|
// GET /api/export/:id?name=<file name>[&hevc=1] — the server's copy as a file
|
|
// download (Content-Disposition: attachment), Range-aware so a big file can
|
|
// resume. Serves a cached YouTube copy or one of the server's own uploads, so a
|
|
// device can put it in its gallery / Files app without going through OPFS.
|
|
app.get('/api/export/:id', async (c) => {
|
|
const id = c.req.param('id');
|
|
let path = null, mime = 'video/mp4', ext = 'mp4';
|
|
if (isUpload(id)) {
|
|
const u = await notesDb.getUpload(id);
|
|
if (u) { path = uploads.filePath(u); mime = u.mime; ext = u.ext; }
|
|
} else {
|
|
const row = await notesDb.getMedia(id);
|
|
// An HEVC copy only goes to a client that said it can play one.
|
|
if (row && row.status === 'ready' && (row.vcodec !== 'hevc' || c.req.query('hevc') === '1')) path = await media.filePath(id, null, 'mp4');
|
|
}
|
|
if (!path) return c.json({ ok: false, error: 'not cached on the server' }, 404);
|
|
media.touch?.(id);
|
|
const base = String(c.req.query('name') || id).replace(/\.[a-z0-9]{1,5}$/i, '')
|
|
.replace(/[\u0000-\u001f\\/:*?"<>|]+/g, ' ').replace(/\s+/g, ' ').trim().slice(0, 120) || id;
|
|
const res = rangeFileResponse(c, path, mime, 'no-store');
|
|
res.headers.set('Content-Disposition', `attachment; filename="${base.replace(/[^\x20-\x7e]/g, '_').replace(/"/g, '')}.${ext}"; filename*=UTF-8''${encodeURIComponent(base)}.${ext}`);
|
|
return res;
|
|
});
|
|
|
|
// GET /api/media/:id/status
|
|
// Where the bytes for an id live: the validated YouTube copy in the media
|
|
// cache, or one of the server's own uploads.
|
|
async function audioSourcePath(id) {
|
|
if (isUpload(id)) {
|
|
const u = await notesDb.getUpload(id);
|
|
return u ? uploads.filePath(u) : null;
|
|
}
|
|
return (await media.filePath(id, null, 'm4a')) || (await media.filePath(id, null, 'mp4'));
|
|
}
|
|
async function videoSourcePath(id) {
|
|
if (isUpload(id)) {
|
|
const u = await notesDb.getUpload(id);
|
|
return u && u.kind === 'video' ? uploads.filePath(u) : null;
|
|
}
|
|
return media.filePath(id, null, 'mp4');
|
|
}
|
|
|
|
// GET /api/media/:id/peaks — loudness envelope of a server-cached copy for
|
|
// the waveform seek bar: PEAKS_N RMS buckets scaled 0..100. Computed once per
|
|
// file with ffmpeg (mono, 2 kHz is plenty for an envelope) and kept in memory;
|
|
// the file path carries the gen, so a replaced copy is recomputed.
|
|
const PEAKS_N = 400;
|
|
const PEAKS_RATE = 2000;
|
|
const peaksCache = new Map();
|
|
let peaksRunning = 0;
|
|
function computePeaks(path) {
|
|
return new Promise((resolve, reject) => {
|
|
const child = spawn(FFMPEG, ['-v', 'error', '-i', path, '-vn', '-ac', '1', '-ar', String(PEAKS_RATE), '-f', 's16le', '-'],
|
|
{ stdio: ['ignore', 'pipe', 'pipe'] });
|
|
const chunks = [];
|
|
let err = '';
|
|
child.stdout.on('data', (d) => chunks.push(d));
|
|
child.stderr.on('data', (d) => { err = (err + d).slice(-2000); });
|
|
child.on('error', reject);
|
|
child.on('close', (code) => {
|
|
if (code !== 0) return reject(new Error(err.trim() || `ffmpeg exited ${code}`));
|
|
const buf = Buffer.concat(chunks);
|
|
const n = Math.floor(buf.length / 2);
|
|
if (!n) return reject(new Error('no audio'));
|
|
const per = Math.max(1, Math.ceil(n / PEAKS_N));
|
|
const rms = [];
|
|
for (let b = 0; b * per < n; b++) {
|
|
let sum = 0, cnt = 0;
|
|
for (let i = b * per; i < Math.min(n, (b + 1) * per); i++) { const v = buf.readInt16LE(i * 2); sum += v * v; cnt++; }
|
|
rms.push(Math.sqrt(sum / Math.max(1, cnt)));
|
|
}
|
|
const max = Math.max(...rms) || 1;
|
|
// ^0.7 lifts quiet passages so verses stay visible next to a loud bridge.
|
|
resolve({ duration: n / PEAKS_RATE, peaks: rms.map((v) => Math.round(Math.pow(v / max, 0.7) * 100)) });
|
|
});
|
|
});
|
|
}
|
|
|
|
app.get('/api/media/:id/peaks', async (c) => {
|
|
const id = c.req.param('id');
|
|
if (!/^([A-Za-z0-9_-]{11}|upl_[a-f0-9]{12})$/.test(id)) return c.json({ ok: false, error: 'invalid id' }, 400);
|
|
const path = await audioSourcePath(id);
|
|
if (!path) return c.json({ ok: false, error: 'not cached on the server' }, 404);
|
|
let hit = peaksCache.get(path);
|
|
if (!hit) {
|
|
if (peaksRunning >= 2) return c.json({ ok: false, error: 'busy — retry shortly' }, 503);
|
|
peaksRunning++;
|
|
try { hit = await computePeaks(path); } catch (err) { return c.json({ ok: false, error: err.message }, 500); }
|
|
finally { peaksRunning--; }
|
|
peaksCache.set(path, hit);
|
|
if (peaksCache.size > 300) peaksCache.delete(peaksCache.keys().next().value);
|
|
}
|
|
return c.json({ ok: true, ...hit }, 200, { 'Cache-Control': 'public, max-age=3600' });
|
|
});
|
|
|
|
// GET /api/media/:id/gif?t=<sec>&d=<sec>&w=<px> — a short looping GIF of an
|
|
// exact moment, cut from the server-cached copy (timestamp sharing). Palette
|
|
// is generated per clip (palettegen/paletteuse) so it stays small and clean.
|
|
const gifCache = new Map();
|
|
let gifRunning = 0;
|
|
app.get('/api/media/:id/gif', async (c) => {
|
|
const id = c.req.param('id');
|
|
if (!/^([A-Za-z0-9_-]{11}|upl_[a-f0-9]{12})$/.test(id)) return c.json({ ok: false, error: 'invalid id' }, 400);
|
|
const t = Math.max(0, Number(c.req.query('t')) || 0);
|
|
const d = Math.min(6, Math.max(1, Number(c.req.query('d')) || 3));
|
|
const w = Math.min(640, Math.max(240, Math.round((Number(c.req.query('w')) || 480) / 2) * 2));
|
|
const path = await videoSourcePath(id);
|
|
if (!path) return c.json({ ok: false, error: 'this video is not cached on the server yet — play it once, then try again' }, 404);
|
|
const key = `${path}|${t.toFixed(1)}|${d}|${w}`;
|
|
let gif = gifCache.get(key);
|
|
if (!gif) {
|
|
if (gifRunning >= 2) return c.json({ ok: false, error: 'busy — try again in a moment' }, 503);
|
|
gifRunning++;
|
|
try {
|
|
gif = await new Promise((resolve, reject) => {
|
|
const child = spawn(FFMPEG, ['-v', 'error', '-ss', t.toFixed(2), '-t', String(d), '-i', path, '-an',
|
|
'-vf', `fps=12,scale=${w}:-2:flags=lanczos,split[a][b];[a]palettegen=max_colors=128:stats_mode=diff[p];[b][p]paletteuse=dither=bayer:bayer_scale=4`,
|
|
'-loop', '0', '-f', 'gif', 'pipe:1'], { stdio: ['ignore', 'pipe', 'pipe'] });
|
|
const chunks = [];
|
|
let err = '';
|
|
child.stdout.on('data', (b) => chunks.push(b));
|
|
child.stderr.on('data', (b) => { err = (err + b).slice(-1500); });
|
|
child.on('error', reject);
|
|
child.on('close', (code) => {
|
|
const buf = Buffer.concat(chunks);
|
|
if (code !== 0 || buf.length < 100) reject(new Error(err.trim() || 'could not make the GIF'));
|
|
else resolve(buf);
|
|
});
|
|
});
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
} finally {
|
|
gifRunning--;
|
|
}
|
|
gifCache.set(key, gif);
|
|
if (gifCache.size > 40) gifCache.delete(gifCache.keys().next().value);
|
|
}
|
|
return c.body(gif, 200, {
|
|
'Content-Type': 'image/gif',
|
|
'Cache-Control': 'public, max-age=86400',
|
|
'Content-Disposition': `inline; filename="${id}-${Math.floor(t)}s.gif"`,
|
|
});
|
|
});
|
|
|
|
// GET /api/media/:id/clip?start=&end=&fmt=mp3|m4r — a soundbite of the
|
|
// server-cached audio with short fades: MP3 (≤ 60 s) or an iPhone ringtone
|
|
// (.m4r = AAC in an iPod MP4 container, ≤ 40 s — Apple's ringtone limit).
|
|
app.get('/api/media/:id/clip', async (c) => {
|
|
const id = c.req.param('id');
|
|
if (!/^([A-Za-z0-9_-]{11}|upl_[a-f0-9]{12})$/.test(id)) return c.json({ ok: false, error: 'invalid id' }, 400);
|
|
const fmt = c.req.query('fmt') === 'm4r' ? 'm4r' : 'mp3';
|
|
const start = Math.max(0, Number(c.req.query('start')) || 0);
|
|
const maxLen = fmt === 'm4r' ? 40 : 60;
|
|
const len = Math.min(maxLen, Math.max(1, (Number(c.req.query('end')) || start + 20) - start));
|
|
const path = await audioSourcePath(id);
|
|
if (!path) return c.json({ ok: false, error: 'this video is not cached on the server yet — play it once, then try again' }, 404);
|
|
const fade = Math.min(0.5, len / 4);
|
|
const af = `afade=t=in:st=0:d=${fade},afade=t=out:st=${(len - fade).toFixed(2)}:d=${fade}`;
|
|
const args = ['-v', 'error', '-ss', start.toFixed(2), '-t', len.toFixed(2), '-i', path, '-vn', '-af', af];
|
|
// A ringtone must be a regular (non-fragmented) MP4 for iPhone imports,
|
|
// which needs a seekable output — so .m4r goes through a temp file.
|
|
const tmpOut = fmt === 'm4r' ? `${tmpdir()}/ytp-clip-${id}-${Date.now()}-${Math.random().toString(36).slice(2)}.m4r` : null;
|
|
if (fmt === 'mp3') args.push('-c:a', 'libmp3lame', '-b:a', '192k', '-f', 'mp3', 'pipe:1');
|
|
else args.push('-c:a', 'aac', '-b:a', '192k', '-movflags', '+faststart', '-f', 'ipod', '-y', tmpOut);
|
|
try {
|
|
const buf = await new Promise((resolve, reject) => {
|
|
const child = spawn(FFMPEG, args, { stdio: ['ignore', 'pipe', 'pipe'] });
|
|
const chunks = [];
|
|
let err = '';
|
|
child.stdout.on('data', (b) => chunks.push(b));
|
|
child.stderr.on('data', (b) => { err = (err + b).slice(-1500); });
|
|
child.on('error', reject);
|
|
child.on('close', (code) => {
|
|
let out = Buffer.concat(chunks);
|
|
if (tmpOut) {
|
|
try { out = readFileSync(tmpOut); } catch { out = Buffer.alloc(0); }
|
|
try { unlinkSync(tmpOut); } catch { /* gone */ }
|
|
}
|
|
if (code !== 0 || out.length < 200) reject(new Error(err.trim() || 'could not cut the soundbite'));
|
|
else resolve(out);
|
|
});
|
|
});
|
|
return c.body(buf, 200, {
|
|
'Content-Type': fmt === 'mp3' ? 'audio/mpeg' : 'audio/mp4',
|
|
'Content-Disposition': `attachment; filename="${id}-${Math.floor(start)}s.${fmt}"`,
|
|
'Cache-Control': 'public, max-age=86400',
|
|
});
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
app.get('/api/media/:id/status', async (c) => {
|
|
try {
|
|
return c.json({ ok: true, ...(await media.status(c.req.param('id'))) }, 200, { 'Cache-Control': 'no-store' });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// POST /api/media/:id/redownload — the client's "Broken" button. Idempotent:
|
|
// a job already queued/running for this id is just reported back.
|
|
app.post('/api/media/:id/redownload', async (c) => {
|
|
try {
|
|
return c.json({ ok: true, ...(await media.redownload(c.req.param('id'))) });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 400);
|
|
}
|
|
});
|
|
|
|
// Stream the server-cached copy as a save (OPFS writes it on the device).
|
|
function cachedDownloadResponse(videoId, fp, row) {
|
|
const file = Bun.file(`${MEDIA_DIR}/${videoId}.${row.gen}.mp4`);
|
|
if (fp) recordVideoAccess(fp, { id: videoId }).catch(() => {});
|
|
media.touch(videoId);
|
|
return new Response(file, {
|
|
status: 200,
|
|
headers: {
|
|
'Content-Type': 'video/mp4',
|
|
'Content-Length': String(file.size),
|
|
'Content-Disposition': `attachment; filename="${videoId}.mp4"`,
|
|
'Cache-Control': 'no-store',
|
|
'Access-Control-Allow-Origin': '*',
|
|
...(row.sha256 ? { 'X-Content-SHA256': row.sha256, 'Access-Control-Expose-Headers': 'X-Content-SHA256' } : {}),
|
|
},
|
|
});
|
|
}
|
|
|
|
// GET /api/download/:videoId
|
|
// Saves go through the server cache: the fetch runs as a server-owned job
|
|
// (a disconnecting client no longer kills it — the copy lands for next
|
|
// time), then the validated file is streamed so OPFS can store it. The
|
|
// browser never contacts YouTube CDN directly (CORS would block it).
|
|
app.get('/api/download/:videoId', async (c) => {
|
|
const videoId = (c.req.param('videoId') || '').replace(/[/\\:?<>|*"]/g, '').trim();
|
|
if (!videoId) return c.json({ ok: false, error: 'missing videoId' }, 400);
|
|
countView(c, videoId);
|
|
|
|
// Uploads are already a single file on disk — hand it over as-is.
|
|
if (isUpload(videoId)) {
|
|
const u = await notesDb.getUpload(videoId);
|
|
if (!u) return c.json({ ok: false, error: 'upload not found' }, 404);
|
|
const f = Bun.file(uploads.filePath(u));
|
|
return new Response(f, { status: 200, headers: {
|
|
'Content-Type': u.mime, 'Content-Length': String(f.size),
|
|
'Content-Disposition': `attachment; filename="${u.id}.${u.ext}"`, 'Cache-Control': 'no-store',
|
|
} });
|
|
}
|
|
const fp = c.req.query('fp');
|
|
|
|
// ?edit=1&keep=s-e,s-e — "Edit & download" path: download the source, then
|
|
// ffmpeg-trim it to the requested keep segments and stream the custom cut.
|
|
// Requires ffmpeg; there is no progressive fallback because the whole point
|
|
// is the server-side edit. Invalid/empty keep params are rejected up front.
|
|
if (c.req.query('edit') === '1') {
|
|
const keep = parseKeepParam(c.req.query('keep') || '');
|
|
if (!keep.length) return c.json({ ok: false, error: 'missing or invalid keep segments' }, 400);
|
|
try {
|
|
const row = await media.ensureCached(videoId, { priority: HIGH }).catch(() => null);
|
|
const src = row ? await media.filePath(videoId, row.gen) : null;
|
|
return await ytdlpEditedDownloadResponse(videoId, fp, keep, c.req.raw.signal, src);
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
}
|
|
|
|
// Server cache first — plain saves and ?mux=1 both resolve to the same
|
|
// canonical ≤720p H.264 + AAC copy.
|
|
let cacheErr = null;
|
|
try {
|
|
const row = await media.ensureCached(videoId, { priority: HIGH });
|
|
// A device that can't decode HEVC must not save the HEVC copy — give it
|
|
// the legacy H.264 save instead (the ?mux=1 path included).
|
|
if (servableTo(row, c.req.query('hevc') === '1')) return cachedDownloadResponse(videoId, fp, row);
|
|
cacheErr = Object.assign(new Error('cached copy is HEVC; client did not ask for it'), { code: 'SKIPPED' });
|
|
} catch (err) {
|
|
cacheErr = err;
|
|
console.warn(`[ytplayer] cache save unavailable for ${videoId} (${err.code || 'error'}): ${err.message}`);
|
|
}
|
|
|
|
// Legacy per-request saves — only when the cache can't hold this video.
|
|
// A cache job that downloaded and FAILED validation already tried the mux
|
|
// format, so go straight to the progressive-first default in that case.
|
|
// ?mux=1 — "Save before playing" path: bestvideo up to 720p PLUS bestaudio
|
|
// compiled into one mp4 with ffmpeg on the server. Falls back to the
|
|
// progressive single-file save below when ffmpeg is missing or the merge
|
|
// fails.
|
|
if (c.req.query('mux') === '1' && cacheErr && cacheErr.code === 'SKIPPED') {
|
|
try {
|
|
return await ytdlpDownloadResponse(videoId, fp, [
|
|
'-f', SAVE_MUX_FORMAT,
|
|
'--merge-output-format', 'mp4',
|
|
], c.req.raw.signal);
|
|
} catch (err) {
|
|
console.warn(`[ytplayer] mux download failed for ${videoId}, falling back to progressive:`, err.message);
|
|
}
|
|
}
|
|
|
|
// Default save — best progressive (audio+video single-file) format when
|
|
// one exists, otherwise merge bestvideo (≤720p) + bestaudio with ffmpeg.
|
|
// YouTube now serves many videos with NO progressive format at all, which
|
|
// made the old progressive-only selector fail every save with
|
|
// "Requested format is not available".
|
|
try {
|
|
return await ytdlpDownloadResponse(videoId, fp, [
|
|
'-f', SAVE_DEFAULT_FORMAT,
|
|
'--merge-output-format', 'mp4',
|
|
], c.req.raw.signal);
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// ============================================================================
|
|
// Online profiles — named cross-device sync. The NAME IS THE PASSKEY:
|
|
// anyone who knows it can load and overwrite that profile, so clients are
|
|
// encouraged to use the random generator. No other auth by design.
|
|
// ============================================================================
|
|
const PROFILE_NAME_RE = /^[A-Za-z0-9][A-Za-z0-9_-]{2,39}$/;
|
|
const PROFILE_MAX_BYTES = 2_000_000; // full data blob; typical payloads are ~KBs
|
|
|
|
const RAND_ADJ = ['amber', 'brave', 'calm', 'coral', 'crimson', 'dusty', 'gentle', 'golden',
|
|
'hidden', 'ivory', 'jade', 'lunar', 'mellow', 'misty', 'noble', 'quiet',
|
|
'rapid', 'silver', 'solar', 'stormy', 'swift', 'velvet', 'wild', 'zesty'];
|
|
const RAND_NOUN = ['falcon', 'harbor', 'willow', 'ember', 'canyon', 'meadow', 'otter', 'pine',
|
|
'raven', 'reef', 'sparrow', 'summit', 'thicket', 'tundra', 'brook', 'cedar',
|
|
'dune', 'fjord', 'glade', 'heron', 'lagoon', 'maple', 'prairie', 'wren'];
|
|
function randomProfileName() {
|
|
const a = RAND_ADJ[Math.floor(Math.random() * RAND_ADJ.length)];
|
|
const n = RAND_NOUN[Math.floor(Math.random() * RAND_NOUN.length)];
|
|
return `${a}-${n}-${1000 + Math.floor(Math.random() * 9000)}`;
|
|
}
|
|
|
|
// POST /api/profile/create
|
|
// Body: { name?, data? } — empty/absent name asks the server to generate a
|
|
// unique random one. Fails with 409 when the requested name is taken.
|
|
app.post('/api/profile/create', async (c) => {
|
|
let body;
|
|
try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
|
|
|
let name = (body.name || '').trim().toLowerCase();
|
|
const dataJson = JSON.stringify(body.data && typeof body.data === 'object' ? body.data : {});
|
|
if (dataJson.length > PROFILE_MAX_BYTES) return c.json({ ok: false, error: 'profile data too large' }, 413);
|
|
|
|
try {
|
|
if (name) {
|
|
if (!PROFILE_NAME_RE.test(name)) {
|
|
return c.json({ ok: false, error: 'invalid name — 3-40 characters: letters, digits, - or _' }, 400);
|
|
}
|
|
if (!(await createProfile(name, dataJson))) {
|
|
return c.json({ ok: false, error: `“${name}” is already taken — pick another name` }, 409);
|
|
}
|
|
} else {
|
|
let created = false;
|
|
for (let tries = 0; tries < 20 && !created; tries++) {
|
|
name = randomProfileName();
|
|
created = await createProfile(name, dataJson);
|
|
}
|
|
if (!created) return c.json({ ok: false, error: 'could not generate a unique name — try again' }, 500);
|
|
}
|
|
const row = await getProfile(name);
|
|
return c.json({ ok: true, name, updatedAt: row ? row.updatedAt : 0 });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// GET /api/profile/load?name=<name>
|
|
app.get('/api/profile/load', async (c) => {
|
|
const name = (c.req.query('name') || '').trim().toLowerCase();
|
|
if (!name) return c.json({ ok: false, error: 'missing name' }, 400);
|
|
try {
|
|
const row = await getProfile(name);
|
|
if (!row) return c.json({ ok: false, error: 'profile not found' }, 404);
|
|
let data = {};
|
|
try { data = JSON.parse(row.data || '{}'); } catch { /* corrupt blob — hand back empty */ }
|
|
return c.json({ ok: true, name, data, updatedAt: row.updatedAt });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// POST /api/profile/save
|
|
// Body: { name, data } — updates an EXISTING profile only (404 otherwise).
|
|
app.post('/api/profile/save', async (c) => {
|
|
let body;
|
|
try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
|
|
|
const name = (body.name || '').trim().toLowerCase();
|
|
if (!name) return c.json({ ok: false, error: 'missing name' }, 400);
|
|
const dataJson = JSON.stringify(body.data && typeof body.data === 'object' ? body.data : {});
|
|
if (dataJson.length > PROFILE_MAX_BYTES) return c.json({ ok: false, error: 'profile data too large' }, 413);
|
|
|
|
try {
|
|
if (!(await saveProfile(name, dataJson))) {
|
|
return c.json({ ok: false, error: 'profile not found' }, 404);
|
|
}
|
|
const row = await getProfile(name);
|
|
return c.json({ ok: true, updatedAt: row ? row.updatedAt : 0 });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// ============================================================================
|
|
// Shared playlists — single-playlist sharing via a 10-character code.
|
|
// ============================================================================
|
|
const PLAYLIST_CODE_CHARS = 'abcdefghijklmnopqrstuvwxyz0123456789';
|
|
function randomPlaylistCode() {
|
|
let code = '';
|
|
for (let i = 0; i < 10; i++) {
|
|
code += PLAYLIST_CODE_CHARS[Math.floor(Math.random() * PLAYLIST_CODE_CHARS.length)];
|
|
}
|
|
return code;
|
|
}
|
|
|
|
// A shared playlist is the ONLY path by which one user's video objects reach
|
|
// another user's DOM, so the blob is rebuilt field-by-field here rather than
|
|
// stored as sent. Anything not in this whitelist is dropped, and the two fields
|
|
// that end up in HTML attributes (thumbnail, channelUrl) must parse as http(s)
|
|
// URLs — otherwise a crafted `thumbnail` closes the src attribute and injects
|
|
// markup on the importing device. `custom` edits are dropped outright: their
|
|
// media only exists in the sharer's OPFS cache, so they are unplayable anywhere
|
|
// else and would just render as permanently broken entries.
|
|
const SHARED_STR_MAX = 300;
|
|
function safeStr(v, max = SHARED_STR_MAX) {
|
|
return typeof v === 'string' ? v.slice(0, max) : '';
|
|
}
|
|
function safeHttpUrl(v) {
|
|
if (typeof v !== 'string' || v.length > 2000) return '';
|
|
try {
|
|
const u = new URL(v);
|
|
return (u.protocol === 'http:' || u.protocol === 'https:') ? u.href : '';
|
|
} catch { return ''; }
|
|
}
|
|
function sanitizeSharedVideo(v) {
|
|
if (!v || typeof v !== 'object') return null;
|
|
const id = safeStr(v.id, 64);
|
|
if (!id || v.custom) return null;
|
|
const duration = Number(v.duration);
|
|
return {
|
|
id,
|
|
title: safeStr(v.title),
|
|
channel: safeStr(v.channel),
|
|
channelId: safeStr(v.channelId, 64),
|
|
channelUrl: safeHttpUrl(v.channelUrl),
|
|
duration: Number.isFinite(duration) && duration >= 0 ? duration : 0,
|
|
thumbnail: safeHttpUrl(v.thumbnail),
|
|
};
|
|
}
|
|
|
|
// POST /api/playlist/share
|
|
// Body: { playlist: { name, videos: [...] } }
|
|
app.post('/api/playlist/share', async (c) => {
|
|
let body;
|
|
try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
|
|
|
const pl = body?.playlist;
|
|
if (!pl || typeof pl !== 'object') {
|
|
return c.json({ ok: false, error: 'missing playlist' }, 400);
|
|
}
|
|
|
|
const name = typeof pl.name === 'string' ? pl.name.trim() : '';
|
|
if (!name || name.length > 200) {
|
|
return c.json({ ok: false, error: 'invalid name — must be non-empty and <= 200 chars' }, 400);
|
|
}
|
|
|
|
if (!Array.isArray(pl.videos) || pl.videos.length < 1 || pl.videos.length > 500) {
|
|
return c.json({ ok: false, error: 'invalid videos — must be an array of 1 to 500 videos' }, 400);
|
|
}
|
|
|
|
const videos = pl.videos.map(sanitizeSharedVideo).filter(Boolean);
|
|
if (!videos.length) return c.json({ ok: false, error: 'no usable videos in that playlist' }, 400);
|
|
|
|
const dataJson = JSON.stringify({ name, videos });
|
|
if (dataJson.length > PROFILE_MAX_BYTES) {
|
|
return c.json({ ok: false, error: 'playlist data too large' }, 413);
|
|
}
|
|
|
|
try {
|
|
let code = '';
|
|
let created = false;
|
|
for (let tries = 0; tries < 20 && !created; tries++) {
|
|
code = randomPlaylistCode();
|
|
created = await createSharedPlaylist(code, dataJson);
|
|
}
|
|
if (!created) {
|
|
return c.json({ ok: false, error: 'could not generate a unique share code — try again' }, 500);
|
|
}
|
|
return c.json({ ok: true, code });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// ---- Playlist inbox --------------------------------------------------------
|
|
// Send an already-shared playlist to another profile by name, and let that
|
|
// profile pick up what was sent. A delivery stores only the share code, so
|
|
// these endpoints never move video data around.
|
|
//
|
|
// Auth note: a profile name IS the credential in this app (see the profiles
|
|
// table), so `name` alone authorises reading and clearing an inbox. That is
|
|
// the same trust level as /api/profile/load, which already returns a whole
|
|
// profile for a bare name — these routes add no new exposure.
|
|
|
|
const INBOX_MAX_PENDING = 25;
|
|
|
|
// POST /api/playlist/send { to, from, code }
|
|
app.post('/api/playlist/send', async (c) => {
|
|
let body;
|
|
try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
|
|
|
const to = String(body?.to || '').trim();
|
|
const from = String(body?.from || '').trim();
|
|
const code = String(body?.code || '').trim().toLowerCase();
|
|
|
|
if (!to) return c.json({ ok: false, error: 'missing recipient' }, 400);
|
|
if (!PROFILE_NAME_RE.test(to)) {
|
|
return c.json({ ok: false, error: 'invalid username — 3-40 characters: letters, digits, - or _' }, 400);
|
|
}
|
|
if (from && !PROFILE_NAME_RE.test(from)) {
|
|
return c.json({ ok: false, error: 'invalid sender name' }, 400);
|
|
}
|
|
if (to.toLowerCase() === from.toLowerCase()) {
|
|
return c.json({ ok: false, error: 'that is your own username' }, 400);
|
|
}
|
|
if (!code) return c.json({ ok: false, error: 'missing code' }, 400);
|
|
|
|
try {
|
|
// The share code must exist — this is also where the title comes from, so
|
|
// a sender cannot attach arbitrary text to someone else's inbox.
|
|
const shared = await getSharedPlaylist(code);
|
|
if (!shared) return c.json({ ok: false, error: 'shared playlist not found' }, 404);
|
|
let pl = null;
|
|
try { pl = JSON.parse(shared.data || '{}'); } catch { /* corrupt blob */ }
|
|
if (!pl || !pl.name) return c.json({ ok: false, error: 'shared playlist not found' }, 404);
|
|
|
|
if (!(await getProfile(to))) {
|
|
return c.json({ ok: false, error: `no user named “${to}”` }, 404);
|
|
}
|
|
if (await countInbox(to) >= INBOX_MAX_PENDING) {
|
|
return c.json({ ok: false, error: `“${to}” has too many unopened playlists` }, 429);
|
|
}
|
|
|
|
await queueInboxPlaylist({
|
|
id: randomPlaylistCode() + randomPlaylistCode(),
|
|
toName: to,
|
|
fromName: from || null,
|
|
code,
|
|
title: String(pl.name).slice(0, 200),
|
|
});
|
|
return c.json({ ok: true });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// GET /api/playlist/inbox?name=<profile>
|
|
app.get('/api/playlist/inbox', async (c) => {
|
|
const name = (c.req.query('name') || '').trim();
|
|
if (!name) return c.json({ ok: false, error: 'missing name' }, 400);
|
|
if (!PROFILE_NAME_RE.test(name)) return c.json({ ok: true, items: [] });
|
|
try {
|
|
return c.json({ ok: true, items: await listInbox(name) });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// POST /api/playlist/inbox/dismiss { name, id }
|
|
app.post('/api/playlist/inbox/dismiss', async (c) => {
|
|
let body;
|
|
try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
|
const name = String(body?.name || '').trim();
|
|
const id = String(body?.id || '').trim();
|
|
if (!name || !id) return c.json({ ok: false, error: 'missing name or id' }, 400);
|
|
try {
|
|
await deleteInboxItem(name, id);
|
|
return c.json({ ok: true });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// GET /api/playlist/shared?code=<code>
|
|
app.get('/api/playlist/shared', async (c) => {
|
|
const code = (c.req.query('code') || '').trim().toLowerCase();
|
|
if (!code) return c.json({ ok: false, error: 'missing code' }, 400);
|
|
|
|
try {
|
|
const row = await getSharedPlaylist(code);
|
|
if (!row) {
|
|
return c.json({ ok: false, error: 'shared playlist not found' }, 404);
|
|
}
|
|
let pl = null;
|
|
try { pl = JSON.parse(row.data || '{}'); } catch { /* corrupt blob */ }
|
|
if (!pl || typeof pl !== 'object' || !pl.name || !Array.isArray(pl.videos)) {
|
|
return c.json({ ok: false, error: 'shared playlist not found' }, 404);
|
|
}
|
|
return c.json({
|
|
ok: true,
|
|
code,
|
|
playlist: { name: pl.name, videos: pl.videos },
|
|
createdAt: row.createdAt,
|
|
});
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
function isYoutubeHost(hostname) {
|
|
const h = (hostname || '').toLowerCase();
|
|
return h === 'youtube.com' || h.endsWith('.youtube.com') ||
|
|
h === 'youtu.be' || h.endsWith('.youtu.be');
|
|
}
|
|
|
|
function isValidYoutubeUrl(rawUrl) {
|
|
if (typeof rawUrl !== 'string' || !rawUrl.trim()) return false;
|
|
try {
|
|
const u = new URL(rawUrl.trim());
|
|
return (u.protocol === 'http:' || u.protocol === 'https:') && isYoutubeHost(u.hostname);
|
|
} catch {
|
|
return false;
|
|
}
|
|
}
|
|
|
|
// GET /api/playlist/expand?url=<youtube playlist url>
|
|
// Expands a YouTube playlist URL into slim video cards without resolving streams.
|
|
app.get('/api/playlist/expand', async (c) => {
|
|
const url = (c.req.query('url') || '').trim();
|
|
if (!url || !isValidYoutubeUrl(url)) {
|
|
return c.json({ ok: false, error: 'invalid YouTube URL' }, 400);
|
|
}
|
|
|
|
try {
|
|
const out = await runYtdlpResilient([
|
|
url,
|
|
'--dump-json', '--flat-playlist',
|
|
'--no-warnings', '--ignore-errors',
|
|
'--playlist-end', '201',
|
|
], { pooled: true });
|
|
const parsed = parseCards(out);
|
|
let title = '';
|
|
for (const line of out.split('\n')) {
|
|
const t = line.trim();
|
|
if (!t) continue;
|
|
try {
|
|
const j = JSON.parse(t);
|
|
title = title || pick(j, 'playlist_title', 'playlist');
|
|
if (title) break;
|
|
} catch { /* skip */ }
|
|
}
|
|
const truncated = parsed.length > 200;
|
|
const entries = parsed.slice(0, 200);
|
|
const resp = {
|
|
ok: true,
|
|
title: title || '',
|
|
entries,
|
|
};
|
|
if (truncated) resp.truncated = true;
|
|
return c.json(resp);
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 502);
|
|
}
|
|
});
|
|
|
|
// POST /api/user/sync
|
|
// Body: { fingerprint, playlists?, recentVideo?, appVersion? }
|
|
app.post('/api/user/sync', async (c) => {
|
|
let body;
|
|
try { body = await c.req.json(); } catch { return c.json({ ok: false, error: 'invalid JSON' }, 400); }
|
|
|
|
const fp = (body.fingerprint || '').trim();
|
|
if (!fp) return c.json({ ok: false, error: 'missing fingerprint' }, 400);
|
|
|
|
try {
|
|
await upsertUser({ fingerprint: fp, appVersion: body.appVersion, playlists: body.playlists });
|
|
if (body.recentVideo && body.recentVideo.id) {
|
|
await recordVideoAccess(fp, body.recentVideo);
|
|
}
|
|
return c.json({ ok: true });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// GET /api/user/data?fp=<fingerprint>
|
|
app.get('/api/user/data', async (c) => {
|
|
const fp = (c.req.query('fp') || '').trim();
|
|
if (!fp) return c.json({ ok: false, error: 'missing fingerprint' }, 400);
|
|
try {
|
|
const userData = await getUserData(fp);
|
|
return c.json({ ok: true, ...userData, version: APP_VERSION });
|
|
} catch (err) {
|
|
return c.json({ ok: false, error: err.message }, 500);
|
|
}
|
|
});
|
|
|
|
// ============================================================================
|
|
// Phone remote for a desktop instance — see remote.js
|
|
// ============================================================================
|
|
const remote = createRemoteHub({ requireSameNetwork: process.env.REMOTE_SAME_NETWORK === '1' });
|
|
const party = createPartyHub();
|
|
|
|
// Bun allows ONE websocket handler per server: party sockets are tagged
|
|
// (ws.data.hub === 'party'), everything else belongs to the remote relay.
|
|
// …and P2P sockets are tagged ws.data.hub === 'p2p' (p2p-hub.js).
|
|
const pickHub = (ws) => (ws.data && ws.data.hub === 'party' ? party.websocket
|
|
: ws.data && ws.data.hub === 'p2p' ? p2pHub.websocket : remote.websocket);
|
|
const websocketHandler = {
|
|
maxPayloadLength: 256 * 1024,
|
|
idleTimeout: 120,
|
|
open: (ws) => pickHub(ws).open(ws),
|
|
message: (ws, msg) => pickHub(ws).message(ws, msg),
|
|
close: (ws, code, reason) => pickHub(ws).close(ws, code, reason),
|
|
};
|
|
|
|
// The public IP a request came from. Behind Traefik that is the first
|
|
// X-Forwarded-For hop; locally, the socket address.
|
|
function clientIpOf(req, server) {
|
|
const xff = (req.headers.get('x-forwarded-for') || '').split(',')[0].trim();
|
|
if (xff) return xff;
|
|
try { return server?.requestIP?.(req)?.address || ''; } catch { return ''; }
|
|
}
|
|
|
|
app.post('/api/remote/pair', async (c) => {
|
|
let body = {};
|
|
try { body = await c.req.json(); } catch { /* validated below */ }
|
|
const r = remote.pair({ code: body.code, name: body.name, ip: clientIpOf(c.req.raw, c.env) });
|
|
return c.json(r.body, r.status);
|
|
});
|
|
|
|
// QR for the pairing link, drawn server-side so the client needs no library.
|
|
// Only ever encodes this origin's own /?pair=<6 digits> URL.
|
|
app.get('/api/remote/qr/:code', async (c) => {
|
|
const code = c.req.param('code');
|
|
if (!/^\d{6}$/.test(code)) return c.text('bad code', 400);
|
|
const proto = c.req.header('x-forwarded-proto') || new URL(c.req.url).protocol.replace(':', '');
|
|
const host = c.req.header('x-forwarded-host') || c.req.header('host') || new URL(c.req.url).host;
|
|
// ?kind=present → the presenter-screen link instead of the phone-remote one.
|
|
const param = c.req.query('kind') === 'present' ? 'present' : 'pair';
|
|
const svg = await QRCode.toString(`${proto}://${host}/?${param}=${code}`, { type: 'svg', margin: 1, errorCorrectionLevel: 'M' });
|
|
return c.body(svg, 200, { 'Content-Type': 'image/svg+xml', 'Cache-Control': 'no-store' });
|
|
});
|
|
|
|
// ============================================================================
|
|
// Shared lyrics + chapters, API tokens, admin page — see notes.js
|
|
// ============================================================================
|
|
const notes = registerNoteRoutes(app, {
|
|
db: notesDb,
|
|
getProfile,
|
|
profileNameRe: PROFILE_NAME_RE,
|
|
runYtdlp: runYtdlpResilient,
|
|
adminPassword: process.env.ADMIN_PASSWORD || '',
|
|
backupDir: pathJoin(dirname(process.env.DB_PATH || './data/ytplayer.db'), 'backups'),
|
|
adminHtmlPath: './public/admin.html',
|
|
workerToken: process.env.LYRICS_WORKER_TOKEN || '',
|
|
});
|
|
|
|
// ============================================================================
|
|
// Uploads — the server's own searchable media library (see uploads.js)
|
|
// ============================================================================
|
|
const UPLOAD_DIR = process.env.UPLOAD_DIR || pathJoin(dirname(process.env.DB_PATH || './data/ytplayer.db'), 'uploads');
|
|
const uploads = registerUploadRoutes(app, {
|
|
db: notesDb,
|
|
ffmpeg: FFMPEG,
|
|
ffprobe: FFPROBE,
|
|
uploadDir: UPLOAD_DIR,
|
|
rangeFileResponse,
|
|
requireAdminOrToken: notes.requireAdminOrToken,
|
|
parseLrc,
|
|
sanitizeLyrics,
|
|
});
|
|
const isUpload = (id) => uploads.isUploadId(id);
|
|
|
|
// ============================================================================
|
|
// Peer-to-peer sharing — devices + holdings (docs/p2p-architecture.md)
|
|
// ============================================================================
|
|
async function fileForCid(cid) {
|
|
const row = await p2pDb.findMediaByCid(cid);
|
|
if (!row) return null;
|
|
const path = `${MEDIA_DIR}/${row.video_id}.${row.gen}.mp4`;
|
|
const f = Bun.file(path);
|
|
return (await f.exists()) ? { path, size: f.size } : null;
|
|
}
|
|
const p2p = registerP2pRoutes(app, { cfg: P2P, p2pDb, fileForCid, sha256Range });
|
|
const p2pHub = createP2pHub({ getDevice: p2pDb.getDevice, enabled: () => P2P.enabled });
|
|
|
|
// GET /api/p2p/holders?v=<videoId>|cid=<sha256> — who holds a copy. Rows are
|
|
// persistent; `stale` flags a holder not re-verified for P2P_STALE_DAYS.
|
|
app.get('/api/p2p/holders', p2p.gate, async (c) => {
|
|
const v = (c.req.query('v') || '').trim();
|
|
const cid = (c.req.query('cid') || '').trim().toLowerCase();
|
|
const hasCid = /^[0-9a-f]{64}$/.test(cid);
|
|
if (!hasCid && !/^[A-Za-z0-9_-]{6,64}$/.test(v)) return c.json({ ok: false, error: 'missing v or cid' }, 400);
|
|
const payload = await holdersPayload({
|
|
videoId: v, cid: hasCid ? cid : null, p2pDb, isOnline: p2pHub.isOnline,
|
|
staleDays: P2P.staleDays, serverHas: fileForCid,
|
|
});
|
|
return c.json(payload, 200, { 'Cache-Control': 'no-store' });
|
|
});
|
|
|
|
// GET /api/p2p/available — the "On other devices" playlist: every verified file
|
|
// that at least one sharing device holds, online ones first. Holders are never
|
|
// dropped for age; `stale` only flags an old check.
|
|
app.get('/api/p2p/available', p2p.gate, async (c) => {
|
|
const limit = Math.min(200, Math.max(1, Number(c.req.query('limit')) || 100));
|
|
const rows = await p2pDb.listSharedContent(limit * 2);
|
|
const out = [];
|
|
for (const r of rows) {
|
|
if (!r.holders) continue;
|
|
const hs = (await p2pDb.listHolders(r.cid, 50)).filter((h) => Number(h.share) === 1);
|
|
const online = hs.filter((h) => p2pHub.isOnline(h.device_id)).length;
|
|
const newest = Math.max(0, ...hs.map((h) => Number(h.last_verified_at) || 0));
|
|
let meta = {};
|
|
try { meta = JSON.parse(r.meta || '{}'); } catch { /* no meta */ }
|
|
out.push({
|
|
id: r.video_id, cid: r.cid, title: meta.title || '', channel: meta.channel || '',
|
|
duration: Number(r.duration) || 0, size: Number(r.size) || 0, height: Number(r.height) || 0,
|
|
holders: hs.length, online, lastVerifiedAt: newest,
|
|
stale: Date.now() - newest > P2P.staleDays * 86400_000,
|
|
serverHas: !!(await fileForCid(r.cid)),
|
|
});
|
|
}
|
|
out.sort((a, b) => (b.online - a.online) || (b.lastVerifiedAt - a.lastVerifiedAt));
|
|
return c.json({ ok: true, staleDays: P2P.staleDays, items: out.slice(0, limit) }, 200, { 'Cache-Control': 'no-store' });
|
|
});
|
|
|
|
// Device → server intake: the server hashes, validates (and scans when
|
|
// P2P_MALWARE_SCAN=1) before a cid is admitted. docs/p2p-architecture.md flow 7.
|
|
const intake = registerIntakeRoutes(app, {
|
|
cfg: P2P, p2pDb, gate: p2p.gate, requireDevice: p2p.requireDevice, admitFile,
|
|
validateMedia: (path, expected, opts) => validateMedia(path, expected,
|
|
{ ...opts, ffmpeg: FFMPEG, ffprobe: process.env.FFPROBE_PATH || 'ffprobe' }),
|
|
adopt: (videoId, path, info) => media.adoptFile(videoId, path, info),
|
|
serverHasCid: async (cid) => !!(await fileForCid(cid)),
|
|
});
|
|
|
|
// Source gone + server copy evicted → ask one online holder to send it back
|
|
// through intake (docs/p2p-architecture.md flow 8).
|
|
const p2pRehydrate = createRehydrator({
|
|
p2pDb, hub: p2pHub, enabled: () => P2P.enabled,
|
|
hasServerCopy: async (id) => !!(await media.getReady(id).catch(() => null)),
|
|
});
|
|
// Admin: P2P overview + revoke (frontend/admin.html → Peer-to-peer).
|
|
app.get('/api/admin/p2p', notes.requireAdminOrToken, async (c) => c.json({
|
|
ok: true,
|
|
config: {
|
|
enabled: P2P.enabled, malwareScan: P2P.malwareScan, staleDays: P2P.staleDays,
|
|
keepMinViews: P2P.keepMinViews, keepDays: P2P.keepDays, keepRecentDays: P2P.keepRecentDays,
|
|
},
|
|
stats: { ...(await p2pDb.p2pStats()), online: p2pHub.onlineCount(), intakeOpen: intake.openTickets(), intakeActive: intake.active() },
|
|
recent: await p2pDb.recentContent(30),
|
|
}, 200, { 'Cache-Control': 'no-store' }));
|
|
app.post('/api/admin/p2p/revoke', notes.requireAdminOrToken, async (c) => {
|
|
const body = await c.req.json().catch(() => ({}));
|
|
const cid = String(body.cid || '').toLowerCase();
|
|
if (!/^[0-9a-f]{64}$/.test(cid)) return c.json({ ok: false, error: 'bad cid' }, 400);
|
|
await p2pDb.revokeContent(cid);
|
|
console.warn(`[p2p] admin revoked ${cid.slice(0, 12)}`);
|
|
return c.json({ ok: true });
|
|
});
|
|
|
|
// ============================================================================
|
|
// GET /sw.js — serve the service worker with BUILD_TAG injected
|
|
//
|
|
// The raw sw.js file contains the placeholder `__BUILD_TAG__` which is
|
|
// replaced here with the actual BUILD_TAG string so the SW's cache name
|
|
// tracks the deployment automatically — no manual version bump needed.
|
|
// Served with no-store cache headers so browsers always re-fetch it and
|
|
// pick up the substituted value rather than a browser-cached stale copy.
|
|
// ============================================================================
|
|
let _swSource = null;
|
|
app.get('/sw.js', (c) => {
|
|
if (!_swSource) {
|
|
try {
|
|
_swSource = readFileSync('./public/sw.js', 'utf8');
|
|
} catch {
|
|
return c.text('Service worker not found', 404);
|
|
}
|
|
}
|
|
// Inject the build tag: replace the whole fallback expression with the
|
|
// real value. Matched by REGEX, not an exact string — an exact match broke
|
|
// the moment the fallback literal in sw.js was bumped ('v1.0.3' → 'v1.0.4'),
|
|
// after which the replacement silently did nothing, the SW version froze,
|
|
// and clients never saw another update no matter how many times we deployed.
|
|
const src = _swSource.replace(
|
|
/typeof __BUILD_TAG__ !== 'undefined' \? __BUILD_TAG__ : '[^']*'/,
|
|
JSON.stringify(BUILD_TAG)
|
|
);
|
|
if (src === _swSource) {
|
|
console.error('[sw] BUILD_TAG injection failed — placeholder not found in sw.js');
|
|
}
|
|
return sendCompressed(c, compressedEntry('sw.js', Buffer.from(src), MIME.js), 'no-store, no-cache, must-revalidate');
|
|
});
|
|
|
|
// ============================================================================
|
|
// Static files — serve the frontend/public directory
|
|
// Must come AFTER all /api routes so API takes priority
|
|
// ============================================================================
|
|
// index.html carries the build tag it belongs to (<meta name="ytp-build">),
|
|
// stamped here at request time the same way /sw.js gets it. The page compares
|
|
// that against /api/version, so "Update available" only ever shows when the
|
|
// running build really differs from the deployed one (see app.js).
|
|
// Compressed + ETagged static text. The shell is ~620 KB raw / ~160 KB gzip
|
|
// and every byte crosses the slow VPS→homelab link, so compress once per
|
|
// file content and keep it in memory. ETag = sha256 of the RAW bytes, so a
|
|
// `no-cache` revalidation costs a 304 instead of the whole file.
|
|
const COMPRESSIBLE = /\.(js|css|html|json|webmanifest|svg|txt)$/i;
|
|
const compressedCache = new Map(); // key -> { etag, raw, gz, br, type }
|
|
function compressedEntry(key, raw, type) {
|
|
let e = compressedCache.get(key);
|
|
const etag = '"' + createHash('sha256').update(raw).digest('hex').slice(0, 32) + '"';
|
|
if (e && e.etag === etag) return e;
|
|
e = {
|
|
etag, raw, type,
|
|
gz: Bun.gzipSync(raw, { level: 9 }),
|
|
br: brotliCompressSync(raw, { params: { [zlibConstants.BROTLI_PARAM_QUALITY]: 11 } }),
|
|
};
|
|
compressedCache.set(key, e);
|
|
return e;
|
|
}
|
|
function sendCompressed(c, e, cacheControl) {
|
|
const headers = { 'Content-Type': e.type, 'Cache-Control': cacheControl, ETag: e.etag, Vary: 'Accept-Encoding' };
|
|
const inm = c.req.header('if-none-match') || '';
|
|
if (inm.split(',').map((s) => s.trim()).includes(e.etag)) return new Response(null, { status: 304, headers });
|
|
const ae = c.req.header('accept-encoding') || '';
|
|
if (/\bbr\b/.test(ae)) return new Response(e.br, { headers: { ...headers, 'Content-Encoding': 'br' } });
|
|
if (/\bgzip\b/.test(ae)) return new Response(e.gz, { headers: { ...headers, 'Content-Encoding': 'gzip' } });
|
|
return new Response(e.raw, { headers });
|
|
}
|
|
const MIME = { js: 'text/javascript; charset=utf-8', css: 'text/css; charset=utf-8', html: 'text/html; charset=utf-8',
|
|
json: 'application/json', webmanifest: 'application/manifest+json', svg: 'image/svg+xml', txt: 'text/plain; charset=utf-8' };
|
|
|
|
let _indexSource = null;
|
|
function indexHtml(c) {
|
|
if (_indexSource === null) {
|
|
try { _indexSource = readFileSync('./public/index.html', 'utf8'); }
|
|
catch { return c.text('index.html not found', 404); }
|
|
}
|
|
const html = _indexSource.replace('__BUILD_TAG__', BUILD_TAG);
|
|
return sendCompressed(c, compressedEntry('index.html', Buffer.from(html), MIME.html), 'no-cache');
|
|
}
|
|
app.get('/', indexHtml);
|
|
app.get('/index.html', indexHtml);
|
|
|
|
// no-cache (revalidate every time) on the shell: the service worker is the
|
|
// only cache that should hold app files. With no headers at all, a browser
|
|
// may heuristically cache app.js and hand a new worker's install the old one.
|
|
app.get('/*', async (c, next) => {
|
|
const p = decodeURIComponent(new URL(c.req.url).pathname);
|
|
if (!COMPRESSIBLE.test(p) || p.includes('..') || p.startsWith('/api/')) return next();
|
|
const file = Bun.file('./public' + p);
|
|
if (!(await file.exists())) return next();
|
|
const raw = Buffer.from(await file.arrayBuffer());
|
|
const ext = p.slice(p.lastIndexOf('.') + 1).toLowerCase();
|
|
return sendCompressed(c, compressedEntry(p, raw, MIME[ext] || 'application/octet-stream'), 'no-cache');
|
|
});
|
|
app.use('/*', serveStatic({ root: './public', onFound: (_path, c) => { c.header('Cache-Control', 'no-cache'); } }));
|
|
// SPA fallback — return index.html for any unmatched path
|
|
app.get('/*', indexHtml);
|
|
|
|
// ============================================================================
|
|
// Boot
|
|
// ============================================================================
|
|
async function main() {
|
|
await initDb();
|
|
await p2pDb.initP2pSchema();
|
|
console.log(`[ytplayer] DB ready`);
|
|
await media.init();
|
|
notes.startBackups();
|
|
console.log(`[ytplayer] Starting on port ${PORT}`);
|
|
|
|
// Bun.serve is the native Bun HTTP server
|
|
Bun.serve({
|
|
port: PORT,
|
|
// /ws/remote is upgraded here, before Hono — see remote.js.
|
|
fetch(req, server) {
|
|
const path = new URL(req.url).pathname;
|
|
if (path === '/ws/remote') return remote.upgrade(req, server, clientIpOf(req, server));
|
|
if (path === '/ws/party') return party.upgrade(req, server, clientIpOf(req, server));
|
|
if (path === '/ws/p2p') return p2pHub.upgrade(req, server);
|
|
return app.fetch(req, server);
|
|
},
|
|
websocket: websocketHandler,
|
|
// Default is 10s, which killed /api/download proxy streams whenever the
|
|
// connection went idle mid-transfer. 240s then killed every save whose
|
|
// server-side yt-dlp phase (no bytes sent yet) ran longer than 4 min —
|
|
// long videos at the rate cap — and the client retried in a loop.
|
|
// 0 disables the idle timeout entirely (verified with a 300s request);
|
|
// runaway downloads are bounded by MAX_SAVE_SECONDS + kill-on-disconnect.
|
|
idleTimeout: 0,
|
|
// Bun's default request-body cap is 128 MB, which silently rejected every
|
|
// upload / device intake of a long video with a 413. Streaming routes
|
|
// enforce their own per-route limits (uploads 4 GB, intake P2P_INTAKE_MAX_BYTES).
|
|
maxRequestBodySize: 5 * 1024 ** 3,
|
|
});
|
|
|
|
console.log(`[ytplayer] Listening → http://localhost:${PORT}`);
|
|
}
|
|
|
|
main().catch((err) => {
|
|
console.error('[ytplayer] fatal:', err);
|
|
process.exit(1);
|
|
});
|