/* ============================================================================ * db.js — libsql (embedded SQLite) database setup and query helpers * * The DB file lives at $DB_PATH (mounted volume in Docker so data persists * across container rebuilds). libsql is an open-source SQLite fork by Turso * with WAL-by-default for concurrent reads. * * Schema (three small tables): * users — one row per fingerprint; tracks last app version seen * playlists — one row per fingerprint; full playlist JSON blob * video_history — one row per (fingerprint, video_id); recent 200 entries * media_cache — one row per server-cached video file (see media-cache.js) * ========================================================================== */ import { createClient } from '@libsql/client'; import { join } from 'node:path'; const DB_PATH = process.env.DB_PATH || join(process.cwd(), 'data', 'ytplayer.db'); export const db = createClient({ url: 'file:' + DB_PATH, }); // ---- Schema ---------------------------------------------------------------- export async function initDb() { await db.executeMultiple(` PRAGMA journal_mode = WAL; PRAGMA synchronous = NORMAL; PRAGMA foreign_keys = ON; CREATE TABLE IF NOT EXISTS users ( fingerprint TEXT PRIMARY KEY, last_version TEXT, last_seen INTEGER NOT NULL DEFAULT (unixepoch()) ); CREATE TABLE IF NOT EXISTS playlists ( fingerprint TEXT PRIMARY KEY, data TEXT NOT NULL DEFAULT '[]', updated_at INTEGER NOT NULL DEFAULT (unixepoch()) ); CREATE TABLE IF NOT EXISTS video_history ( fingerprint TEXT NOT NULL, video_id TEXT NOT NULL, title TEXT, channel TEXT, thumbnail TEXT, duration REAL DEFAULT 0, accessed_at INTEGER NOT NULL DEFAULT (unixepoch()), PRIMARY KEY (fingerprint, video_id) ); CREATE INDEX IF NOT EXISTS idx_vh_fp_time ON video_history (fingerprint, accessed_at DESC); CREATE TABLE IF NOT EXISTS profiles ( name TEXT PRIMARY KEY, data TEXT NOT NULL DEFAULT '{}', created_at INTEGER NOT NULL DEFAULT (unixepoch()), updated_at INTEGER NOT NULL DEFAULT (unixepoch()) ); CREATE TABLE IF NOT EXISTS shared_playlists ( code TEXT PRIMARY KEY, data TEXT NOT NULL DEFAULT '{}', created_at INTEGER NOT NULL DEFAULT (unixepoch()) ); -- Playlists sent from one profile to another. Holds only a pointer to the -- shared_playlists row, never a copy of the videos, so a delivery costs a -- few bytes and the recipient decides whether to keep it. CREATE TABLE IF NOT EXISTS playlist_inbox ( id TEXT PRIMARY KEY, to_name TEXT NOT NULL, from_name TEXT, code TEXT NOT NULL, title TEXT NOT NULL, created_at INTEGER NOT NULL DEFAULT (unixepoch()) ); CREATE INDEX IF NOT EXISTS idx_inbox_to ON playlist_inbox (to_name, created_at DESC); -- Server-side media cache: one validated mp4 per video under MEDIA_DIR. -- The file on disk is ..mp4 (+ .m4a audio sidecar); gen -- bumps on every replacement so a URL carrying ?g= always maps to the -- same bytes. Timestamps here are ms epochs (backoff math needs ms). CREATE TABLE IF NOT EXISTS media_cache ( video_id TEXT PRIMARY KEY, status TEXT NOT NULL DEFAULT 'queued', -- queued|downloading|validating|ready|failed gen INTEGER NOT NULL DEFAULT 0, size INTEGER NOT NULL DEFAULT 0, -- mp4 + m4a bytes height INTEGER NOT NULL DEFAULT 0, vcodec TEXT, acodec TEXT, duration REAL NOT NULL DEFAULT 0, optimized INTEGER NOT NULL DEFAULT 0, meta TEXT NOT NULL DEFAULT '{}', priority INTEGER NOT NULL DEFAULT 1, auto INTEGER NOT NULL DEFAULT 0, attempts INTEGER NOT NULL DEFAULT 0, error TEXT, retry_at INTEGER NOT NULL DEFAULT 0, created_at INTEGER NOT NULL DEFAULT 0, updated_at INTEGER NOT NULL DEFAULT 0, last_access INTEGER NOT NULL DEFAULT 0, hits INTEGER NOT NULL DEFAULT 0 ); CREATE INDEX IF NOT EXISTS idx_media_lru ON media_cache (status, last_access); -- Shared per-video documents (kind = lyrics | chapters), visible to every -- user. The live copy is here; every save also lands in video_note_revs -- as a full snapshot, which is the server-side backup and undo history. CREATE TABLE IF NOT EXISTS video_notes ( video_id TEXT NOT NULL, kind TEXT NOT NULL, data TEXT NOT NULL, rev INTEGER NOT NULL, source TEXT NOT NULL DEFAULT 'user', -- user | auto | api | restore updated_by TEXT, updated_at INTEGER NOT NULL DEFAULT (unixepoch()), PRIMARY KEY (video_id, kind) ); CREATE TABLE IF NOT EXISTS video_note_revs ( id INTEGER PRIMARY KEY AUTOINCREMENT, video_id TEXT NOT NULL, kind TEXT NOT NULL, rev INTEGER NOT NULL, data TEXT NOT NULL, source TEXT NOT NULL DEFAULT 'user', updated_by TEXT, created_at INTEGER NOT NULL DEFAULT (unixepoch()) ); CREATE UNIQUE INDEX IF NOT EXISTS idx_note_revs ON video_note_revs (video_id, kind, rev); CREATE INDEX IF NOT EXISTS idx_note_revs_time ON video_note_revs (created_at DESC); -- API tokens for scripts (lyrics injection etc.). Only a SHA-256 of the -- token is stored; the plaintext is shown once when it is created. CREATE TABLE IF NOT EXISTS api_tokens ( id TEXT PRIMARY KEY, label TEXT NOT NULL, token_hash TEXT NOT NULL UNIQUE, created_at INTEGER NOT NULL DEFAULT (unixepoch()), last_used_at INTEGER ); `); } // ---- Shared video notes (lyrics / chapters) --------------------------------- function noteRow(r) { return { videoId: r.video_id, kind: r.kind, data: JSON.parse(r.data), rev: Number(r.rev), source: r.source, updatedBy: r.updated_by || null, updatedAt: Number(r.updated_at), }; } export async function getNotes(videoId) { const r = await db.execute({ sql: 'SELECT * FROM video_notes WHERE video_id = ?', args: [videoId] }); const out = {}; for (const row of r.rows) out[row.kind] = noteRow(row); return out; } // Write a new revision when `baseRev` still matches the live one (0 = none // yet). Returns { ok: true, rev } or { ok: false, current } on a conflict, so // two people editing at once never silently overwrite each other. `force` // skips the check (admin restore, API callers that ask for it). export async function saveNote({ videoId, kind, data, baseRev, source, updatedBy, force = false }) { const json = JSON.stringify(data); const tx = await db.transaction('write'); try { const cur = await tx.execute({ sql: 'SELECT rev FROM video_notes WHERE video_id = ? AND kind = ?', args: [videoId, kind], }); const liveRev = cur.rows[0] ? Number(cur.rows[0].rev) : 0; if (!force && Number(baseRev || 0) !== liveRev) { await tx.rollback(); const all = await getNotes(videoId); return { ok: false, current: all[kind] || null }; } const rev = liveRev + 1; await tx.execute({ sql: `INSERT INTO video_notes (video_id, kind, data, rev, source, updated_by, updated_at) VALUES (?, ?, ?, ?, ?, ?, unixepoch()) ON CONFLICT(video_id, kind) DO UPDATE SET data = excluded.data, rev = excluded.rev, source = excluded.source, updated_by = excluded.updated_by, updated_at = excluded.updated_at`, args: [videoId, kind, json, rev, source, updatedBy || null], }); await tx.execute({ sql: `INSERT INTO video_note_revs (video_id, kind, rev, data, source, updated_by, created_at) VALUES (?, ?, ?, ?, ?, ?, unixepoch())`, args: [videoId, kind, rev, json, source, updatedBy || null], }); await tx.commit(); return { ok: true, rev }; } catch (err) { try { await tx.rollback(); } catch { /* already closed */ } throw err; } finally { tx.close(); } } export async function listNoteRevs(videoId, kind, limit = 50) { const r = await db.execute({ sql: `SELECT rev, source, updated_by, created_at, length(data) AS bytes FROM video_note_revs WHERE video_id = ? AND kind = ? ORDER BY rev DESC LIMIT ?`, args: [videoId, kind, limit], }); return r.rows.map((x) => ({ rev: Number(x.rev), source: x.source, updatedBy: x.updated_by || null, createdAt: Number(x.created_at), bytes: Number(x.bytes), })); } export async function getNoteRev(videoId, kind, rev) { const r = await db.execute({ sql: 'SELECT * FROM video_note_revs WHERE video_id = ? AND kind = ? AND rev = ?', args: [videoId, kind, rev], }); const x = r.rows[0]; if (!x) return null; return { rev: Number(x.rev), data: JSON.parse(x.data), source: x.source, updatedBy: x.updated_by || null, createdAt: Number(x.created_at), }; } // Latest edits across every video — the admin page's moderation feed. export async function recentNoteRevs(limit = 100) { const r = await db.execute({ sql: `SELECT video_id, kind, rev, source, updated_by, created_at, length(data) AS bytes FROM video_note_revs ORDER BY created_at DESC, id DESC LIMIT ?`, args: [limit], }); return r.rows.map((x) => ({ videoId: x.video_id, kind: x.kind, rev: Number(x.rev), source: x.source, updatedBy: x.updated_by || null, createdAt: Number(x.created_at), bytes: Number(x.bytes), })); } export async function allNotes() { const r = await db.execute('SELECT * FROM video_notes ORDER BY video_id, kind'); return r.rows.map(noteRow); } // ---- API tokens --------------------------------------------------------------- export async function createApiToken({ id, label, tokenHash }) { await db.execute({ sql: 'INSERT INTO api_tokens (id, label, token_hash, created_at) VALUES (?, ?, ?, unixepoch())', args: [id, label, tokenHash], }); } export async function listApiTokens() { const r = await db.execute('SELECT id, label, created_at, last_used_at FROM api_tokens ORDER BY created_at DESC'); return r.rows.map((x) => ({ id: x.id, label: x.label, createdAt: Number(x.created_at), lastUsedAt: x.last_used_at == null ? null : Number(x.last_used_at), })); } export async function deleteApiToken(id) { const r = await db.execute({ sql: 'DELETE FROM api_tokens WHERE id = ?', args: [id] }); return (r.rowsAffected || 0) > 0; } // Resolve a token hash to its row and stamp last use. null = unknown token. export async function useApiToken(tokenHash) { const r = await db.execute({ sql: 'SELECT id, label FROM api_tokens WHERE token_hash = ?', args: [tokenHash] }); const x = r.rows[0]; if (!x) return null; db.execute({ sql: 'UPDATE api_tokens SET last_used_at = unixepoch() WHERE id = ?', args: [x.id] }).catch(() => {}); return { id: x.id, label: x.label }; } // ---- Media cache ------------------------------------------------------------- const MEDIA_COLS = new Set([ 'status', 'gen', 'size', 'height', 'vcodec', 'acodec', 'duration', 'optimized', 'meta', 'priority', 'auto', 'attempts', 'error', 'retry_at', 'created_at', 'updated_at', 'last_access', 'hits', ]); function rowToObj(res, row) { const o = {}; res.columns.forEach((c, i) => { o[c] = row[i]; }); return o; } export async function getMedia(videoId) { const r = await db.execute({ sql: 'SELECT * FROM media_cache WHERE video_id = ?', args: [videoId] }); return r.rows[0] ? rowToObj(r, r.rows[0]) : null; } // Insert or partially update a media row. Unknown keys are ignored. export async function upsertMedia(videoId, fields) { const cols = Object.keys(fields).filter((k) => MEDIA_COLS.has(k)); const all = ['video_id', ...cols]; const set = cols.length ? 'DO UPDATE SET ' + cols.map((c) => `${c} = excluded.${c}`).join(', ') : 'DO NOTHING'; await db.execute({ sql: `INSERT INTO media_cache (${all.join(', ')}) VALUES (${all.map(() => '?').join(', ')}) ON CONFLICT(video_id) ${set}`, args: [videoId, ...cols.map((c) => fields[c] ?? null)], }); } export async function deleteMedia(videoId) { await db.execute({ sql: 'DELETE FROM media_cache WHERE video_id = ?', args: [videoId] }); } export async function listMedia() { const r = await db.execute('SELECT * FROM media_cache'); return r.rows.map((row) => rowToObj(r, row)); } // Ready rows, least-recently-played first — the eviction order. export async function listMediaLru() { const r = await db.execute( "SELECT video_id, size, last_access FROM media_cache WHERE status = 'ready' ORDER BY last_access ASC", ); return r.rows.map((row) => rowToObj(r, row)); } export async function touchMedia(videoId, now) { await db.execute({ sql: 'UPDATE media_cache SET last_access = ?, hits = hits + 1 WHERE video_id = ?', args: [now, videoId], }); } export async function mediaStats() { const r = await db.execute( `SELECT SUM(CASE WHEN status = 'ready' THEN 1 ELSE 0 END) AS count, SUM(CASE WHEN status = 'ready' THEN size ELSE 0 END) AS bytes, SUM(CASE WHEN status IN ('queued','downloading','validating') THEN 1 ELSE 0 END) AS pending, SUM(CASE WHEN status = 'failed' THEN 1 ELSE 0 END) AS failed, SUM(CASE WHEN status = 'ready' AND optimized > 0 THEN 1 ELSE 0 END) AS optimized, SUM(CASE WHEN status = 'ready' AND vcodec = 'hevc' THEN 1 ELSE 0 END) AS hevc FROM media_cache`, ); const o = r.rows[0] ? rowToObj(r, r.rows[0]) : {}; return { count: o.count || 0, bytes: o.bytes || 0, pending: o.pending || 0, failed: o.failed || 0, optimized: o.optimized || 0, hevc: o.hevc || 0, }; } // ---- Profiles (named cross-device sync; the name acts as the passkey) ------ // Insert a new profile. Returns false when the name is already taken. export async function createProfile(name, dataJson) { try { await db.execute({ sql: `INSERT INTO profiles (name, data, created_at, updated_at) VALUES (?, ?, unixepoch(), unixepoch())`, args: [name, dataJson], }); return true; } catch (err) { const msg = String(err && err.message || err); if (msg.includes('UNIQUE') || msg.includes('PRIMARY KEY')) return false; throw err; } } export async function getProfile(name) { const r = await db.execute({ sql: 'SELECT data, updated_at FROM profiles WHERE name = ?', args: [name], }); const row = r.rows[0]; if (!row) return null; return { data: row.data, updatedAt: Number(row.updated_at) }; } // Update an EXISTING profile's data blob. Returns false when it doesn't exist // (saving must never implicitly create a profile — creation is a deliberate, // uniqueness-checked act). export async function saveProfile(name, dataJson) { const r = await db.execute({ sql: 'UPDATE profiles SET data = ?, updated_at = unixepoch() WHERE name = ?', args: [dataJson, name], }); return (r.rowsAffected || 0) > 0; } // ---- Shared Playlists ------------------------------------------------------ // Insert a new shared playlist. Returns false when the code is already taken. export async function createSharedPlaylist(code, dataJson) { try { await db.execute({ sql: `INSERT INTO shared_playlists (code, data, created_at) VALUES (?, ?, unixepoch())`, args: [code, dataJson], }); return true; } catch (err) { const msg = String(err && err.message || err); if (msg.includes('UNIQUE') || msg.includes('PRIMARY KEY')) return false; throw err; } } export async function getSharedPlaylist(code) { const r = await db.execute({ sql: 'SELECT data, created_at FROM shared_playlists WHERE code = ?', args: [code], }); const row = r.rows[0]; if (!row) return null; return { data: row.data, createdAt: Number(row.created_at) }; } // ---- Playlist inbox (profile-to-profile sends) ------------------------------ // Queue a shared playlist for a recipient profile. Caller must have verified // the recipient exists; `code` must already be a row in shared_playlists. export async function queueInboxPlaylist({ id, toName, fromName, code, title }) { await db.execute({ sql: `INSERT INTO playlist_inbox (id, to_name, from_name, code, title, created_at) VALUES (?, ?, ?, ?, ?, unixepoch())`, args: [id, toName, fromName || null, code, title], }); } // Pending deliveries for a profile, oldest first so they are offered in the // order they were sent. Capped so a flooded inbox can't blow up the response. export async function listInbox(toName, limit = 20) { const r = await db.execute({ sql: `SELECT id, from_name, code, title, created_at FROM playlist_inbox WHERE to_name = ? ORDER BY created_at ASC LIMIT ?`, args: [toName, limit], }); return r.rows.map((row) => ({ id: row.id, from: row.from_name || null, code: row.code, title: row.title, createdAt: Number(row.created_at), })); } // Remove one delivery once the recipient has kept or dismissed it. Scoped by // to_name so knowing an id alone cannot clear someone else's inbox. export async function deleteInboxItem(toName, id) { const r = await db.execute({ sql: 'DELETE FROM playlist_inbox WHERE to_name = ? AND id = ?', args: [toName, id], }); return (r.rowsAffected || 0) > 0; } // How many deliveries this recipient already has waiting — used to refuse a // send that would flood an inbox. export async function countInbox(toName) { const r = await db.execute({ sql: 'SELECT COUNT(*) AS n FROM playlist_inbox WHERE to_name = ?', args: [toName], }); return Number(r.rows[0] ? r.rows[0].n : 0); } // ---- Helpers --------------------------------------------------------------- // Upsert the users row and optionally update playlists. export async function upsertUser({ fingerprint, appVersion, playlists }) { await db.execute({ sql: `INSERT INTO users (fingerprint, last_version, last_seen) VALUES (?, ?, unixepoch()) ON CONFLICT(fingerprint) DO UPDATE SET last_version = excluded.last_version, last_seen = excluded.last_seen`, args: [fingerprint, appVersion || null], }); if (Array.isArray(playlists)) { await db.execute({ sql: `INSERT INTO playlists (fingerprint, data, updated_at) VALUES (?, ?, unixepoch()) ON CONFLICT(fingerprint) DO UPDATE SET data = excluded.data, updated_at = excluded.updated_at`, args: [fingerprint, JSON.stringify(playlists)], }); } } // Record a video access (insert or bump accessed_at). export async function recordVideoAccess(fingerprint, video) { await db.execute({ sql: `INSERT INTO video_history (fingerprint, video_id, title, channel, thumbnail, duration, accessed_at) VALUES (?, ?, ?, ?, ?, ?, unixepoch()) ON CONFLICT(fingerprint, video_id) DO UPDATE SET title = excluded.title, channel = excluded.channel, thumbnail = excluded.thumbnail, duration = excluded.duration, accessed_at = excluded.accessed_at`, args: [ fingerprint, video.id, video.title || null, video.channel || null, video.thumbnail || null, video.duration || 0, ], }); // Keep only the 200 most recent entries per user to avoid unbounded growth await db.execute({ sql: `DELETE FROM video_history WHERE fingerprint = ? AND video_id NOT IN ( SELECT video_id FROM video_history WHERE fingerprint = ? ORDER BY accessed_at DESC LIMIT 200 )`, args: [fingerprint, fingerprint], }); } // Fetch stored data for a fingerprint. export async function getUserData(fingerprint) { const [userRow, plRow, histRows] = await Promise.all([ db.execute({ sql: 'SELECT last_version FROM users WHERE fingerprint = ?', args: [fingerprint] }), db.execute({ sql: 'SELECT data FROM playlists WHERE fingerprint = ?', args: [fingerprint] }), db.execute({ sql: `SELECT video_id AS id, title, channel, thumbnail, duration, accessed_at FROM video_history WHERE fingerprint = ? ORDER BY accessed_at DESC LIMIT 50`, args: [fingerprint], }), ]); const lastVersion = userRow.rows[0]?.last_version || null; const playlists = plRow.rows[0]?.data ? JSON.parse(plRow.rows[0].data) : []; const history = histRows.rows.map((r) => ({ id: r.id, title: r.title, channel: r.channel, thumbnail: r.thumbnail, duration: r.duration, })); return { lastVersion, playlists, history }; }