/* ============================================================================ * p2p-db.js — tables + queries for peer-to-peer sharing * (docs/p2p-architecture.md). Shares the libsql client from db.js. * * p2p_content one row per verified file (cid = sha256 of the bytes); never * deleted, only revoked — the catalog grows over time * p2p_devices registered devices (secret stored as sha256) * p2p_holders which device holds which cid; PERSISTENT (no TTL) with * last_verified_at — the UI decides what is "stale" * video_views per-video per-day view counts (retention criteria) * All timestamps are ms epochs. * ========================================================================== */ import { db } from './db.js'; export async function initP2pSchema() { await db.executeMultiple(` CREATE TABLE IF NOT EXISTS p2p_content ( cid TEXT PRIMARY KEY, video_id TEXT NOT NULL, size INTEGER NOT NULL, height INTEGER NOT NULL DEFAULT 0, vcodec TEXT, acodec TEXT, duration REAL NOT NULL DEFAULT 0, meta TEXT NOT NULL DEFAULT '{}', origin TEXT NOT NULL, -- server | intake status TEXT NOT NULL DEFAULT 'verified', -- verified | revoked scan TEXT NOT NULL DEFAULT 'skipped', -- skipped | clean created_at INTEGER NOT NULL, verified_at INTEGER NOT NULL ); CREATE INDEX IF NOT EXISTS idx_p2p_content_video ON p2p_content (video_id, created_at DESC); CREATE TABLE IF NOT EXISTS p2p_devices ( device_id TEXT PRIMARY KEY, secret_hash TEXT NOT NULL, fingerprint TEXT, profile TEXT, share INTEGER NOT NULL DEFAULT 1, created_at INTEGER NOT NULL, last_seen_at INTEGER NOT NULL ); CREATE TABLE IF NOT EXISTS p2p_holders ( cid TEXT NOT NULL, device_id TEXT NOT NULL, status TEXT NOT NULL DEFAULT 'active', -- active | removed trust TEXT NOT NULL DEFAULT 'reported', -- reported | challenged first_reported_at INTEGER NOT NULL, last_verified_at INTEGER NOT NULL, removed_at INTEGER, PRIMARY KEY (cid, device_id) ); CREATE INDEX IF NOT EXISTS idx_p2p_holders_device ON p2p_holders (device_id, status); CREATE TABLE IF NOT EXISTS video_views ( video_id TEXT NOT NULL, day TEXT NOT NULL, n INTEGER NOT NULL DEFAULT 0, PRIMARY KEY (video_id, day) ); `); // media_cache.sha256 — the cid of the current ..mp4 (plan 009). try { await db.execute('ALTER TABLE media_cache ADD COLUMN sha256 TEXT'); } catch (e) { if (!/duplicate column/i.test(String(e.message))) throw e; } await db.execute('CREATE INDEX IF NOT EXISTS idx_media_sha256 ON media_cache (sha256)'); } const rowsOf = (r) => r.rows.map((row) => { const o = {}; r.columns.forEach((c, i) => { o[c] = row[i]; }); return o; }); // ---- content ---------------------------------------------------------------- export async function upsertContent(c) { await db.execute({ sql: `INSERT INTO p2p_content (cid, video_id, size, height, vcodec, acodec, duration, meta, origin, status, scan, created_at, verified_at) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, 'verified', ?, ?, ?) ON CONFLICT(cid) DO UPDATE SET verified_at = excluded.verified_at, scan = excluded.scan`, args: [c.cid, c.videoId, c.size, c.height || 0, c.vcodec || null, c.acodec || null, c.duration || 0, JSON.stringify(c.meta || {}), c.origin, c.scan || 'skipped', c.now, c.now], }); } export async function getContent(cid) { const r = await db.execute({ sql: 'SELECT * FROM p2p_content WHERE cid = ?', args: [cid] }); return rowsOf(r)[0] || null; } export async function listContentForVideo(videoId) { const r = await db.execute({ sql: "SELECT * FROM p2p_content WHERE video_id = ? AND status = 'verified' ORDER BY created_at DESC LIMIT 10", args: [videoId], }); return rowsOf(r); } export async function knownCids(cids) { if (!cids.length) return new Set(); const out = new Set(); for (let i = 0; i < cids.length; i += 200) { const part = cids.slice(i, i + 200); const r = await db.execute({ sql: `SELECT cid FROM p2p_content WHERE status = 'verified' AND cid IN (${part.map(() => '?').join(',')})`, args: part, }); for (const row of r.rows) out.add(row[0]); } return out; } export async function revokeContent(cid) { await db.execute({ sql: "UPDATE p2p_content SET status = 'revoked' WHERE cid = ?", args: [cid] }); } // ---- devices ------------------------------------------------------------------ export async function createDevice({ deviceId, secretHash, fingerprint, profile, now }) { await db.execute({ sql: `INSERT INTO p2p_devices (device_id, secret_hash, fingerprint, profile, share, created_at, last_seen_at) VALUES (?, ?, ?, ?, 1, ?, ?)`, args: [deviceId, secretHash, fingerprint || null, profile || null, now, now], }); } export async function getDevice(deviceId) { const r = await db.execute({ sql: 'SELECT * FROM p2p_devices WHERE device_id = ?', args: [deviceId] }); return rowsOf(r)[0] || null; } export async function touchDevice(deviceId, { now, share, profile } = {}) { await db.execute({ sql: `UPDATE p2p_devices SET last_seen_at = ?, share = COALESCE(?, share), profile = COALESCE(?, profile) WHERE device_id = ?`, args: [now, share === undefined ? null : (share ? 1 : 0), profile || null, deviceId], }); } // ---- holders (persistent; never expired by time) ------------------------------ export async function upsertHolder({ cid, deviceId, trust = 'reported', now }) { await db.execute({ sql: `INSERT INTO p2p_holders (cid, device_id, status, trust, first_reported_at, last_verified_at) VALUES (?, ?, 'active', ?, ?, ?) ON CONFLICT(cid, device_id) DO UPDATE SET status = 'active', removed_at = NULL, last_verified_at = excluded.last_verified_at, trust = CASE WHEN p2p_holders.trust = 'challenged' OR excluded.trust = 'challenged' THEN 'challenged' ELSE 'reported' END`, args: [cid, deviceId, trust, now, now], }); } export async function setHolderTrust({ cid, deviceId, trust, now }) { await db.execute({ sql: 'UPDATE p2p_holders SET trust = ?, last_verified_at = ? WHERE cid = ? AND device_id = ?', args: [trust, now, cid, deviceId], }); } export async function removeHolder({ cid, deviceId, now }) { await db.execute({ sql: "UPDATE p2p_holders SET status = 'removed', removed_at = ? WHERE cid = ? AND device_id = ? AND status = 'active'", args: [now, cid, deviceId], }); } // A full report: every active holding of this device NOT in `keep` is removed. export async function removeHoldersExcept({ deviceId, keep, now }) { const r = await db.execute({ sql: "SELECT cid FROM p2p_holders WHERE device_id = ? AND status = 'active'", args: [deviceId], }); const keepSet = new Set(keep); let removed = 0; for (const row of r.rows) { if (keepSet.has(row[0])) continue; await removeHolder({ cid: row[0], deviceId, now }); removed++; } return removed; } export async function activeHoldingsOf(deviceId) { const r = await db.execute({ sql: "SELECT cid FROM p2p_holders WHERE device_id = ? AND status = 'active'", args: [deviceId], }); return r.rows.map((row) => row[0]); } // Holders of one cid, joined with the device's share flag. Newest check first. export async function listHolders(cid, limit = 50) { const r = await db.execute({ sql: `SELECT h.device_id, h.trust, h.first_reported_at, h.last_verified_at, d.share FROM p2p_holders h JOIN p2p_devices d ON d.device_id = h.device_id WHERE h.cid = ? AND h.status = 'active' ORDER BY h.last_verified_at DESC LIMIT ?`, args: [cid, limit], }); return rowsOf(r); } // ---- views + stats -------------------------------------------------------------- export function dayKey(ms) { return new Date(ms).toISOString().slice(0, 10); } export async function addView(videoId, now) { await db.execute({ sql: `INSERT INTO video_views (video_id, day, n) VALUES (?, ?, 1) ON CONFLICT(video_id, day) DO UPDATE SET n = n + 1`, args: [videoId, dayKey(now)], }); } export async function viewsSince(videoId, sinceMs) { const r = await db.execute({ sql: 'SELECT COALESCE(SUM(n), 0) FROM video_views WHERE video_id = ? AND day >= ?', args: [videoId, dayKey(sinceMs)], }); return Number(r.rows[0][0]) || 0; } export async function p2pStats() { const one = async (sql) => Number((await db.execute(sql)).rows[0][0]) || 0; return { content: await one("SELECT COUNT(*) FROM p2p_content WHERE status = 'verified'"), revoked: await one("SELECT COUNT(*) FROM p2p_content WHERE status = 'revoked'"), devices: await one('SELECT COUNT(*) FROM p2p_devices'), holders: await one("SELECT COUNT(*) FROM p2p_holders WHERE status = 'active'"), heldCids: await one("SELECT COUNT(DISTINCT cid) FROM p2p_holders WHERE status = 'active'"), }; } // ---- retention (plan 010) ------------------------------------------------------- // Server copies in the order they should be evicted when the cache needs room: // first the ones that are neither "top" (≥ keepMinViews views in keepDays) nor // "recent" (played within keepRecentDays) — fewest views, then oldest — and // only then the qualifying ones, least recently played first. export async function listMediaEvictionOrder({ now, keepMinViews, keepDays, keepRecentDays }) { const r = await db.execute({ sql: `SELECT m.video_id, m.size, m.last_access, COALESCE((SELECT SUM(v.n) FROM video_views v WHERE v.video_id = m.video_id AND v.day >= ?), 0) AS views FROM media_cache m WHERE m.status = 'ready'`, args: [dayKey(now - keepDays * 86400_000)], }); const recentCut = now - keepRecentDays * 86400_000; const rows = rowsOf(r).map((x) => ({ ...x, views: Number(x.views) || 0, qualifies: (Number(x.views) || 0) >= keepMinViews || Number(x.last_access) >= recentCut, })); rows.sort((a, b) => (a.qualifies - b.qualifies) || (a.qualifies ? 0 : a.views - b.views) || (a.last_access - b.last_access)); return rows; } // The server's own ready copy whose mp4 has this content id (or null). export async function findMediaByCid(cid) { const r = await db.execute({ sql: "SELECT video_id, gen FROM media_cache WHERE sha256 = ? AND status = 'ready' LIMIT 1", args: [cid], }); return rowsOf(r)[0] || null; } // Admin panel (plan 019): newest verified/revoked files with their holder counts. export async function recentContent(limit = 30) { const r = await db.execute({ sql: `SELECT c.cid, c.video_id, c.size, c.height, c.vcodec, c.origin, c.status, c.scan, c.verified_at, c.meta, (SELECT COUNT(*) FROM p2p_holders h WHERE h.cid = c.cid AND h.status = 'active') AS holders FROM p2p_content c ORDER BY c.verified_at DESC LIMIT ?`, args: [limit], }); return rowsOf(r).map((x) => ({ ...x, holders: Number(x.holders) || 0 })); } // One verified file per video (newest), with how many sharing devices hold it — // the "On other devices" playlist. Holders themselves come from listHolders. export async function listSharedContent(limit = 100) { const r = await db.execute({ sql: `SELECT c.cid, c.video_id, c.size, c.height, c.vcodec, c.duration, c.meta, c.verified_at, (SELECT COUNT(*) FROM p2p_holders h JOIN p2p_devices d ON d.device_id = h.device_id WHERE h.cid = c.cid AND h.status = 'active' AND d.share = 1) AS holders FROM p2p_content c WHERE c.status = 'verified' AND c.created_at = (SELECT MAX(c2.created_at) FROM p2p_content c2 WHERE c2.video_id = c.video_id AND c2.status = 'verified') ORDER BY c.verified_at DESC LIMIT ?`, args: [limit], }); return rowsOf(r).map((x) => ({ ...x, holders: Number(x.holders) || 0 })); }