Files
ytplayer/server/db.js

391 lines
14 KiB
JavaScript

/* ============================================================================
* 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 <video_id>.<gen>.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);
`);
}
// ---- 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 = 1 THEN 1 ELSE 0 END) AS optimized
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,
};
}
// ---- 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 };
}