Files
ytplayer/server/db.js

769 lines
29 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);
-- Persistent search cache: the YouTube part of a search (up to 200 cards,
-- JSON), keyed by the lowercased query. Trimmed LRU by count and bytes.
CREATE TABLE IF NOT EXISTS search_cache (
q TEXT PRIMARY KEY,
results TEXT NOT NULL,
size INTEGER NOT NULL DEFAULT 0,
fetched_at INTEGER NOT NULL DEFAULT 0,
last_access INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX IF NOT EXISTS idx_search_cache_lru ON search_cache (last_access);
-- Every video the server has ever seen in a search: lets a search show
-- known videos at once while YouTube is still loading.
CREATE TABLE IF NOT EXISTS video_meta (
id TEXT PRIMARY KEY,
card TEXT NOT NULL, -- the result card, JSON
hay TEXT NOT NULL, -- lowercased title + channel, for matching
seen INTEGER NOT NULL DEFAULT 1, -- times it appeared in a search
updated_at INTEGER NOT NULL DEFAULT 0
);
CREATE INDEX IF NOT EXISTS idx_video_meta_upd ON video_meta (updated_at);
-- 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);
-- Admin uploads: the server's own video/audio library, searched next to
-- YouTube. Files live in UPLOAD_DIR as <id>.<ext> (+ <id>.art.jpg).
CREATE TABLE IF NOT EXISTS uploads (
id TEXT PRIMARY KEY, -- upl_<hex>
kind TEXT NOT NULL, -- video | audio
title TEXT NOT NULL,
artist TEXT,
album TEXT,
duration REAL NOT NULL DEFAULT 0,
ext TEXT NOT NULL,
mime TEXT NOT NULL,
size INTEGER NOT NULL DEFAULT 0,
art TEXT, -- 'embedded' | 'file' | NULL
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
plays INTEGER NOT NULL DEFAULT 0
);
-- "This lyric line is wrong" reports from service mode. One open row per
-- (video, line text, reporter); the admin lyric editor lists them and a
-- lyric edit that changes the line resolves the report automatically.
CREATE TABLE IF NOT EXISTS lyric_flags (
id INTEGER PRIMARY KEY AUTOINCREMENT,
video_id TEXT NOT NULL,
line_index INTEGER NOT NULL DEFAULT -1,
line_text TEXT NOT NULL,
rev INTEGER NOT NULL DEFAULT 0, -- lyrics rev the reporter saw
reason TEXT NOT NULL DEFAULT '', -- words | timing | typo | other | ''
note TEXT NOT NULL DEFAULT '', -- optional description
reporter TEXT NOT NULL, -- profile name or api label
status TEXT NOT NULL DEFAULT 'open', -- open | resolved
resolved_by TEXT,
resolved_at INTEGER,
created_at INTEGER NOT NULL DEFAULT (unixepoch()),
updated_at INTEGER NOT NULL DEFAULT (unixepoch())
);
CREATE INDEX IF NOT EXISTS idx_flags_video ON lyric_flags (video_id, status);
CREATE INDEX IF NOT EXISTS idx_flags_status ON lyric_flags (status, 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);
}
// ---- Lyric line reports ("this line is wrong") --------------------------------
function flagRow(x) {
return {
id: Number(x.id), videoId: x.video_id, index: Number(x.line_index), text: x.line_text,
rev: Number(x.rev), reason: x.reason || '', note: x.note || '', reporter: x.reporter,
status: x.status, resolvedBy: x.resolved_by || null,
resolvedAt: x.resolved_at == null ? null : Number(x.resolved_at),
createdAt: Number(x.created_at), updatedAt: Number(x.updated_at),
};
}
// One open report per (video, line text, reporter): reporting the same line
// again refreshes the reason/description instead of piling up duplicates.
export async function upsertFlag({ videoId, index, text, rev, reason, note, reporter }) {
const ex = await db.execute({
sql: "SELECT id FROM lyric_flags WHERE video_id = ? AND line_text = ? AND reporter = ? AND status = 'open'",
args: [videoId, text, reporter],
});
if (ex.rows[0]) {
const id = Number(ex.rows[0].id);
await db.execute({
sql: 'UPDATE lyric_flags SET line_index = ?, rev = ?, reason = ?, note = ?, updated_at = unixepoch() WHERE id = ?',
args: [index, rev, reason, note, id],
});
return { id, created: false };
}
const r = await db.execute({
sql: `INSERT INTO lyric_flags (video_id, line_index, line_text, rev, reason, note, reporter)
VALUES (?, ?, ?, ?, ?, ?, ?)`,
args: [videoId, index, text, rev, reason, note, reporter],
});
return { id: Number(r.lastInsertRowid), created: true };
}
export async function listFlagsForVideo(videoId, status = 'open') {
const r = await db.execute({
sql: 'SELECT * FROM lyric_flags WHERE video_id = ? AND status = ? ORDER BY line_index, id',
args: [videoId, status],
});
return r.rows.map(flagRow);
}
// Newest first across every video; status 'all' returns both.
export async function listAllFlags({ status = 'open', limit = 200 } = {}) {
const where = status === 'all' ? '' : 'WHERE status = ?';
const r = await db.execute({
sql: `SELECT * FROM lyric_flags ${where} ORDER BY created_at DESC, id DESC LIMIT ?`,
args: status === 'all' ? [limit] : [status, limit],
});
return r.rows.map(flagRow);
}
export async function getFlag(id) {
const r = await db.execute({ sql: 'SELECT * FROM lyric_flags WHERE id = ?', args: [id] });
return r.rows[0] ? flagRow(r.rows[0]) : null;
}
export async function countOpenFlags() {
const r = await db.execute("SELECT COUNT(*) AS n FROM lyric_flags WHERE status = 'open'");
return Number(r.rows[0].n) || 0;
}
export async function setFlagStatus(id, status, by) {
const r = await db.execute({
sql: `UPDATE lyric_flags SET status = ?, resolved_by = ?, resolved_at = ?, updated_at = unixepoch() WHERE id = ?`,
args: [status, status === 'resolved' ? by : null, status === 'resolved' ? Math.floor(Date.now() / 1000) : null, id],
});
return (r.rowsAffected || 0) > 0;
}
export async function deleteFlag(id) {
const r = await db.execute({ sql: 'DELETE FROM lyric_flags WHERE id = ?', args: [id] });
return (r.rowsAffected || 0) > 0;
}
// A reporter withdrawing their own open report.
export async function withdrawFlag(id, reporter) {
const r = await db.execute({
sql: "DELETE FROM lyric_flags WHERE id = ? AND reporter = ? AND status = 'open'",
args: [id, reporter],
});
return (r.rowsAffected || 0) > 0;
}
// After a lyric edit: every open report whose line no longer exists (its text
// changed or was removed) is resolved automatically. `texts` = the new lines.
export async function autoResolveFlags(videoId, texts, by = 'auto (line changed)') {
const open = await listFlagsForVideo(videoId, 'open');
let n = 0;
for (const f of open) {
if (texts.has(f.text)) continue;
if (await setFlagStatus(f.id, 'resolved', by)) n++;
}
return n;
}
// ---- 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', 'sha256',
]);
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 };
}
// ---- Uploads (the server's own media library) ---------------------------------
const uploadRow = (r) => ({
id: r.id, kind: r.kind, title: r.title, artist: r.artist || '', album: r.album || '',
duration: Number(r.duration) || 0, ext: r.ext, mime: r.mime, size: Number(r.size) || 0,
art: r.art || null, createdAt: Number(r.created_at), plays: Number(r.plays) || 0,
});
export async function createUpload(u) {
await db.execute({
sql: `INSERT INTO uploads (id, kind, title, artist, album, duration, ext, mime, size, art, created_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, unixepoch())`,
args: [u.id, u.kind, u.title, u.artist || null, u.album || null, u.duration || 0, u.ext, u.mime, u.size || 0, u.art || null],
});
}
export async function setUploadArt(id, art) {
await db.execute({ sql: 'UPDATE uploads SET art = ? WHERE id = ?', args: [art, id] });
}
export async function getUpload(id) {
const r = await db.execute({ sql: 'SELECT * FROM uploads WHERE id = ?', args: [id] });
return r.rows[0] ? uploadRow(r.rows[0]) : null;
}
// Newest first; `q` matches title/artist/album (case-insensitive).
export async function listUploads({ q = '', limit = 100 } = {}) {
const like = `%${String(q).toLowerCase()}%`;
const r = q
? await db.execute({
sql: `SELECT * FROM uploads
WHERE lower(title) LIKE ? OR lower(ifnull(artist, '')) LIKE ? OR lower(ifnull(album, '')) LIKE ?
ORDER BY created_at DESC LIMIT ?`,
args: [like, like, like, limit],
})
: await db.execute({ sql: 'SELECT * FROM uploads ORDER BY created_at DESC LIMIT ?', args: [limit] });
return r.rows.map(uploadRow);
}
export async function deleteUpload(id) {
const r = await db.execute({ sql: 'DELETE FROM uploads WHERE id = ?', args: [id] });
return (r.rowsAffected || 0) > 0;
}
export async function touchUpload(id) {
db.execute({ sql: 'UPDATE uploads SET plays = plays + 1 WHERE id = ?', args: [id] }).catch(() => {});
}