Channels: one loader tries, in order, the Videos/Live/Shorts tabs (normal channels), the uploads playlist (Topic channels), the Releases/Playlists tabs expanded into their tracks (artist/music channels whose only content is albums, e.g. Cathedral of Praise Worship), then the home page; bare names are resolved to an id first. Nothing public → ok with a message, never a raw yt-dlp error. Saves: SaveSlots caps concurrent offline saves (Settings → Offline cache → Videos saved at once, 1–8, default 4); queued saves show "Waiting for a free slot", and resumed saves now run up to the limit instead of one at a time.
2517 lines
113 KiB
JavaScript
2517 lines
113 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, setProfileSecret, createSharedPlaylist, getSharedPlaylist, queueInboxPlaylist, listInbox, deleteInboxItem, countInbox,
|
||
getMedia, upsertMedia, deleteMedia, listMedia, listMediaLru, touchMedia, mediaStats } from './db.js';
|
||
import { createMediaCache, HIGH, LOW, validateMedia, MediaSkip } 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 * as searchCacheDb from './search-cache.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;
|
||
}
|
||
|
||
// Auto-generated "<Artist> - Topic" channels (YouTube Music) have no Videos
|
||
// tab, so /videos fails with "This channel does not have a videos tab". Their
|
||
// tracks are still in the channel's uploads playlist (UC… → UU…) and on the
|
||
// channel home page, so those are tried next.
|
||
function channelUrlCandidates(c) {
|
||
const videos = channelToUrl(c);
|
||
const root = videos.replace(/\/videos$/, '');
|
||
const m = /(?:^|\/channel\/)(UC[\w-]{22})(?:[/?#]|$)/.exec(c.trim()) || /\/channel\/(UC[\w-]{22})/.exec(root);
|
||
const out = [videos];
|
||
if (m) out.push(`https://www.youtube.com/playlist?list=UU${m[1].slice(2)}`);
|
||
out.push(root);
|
||
return out;
|
||
}
|
||
|
||
// Name → channel id: the first search hit published by a channel with exactly
|
||
// that name (case-insensitive). "<Artist> - Topic" names come from YouTube Music.
|
||
async function channelIdForName(name) {
|
||
const want = name.trim().toLowerCase();
|
||
const out = await runYtdlpResilient([`ytsearch10:${name}`, '--dump-json', '--flat-playlist', '--no-warnings', '--ignore-errors'], { pooled: true });
|
||
for (const line of out.split('\n')) {
|
||
try {
|
||
const j = JSON.parse(line);
|
||
const ch = String(j.channel || j.uploader || '').trim().toLowerCase();
|
||
const id = j.channel_id || j.uploader_id;
|
||
if (ch === want && /^UC[\w-]{22}$/.test(id || '')) return id;
|
||
} catch { /* not JSON */ }
|
||
}
|
||
return null;
|
||
}
|
||
|
||
// 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 = 200; // hot in-memory layer
|
||
const SEARCH_CACHE_TTL_MS = 10 * 60_000; // fresh: answered from memory
|
||
const SEARCH_CACHE_STALE_MS = 15 * 24 * 3600_000; // persistent copy: stale-while-revalidate for 15 days
|
||
const searchCache = new Map(); // lowercased q -> { yt, fetchedAt } (YouTube cards only)
|
||
const inflightSearches = new Map(); // lowercased q -> Promise<{ yt, youtubeError? }>
|
||
|
||
// The YouTube side of one search (up to 200 cards), shared by every caller asking
|
||
// the same thing at the same time, and written to memory + the persistent cache.
|
||
function fetchYoutube(q) {
|
||
const key = q.toLowerCase();
|
||
if (inflightSearches.has(key)) return inflightSearches.get(key);
|
||
const p = (async () => {
|
||
let yt = [];
|
||
if (process.env.SEARCH_INNERTUBE !== '0') {
|
||
try {
|
||
yt = await innertube.searchDeep(q, { limit: searchCacheDb.PER_QUERY });
|
||
} 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 fetchedAt = Date.now();
|
||
if (searchCache.size >= SEARCH_CACHE_MAX) searchCache.delete(searchCache.keys().next().value);
|
||
searchCache.set(key, { yt, fetchedAt });
|
||
if (yt.length) searchCacheDb.rememberVideos(yt).catch(() => {});
|
||
if (yt.length) searchCacheDb.put(q, yt).catch((e) => console.warn(`[search] cache write: ${e.message}`));
|
||
return { yt, fetchedAt };
|
||
})().finally(() => inflightSearches.delete(key));
|
||
inflightSearches.set(key, p);
|
||
return p;
|
||
}
|
||
|
||
const searchLibrary = async (q) => {
|
||
// This server's own library first — it still answers when YouTube is unreachable.
|
||
try { return (await notesDb.listUploads({ q, limit: 20, listedOnly: true })).map(uploads.card); } catch { return []; }
|
||
};
|
||
|
||
// GET /api/search/local?q=<query> — videos this server already knows (from every
|
||
// earlier search) plus its own uploads; answers in milliseconds so the app can
|
||
// show something while the real search is still loading.
|
||
app.get('/api/search/local', async (c) => {
|
||
const q = (c.req.query('q') || '').trim();
|
||
if (!q) return c.json({ ok: false, error: 'empty query' }, 400);
|
||
const [mine, known] = await Promise.all([searchLibrary(q), searchCacheDb.searchVideos(q).catch(() => [])]);
|
||
return c.json({ ok: true, results: [...mine, ...known], local: true });
|
||
});
|
||
|
||
// GET /api/search?q=<query>[&refresh=1]
|
||
// Memory → persistent cache (fresh 10 min, then stale-while-revalidate up to 15
|
||
// days: answered instantly, refreshed in the background) → YouTube. `refresh=1`
|
||
// (the app's "Update search" button) skips every cache. Only a full success is
|
||
// cached, so a transient failure gets retried on the next request.
|
||
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 key = q.toLowerCase();
|
||
const mine = await searchLibrary(q);
|
||
const answer = (yt, extra = {}) => c.json({ ok: true, results: [...mine, ...yt], ...extra });
|
||
if (c.req.query('refresh') !== '1') {
|
||
let hit = searchCache.get(key);
|
||
if (!hit) {
|
||
const row = await searchCacheDb.get(q).catch(() => null);
|
||
if (row) { hit = { yt: row.results, fetchedAt: row.fetchedAt }; searchCache.set(key, hit); }
|
||
}
|
||
const age = hit ? Date.now() - hit.fetchedAt : Infinity;
|
||
if (hit && age < SEARCH_CACHE_TTL_MS) return answer(hit.yt, { fetchedAt: hit.fetchedAt });
|
||
if (hit && age < SEARCH_CACHE_STALE_MS) {
|
||
fetchYoutube(q).catch(() => { /* keep serving the stale copy */ });
|
||
return answer(hit.yt, { stale: true, fetchedAt: hit.fetchedAt });
|
||
}
|
||
}
|
||
try {
|
||
const { yt, fetchedAt } = await fetchYoutube(q);
|
||
return answer(yt, { fetchedAt });
|
||
} 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>
|
||
// Channel listing that works for every kind of channel, not one-off cases:
|
||
// 1. /videos, /streams, /shorts tabs — normal channels
|
||
// 2. uploads playlist (UC… → UU…) — auto-generated "- Topic" channels
|
||
// 3. /releases, /playlists → their tracks — artist/music channels whose only
|
||
// content is albums or playlists (no Videos tab, no uploads playlist)
|
||
// 4. channel home page
|
||
// A bare channel name is first resolved to its id with a search. A channel with
|
||
// nothing public answers ok with no results (never a raw yt-dlp error).
|
||
const ytdlpFlat = (url, end) => runYtdlpResilient([url, '--dump-json', '--flat-playlist', '--no-warnings', '--ignore-errors',
|
||
'--playlist-end', String(end)], { pooled: true });
|
||
|
||
// Playlist links (albums, playlists) in a flat tab listing.
|
||
function playlistUrlsIn(out) {
|
||
const urls = [];
|
||
for (const line of out.split('\n')) {
|
||
try {
|
||
const j = JSON.parse(line);
|
||
const u = String(j.url || j.webpage_url || '');
|
||
if (/[?&]list=/.test(u)) urls.push(u);
|
||
} catch { /* not JSON */ }
|
||
}
|
||
return urls;
|
||
}
|
||
|
||
async function loadChannel(chan) {
|
||
let candidates = channelUrlCandidates(chan);
|
||
if (!/^(https?:|@|UC[\w-]{22}$)/.test(chan)) {
|
||
const id = await channelIdForName(chan).catch(() => null);
|
||
if (id) candidates = channelUrlCandidates(id);
|
||
}
|
||
const root = candidates[candidates.length - 1];
|
||
// 1, 2, 4 — tabs, uploads playlist, home (home is last in the candidate list)
|
||
const direct = [candidates[0], `${root}/streams`, `${root}/shorts`, ...candidates.slice(1)];
|
||
for (const url of direct) {
|
||
try {
|
||
const out = await ytdlpFlat(url, CHANNEL_LIMIT);
|
||
const results = parseCards(out);
|
||
if (results.length) return { out, results };
|
||
} catch { /* next source */ }
|
||
}
|
||
// 3 — albums / playlists, expanded into tracks (a few at a time).
|
||
for (const tab of ['releases', 'playlists']) {
|
||
let lists = [];
|
||
try { lists = playlistUrlsIn(await ytdlpFlat(`${root}/${tab}`, 15)); } catch { continue; }
|
||
if (!lists.length) continue;
|
||
const results = [], seen = new Set();
|
||
let out = '';
|
||
for (let i = 0; i < lists.length && results.length < CHANNEL_LIMIT; i += 3) {
|
||
const batch = await Promise.all(lists.slice(i, i + 3).map((u) => ytdlpFlat(u, 40).catch(() => '')));
|
||
for (const o of batch) {
|
||
out += o + '\n';
|
||
for (const card of parseCards(o)) {
|
||
if (seen.has(card.id)) continue;
|
||
seen.add(card.id);
|
||
results.push(card);
|
||
}
|
||
}
|
||
}
|
||
if (results.length) return { out, results: results.slice(0, CHANNEL_LIMIT) };
|
||
}
|
||
return { out: '', results: [] };
|
||
}
|
||
|
||
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 { out, results } = await loadChannel(chan);
|
||
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 || (results[0]?.channel || ''), channelUrl, results,
|
||
...(results.length ? {} : { message: 'This channel has no public videos to show.' }) });
|
||
} catch (err) {
|
||
return c.json({ ok: false, error: 'Could not load this channel right now — try again.' }, 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';
|
||
// Set when MEDIA_DIR sits on a removable drive: the cache only uses the
|
||
// directory while this marker (created once on the drive itself) exists.
|
||
const MEDIA_VOLUME_MARKER = process.env.MEDIA_VOLUME_MARKER || null;
|
||
if (!MEDIA_VOLUME_MARKER || existsSync(MEDIA_VOLUME_MARKER)) mkdirSync(MEDIA_DIR, { recursive: true });
|
||
|
||
const media = createMediaCache({
|
||
dir: MEDIA_DIR,
|
||
volumeMarker: MEDIA_VOLUME_MARKER,
|
||
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).
|
||
// `etag` (strong, quoted) lets a resuming client send If-Range: when the file
|
||
// changed under it, the Range is ignored and the whole new file comes back
|
||
// with 200, so stale and fresh bytes are never stitched together.
|
||
function rangeFileResponse(c, path, contentType, cacheControl, { etag = null, headers: extra = {} } = {}) {
|
||
const file = Bun.file(path);
|
||
const total = file.size;
|
||
const headers = {
|
||
'Content-Type': contentType,
|
||
'Accept-Ranges': 'bytes',
|
||
'Cache-Control': cacheControl,
|
||
...(etag ? { ETag: etag } : {}),
|
||
...extra,
|
||
};
|
||
let range = c.req.header('range');
|
||
const ifRange = c.req.header('if-range');
|
||
if (range && ifRange && etag && ifRange.trim() !== etag) {
|
||
// Stale partial: send the whole current file. As a stream, because Bun
|
||
// applies the request's Range to a Bun.file body on its own.
|
||
return new Response(file.stream(), { status: 200, headers: { ...headers, 'Content-Length': String(total) } });
|
||
}
|
||
if (!range) return new Response(file, { status: 200, headers: { ...headers, 'Content-Length': String(total) } });
|
||
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).
|
||
// Saves fetch the copy in byte ranges (docs/resumable-downloads-plan.md):
|
||
// the ETag pins the generation, so a resumed save never mixes two copies.
|
||
const prepareSkips = new Map(); // videoId → { reason, at } for /prepare
|
||
const DOWNLOAD_EXPOSE = 'X-Content-SHA256, Content-Range, Content-Length, ETag, Accept-Ranges';
|
||
function cachedDownloadResponse(c, videoId, fp, row) {
|
||
if (fp && !c.req.header('range')) recordVideoAccess(fp, { id: videoId }).catch(() => {});
|
||
media.touch(videoId);
|
||
return rangeFileResponse(c, `${MEDIA_DIR}/${videoId}.${row.gen}.mp4`, 'video/mp4', 'no-store', {
|
||
etag: `"${videoId}.${row.gen}"`,
|
||
headers: {
|
||
'Content-Disposition': `attachment; filename="${videoId}.mp4"`,
|
||
'Access-Control-Allow-Origin': '*',
|
||
'Access-Control-Expose-Headers': DOWNLOAD_EXPOSE,
|
||
...(row.sha256 ? { 'X-Content-SHA256': row.sha256 } : {}),
|
||
},
|
||
});
|
||
}
|
||
|
||
// POST /share-target — the manifest share target when no service worker is in
|
||
// control yet (first visit): links and text still work; files need the worker.
|
||
app.post('/share-target', async (c) => {
|
||
let body = {};
|
||
try { body = await c.req.parseBody(); } catch { /* not a form */ }
|
||
const pick = (k) => (typeof body[k] === 'string' ? body[k].trim() : '');
|
||
const text = [pick('url'), pick('text'), pick('title')].filter(Boolean).join(' ').slice(0, 2000);
|
||
return c.redirect('/?shared=' + encodeURIComponent(text) + (body.media ? '&sharefail=1' : ''), 303);
|
||
});
|
||
|
||
// GET /api/download/:videoId/prepare[?hevc=1] — never blocks. Starts (or
|
||
// joins) the server-side fetch and reports where it is:
|
||
// ready → { gen, size, sha256 }: fetch /api/download/:id in ranges
|
||
// working → poll again (the server is still getting it from YouTube)
|
||
// legacy → this copy can't be served in ranges (too long for the cache,
|
||
// HEVC the device can't play, cache offline): use the old
|
||
// single-request save
|
||
app.get('/api/download/:videoId/prepare', async (c) => {
|
||
const videoId = (c.req.param('videoId') || '').trim();
|
||
const nocache = { 'Cache-Control': 'no-store' };
|
||
if (isUpload(videoId)) {
|
||
const u = await notesDb.getUpload(videoId);
|
||
if (!u) return c.json({ ok: false, error: 'upload not found' }, 404);
|
||
let size = 0;
|
||
try { size = Bun.file(uploads.filePath(u)).size; } catch { /* missing */ }
|
||
return c.json({ ok: true, state: 'ready', gen: 0, size, sha256: null, etag: `"${u.id}"`, ext: u.ext }, 200, nocache);
|
||
}
|
||
if (!media.isMediaId(videoId)) return c.json({ ok: false, error: 'invalid video id' }, 400);
|
||
const hevcOk = c.req.query('hevc') === '1';
|
||
const ready = await media.getReady(videoId);
|
||
if (ready) {
|
||
if (!servableTo(ready, hevcOk)) return c.json({ ok: true, state: 'legacy', reason: 'hevc' }, 200, nocache);
|
||
media.touch(videoId);
|
||
// The row's size can differ from the mp4 on disk (it is not the file the
|
||
// device fetches byte-for-byte), so report the file's real length.
|
||
const path = await media.filePath(videoId, ready.gen);
|
||
let size = 0;
|
||
try { size = path ? Bun.file(path).size : 0; } catch { /* vanished */ }
|
||
if (!size) return c.json({ ok: true, state: 'legacy', reason: 'cached file unavailable' }, 200, nocache);
|
||
return c.json({ ok: true, state: 'ready', gen: ready.gen, size, sha256: ready.sha256 || null,
|
||
etag: `"${videoId}.${ready.gen}"`, ext: 'mp4' }, 200, nocache);
|
||
}
|
||
// A job refused moments ago (too long for the cache, cache offline…) would
|
||
// just be refused again — answer from memory instead of re-probing YouTube.
|
||
const prior = prepareSkips.get(videoId);
|
||
if (prior && Date.now() - prior.at < 10 * 60_000) return c.json({ ok: true, state: 'legacy', reason: prior.reason }, 200, nocache);
|
||
// Start or join the job; its result is picked up by the next poll.
|
||
let skip = null;
|
||
media.ensureCached(videoId, { priority: HIGH }).catch((err) => {
|
||
skip = err;
|
||
if (err instanceof MediaSkip) {
|
||
prepareSkips.set(videoId, { reason: err.message, at: Date.now() });
|
||
if (prepareSkips.size > 2000) prepareSkips.clear();
|
||
}
|
||
});
|
||
await new Promise((r) => setTimeout(r, 0));
|
||
if (skip) {
|
||
const legacy = skip instanceof MediaSkip;
|
||
return c.json({ ok: legacy, state: legacy ? 'legacy' : 'failed', reason: skip.message }, legacy ? 200 : 500, nocache);
|
||
}
|
||
const st = await media.status(videoId);
|
||
if (st.status === 'failed') return c.json({ ok: true, state: 'legacy', reason: st.error || 'server fetch failed' }, 200, nocache);
|
||
return c.json({ ok: true, state: 'working', status: st.status }, 200, nocache);
|
||
});
|
||
|
||
// 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);
|
||
return rangeFileResponse(c, uploads.filePath(u), u.mime, 'no-store', {
|
||
etag: `"${u.id}"`,
|
||
headers: { 'Content-Disposition': `attachment; filename="${u.id}.${u.ext}"`, 'Access-Control-Expose-Headers': DOWNLOAD_EXPOSE },
|
||
});
|
||
}
|
||
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(c, 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)}`;
|
||
}
|
||
|
||
// ---- Optional protection (password or generated key) ----
|
||
// A profile without a secret behaves as before. With one, load/save/changing it
|
||
// need the secret (X-Profile-Secret header or `secret` in the body). Wrong
|
||
// guesses are throttled per profile+IP, so a name alone can't be brute-forced.
|
||
const profileFails = new Map(); // `${name}|${ip}` → { n, until }
|
||
const PROFILE_FAIL_LIMIT = 10, PROFILE_FAIL_WINDOW = 15 * 60_000;
|
||
const clientIp = (c) => (c.req.header('x-forwarded-for') || c.req.header('x-real-ip') || '').split(',')[0].trim() || 'local';
|
||
async function profileGate(c, name, row, provided) {
|
||
if (!row || !row.secretHash) return null;
|
||
const key = `${name}|${clientIp(c)}`;
|
||
const f = profileFails.get(key);
|
||
if (f && f.n >= PROFILE_FAIL_LIMIT && Date.now() < f.until) {
|
||
return c.json({ ok: false, error: 'too many wrong attempts — try again in 15 minutes' }, 429);
|
||
}
|
||
const secret = String(provided || '');
|
||
if (!secret) return c.json({ ok: false, needSecret: true, kind: row.secretKind, error: 'this profile is protected' }, 401);
|
||
let ok = false;
|
||
try { ok = await Bun.password.verify(secret, row.secretHash); } catch { ok = false; }
|
||
if (ok) { profileFails.delete(key); return null; }
|
||
const n = f && Date.now() < f.until ? f.n + 1 : 1;
|
||
profileFails.set(key, { n, until: Date.now() + PROFILE_FAIL_WINDOW });
|
||
if (profileFails.size > 5000) profileFails.clear();
|
||
return c.json({ ok: false, needSecret: true, wrong: true, kind: row.secretKind,
|
||
error: row.secretKind === 'key' ? 'wrong profile key' : 'wrong password' }, 401);
|
||
}
|
||
|
||
// GET /api/profile/info?name= — exists / protected, never the data.
|
||
app.get('/api/profile/info', async (c) => {
|
||
const name = (c.req.query('name') || '').trim().toLowerCase();
|
||
if (!name) return c.json({ ok: false, error: 'missing name' }, 400);
|
||
const row = await getProfile(name);
|
||
return c.json({ ok: true, exists: !!row, protected: !!(row && row.secretHash), kind: row ? row.secretKind : null },
|
||
200, { 'Cache-Control': 'no-store' });
|
||
});
|
||
|
||
// POST /api/profile/secret
|
||
// Body: { name, current?, secret: <new password/key> | null, kind: 'password'|'key' }
|
||
// Adds, changes or (secret: null) removes protection. Changing or removing an
|
||
// existing one needs the current secret.
|
||
app.post('/api/profile/secret', 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();
|
||
const row = name ? await getProfile(name) : null;
|
||
if (!row) return c.json({ ok: false, error: 'profile not found' }, 404);
|
||
const denied = await profileGate(c, name, row, body.current || c.req.header('x-profile-secret'));
|
||
if (denied) return denied;
|
||
if (body.secret === null) {
|
||
await setProfileSecret(name, null, null);
|
||
return c.json({ ok: true, protected: false });
|
||
}
|
||
const kind = body.kind === 'key' ? 'key' : 'password';
|
||
const secret = String(body.secret || '');
|
||
if (kind === 'password' && (secret.length < 8 || secret.length > 200)) {
|
||
return c.json({ ok: false, error: 'password must be 8–200 characters' }, 400);
|
||
}
|
||
if (kind === 'key' && !/^[A-Za-z0-9_-]{32,128}$/.test(secret)) {
|
||
return c.json({ ok: false, error: 'invalid profile key' }, 400);
|
||
}
|
||
const hash = await Bun.password.hash(secret);
|
||
await setProfileSecret(name, hash, kind);
|
||
return c.json({ ok: true, protected: true, kind });
|
||
});
|
||
|
||
// 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);
|
||
const denied = await profileGate(c, name, row, c.req.header('x-profile-secret'));
|
||
if (denied) return denied;
|
||
let data = {};
|
||
try { data = JSON.parse(row.data || '{}'); } catch { /* corrupt blob — hand back empty */ }
|
||
return c.json({ ok: true, name, data, updatedAt: row.updatedAt, protected: !!row.secretHash, kind: row.secretKind });
|
||
} 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 {
|
||
const existing = await getProfile(name);
|
||
if (!existing) return c.json({ ok: false, error: 'profile not found' }, 404);
|
||
const denied = await profileGate(c, name, existing, body.secret || c.req.header('x-profile-secret'));
|
||
if (denied) return denied;
|
||
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 {
|
||
const denied = await profileGate(c, name.toLowerCase(), await getProfile(name.toLowerCase()), c.req.header('x-profile-secret'));
|
||
if (denied) return denied;
|
||
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 {
|
||
const denied = await profileGate(c, name.toLowerCase(), await getProfile(name.toLowerCase()), body.secret || c.req.header('x-profile-secret'));
|
||
if (denied) return denied;
|
||
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,
|
||
backupDir: process.env.UPLOAD_BACKUP_DIR || null,
|
||
volumeMarker: process.env.UPLOAD_VOLUME_MARKER || null,
|
||
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); }
|
||
}
|
||
// Local css/js get `?v=<build>`: that URL's bytes never change, so the static
|
||
// handler below lets browsers keep it forever; a new build is a new URL.
|
||
const html = _indexSource.replace('__BUILD_TAG__', BUILD_TAG)
|
||
.replace(/((?:href|src)=")((?![a-z]+:|\/)[^"?]+\.(?:css|js))(")/g, `$1$2?v=${BUILD_TAG}$3`);
|
||
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();
|
||
// `?v=` equal to the running build → forever (a new build changes the URL).
|
||
// Any other / missing v keeps revalidating.
|
||
const forever = c.req.query('v') === BUILD_TAG;
|
||
return sendCompressed(c, compressedEntry(p, raw, MIME[ext] || 'application/octet-stream'),
|
||
forever ? 'public, max-age=31536000, immutable' : 'no-cache');
|
||
});
|
||
// Self-hosted fonts and icons rarely change: let the browser keep them for 30 days
|
||
// (the service worker precaches the shell anyway); everything else revalidates.
|
||
app.use('/*', serveStatic({ root: './public', onFound: (path, c) => {
|
||
c.header('Cache-Control', /(^|\/)(fonts|icons)\//.test(path) ? 'public, max-age=2592000' : '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);
|
||
});
|