diff --git a/docker-compose.yml b/docker-compose.yml index efaba12..4c50cc2 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -18,6 +18,15 @@ services: # Shared secret the lyrics-worker uses to list songs and upload lyrics. # Set in Dokploy's Environment tab (same value feeds both services). LYRICS_WORKER_TOKEN: "${LYRICS_WORKER_TOKEN:-}" + # Peer-to-peer sharing (docs/p2p-architecture.md). ON by default. + P2P_ENABLED: "${P2P_ENABLED:-1}" + # Malware scan before a file's hash is admitted. OFF by default; needs an + # image built with INSTALL_CLAMAV=1. Hashing + media validation always run. + P2P_MALWARE_SCAN: "${P2P_MALWARE_SCAN:-0}" + # P2P_STALE_DAYS: "7" # holder shown as stale after this many days unchecked + # P2P_KEEP_MIN_VIEWS: "3" # server keeps copies with ≥ this many views… + # P2P_KEEP_DAYS: "30" # …in this many days + # P2P_KEEP_RECENT_DAYS: "14" # …or played this recently # Optional: force yt-dlp search instead of InnerTube API # SEARCH_INNERTUBE: "0" # force yt-dlp search # Optional: override yt-dlp binary path if you mount a custom one diff --git a/plans/INDEX.md b/plans/INDEX.md index c17d569..48eb36a 100644 --- a/plans/INDEX.md +++ b/plans/INDEX.md @@ -15,7 +15,7 @@ green, app boots with no JS errors, P2P on by default, offline boot works). | 005 | 005-warm-streams-on-intent-47b3d3 | Warm the stream cache for likely next plays | done | Warm the stream cache for likely next plays | cold play 7.5 s → cached 1.2 s | | 006 | 006-innertube-search-48066b | Answer searches from YouTube InnerTube directly with yt-dlp fallback | done | Answer searches from YouTube InnerTube directly with yt-dlp fallback | search 4–5 s → ~0.8 s | | 007 | 007-ytdlp-worker-045800 | Keep one long-lived yt-dlp worker process instead of spawning per call | done | Keep one long-lived yt-dlp worker process instead of spawning per call | ~1 s per yt-dlp call | -| 008 | 008-p2p-schema-and-config-127966 | Add P2P tables, config flags and db helpers | in-progress | | P2P ON, malware scan OFF by default | +| 008 | 008-p2p-schema-and-config-127966 | Add P2P tables, config flags and db helpers | done | Add P2P tables, config flags and db helpers | P2P ON, malware scan OFF by default | | 009 | 009-server-content-hash-186e7f | Hash every validated server copy and register it as verified content | queued | | uses plans/patches/009-* | | 010 | 010-views-and-retention-d0c6ca | Count views and evict server copies by retention criteria before LRU | queued | | | | 011 | 011-browser-sha256-e1793d | Add an incremental SHA-256 library for the browser and node tests | queued | | | diff --git a/plans/active/008-p2p-schema-and-config-127966.md b/plans/done/008-p2p-schema-and-config-127966.md similarity index 97% rename from plans/active/008-p2p-schema-and-config-127966.md rename to plans/done/008-p2p-schema-and-config-127966.md index a0fab7d..1ace206 100644 --- a/plans/active/008-p2p-schema-and-config-127966.md +++ b/plans/done/008-p2p-schema-and-config-127966.md @@ -425,3 +425,9 @@ test('views per day and window sums', async () => { expect(s.devices).toBe(1); }); ``` + +## Execution log + +- Executor: in-session Agent (haiku). Attempts: 1. Fix rounds: 0. +- Orchestrator re-ran Verification: `p2p-db.js`, `p2p-db.test.js`, `p2p-config.js` byte-identical to the plan; p2p-db tests 4 pass; all 8 server test files 0 fail; `SERVER_OK`; server boots and `/api/version` returns 200. +- Executor Findings (verbatim): All steps executed successfully. P2P schema, config, and tests created verbatim from plan appendices. Server builds without errors. No deviations from plan requirements. diff --git a/server/db.js b/server/db.js index a4cbae1..4a62941 100644 --- a/server/db.js +++ b/server/db.js @@ -308,7 +308,7 @@ export async function useApiToken(tokenHash) { 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', + 'updated_at', 'last_access', 'hits', 'sha256', ]); function rowToObj(res, row) { diff --git a/server/p2p-config.js b/server/p2p-config.js new file mode 100644 index 0000000..753f9da --- /dev/null +++ b/server/p2p-config.js @@ -0,0 +1,23 @@ +/* p2p-config.js — peer-to-peer settings (see docs/p2p-architecture.md). + * P2P is ON unless P2P_ENABLED=0. The malware scan is OFF unless + * P2P_MALWARE_SCAN=1. Hashing + validateMedia are never optional. */ +import { dirname, join } from 'node:path'; + +const num = (v, d) => (Number.isFinite(Number(v)) && String(v).trim() !== '' ? Number(v) : d); + +export function loadP2pConfig(env = process.env) { + const dbDir = dirname(env.DB_PATH || './data/ytplayer.db'); + return { + enabled: env.P2P_ENABLED !== '0', + malwareScan: env.P2P_MALWARE_SCAN === '1', + scanCmd: (env.P2P_SCAN_CMD || 'clamscan --no-summary --infected').trim(), + staleDays: num(env.P2P_STALE_DAYS, 7), + keepMinViews: num(env.P2P_KEEP_MIN_VIEWS, 3), + keepDays: num(env.P2P_KEEP_DAYS, 30), + keepRecentDays: num(env.P2P_KEEP_RECENT_DAYS, 14), + intakeDir: env.P2P_INTAKE_DIR || join(dbDir, 'p2p-intake'), + intakeMaxBytes: num(env.P2P_INTAKE_MAX_BYTES, 3 * 1024 ** 3), + }; +} + +export const P2P = loadP2pConfig(); diff --git a/server/p2p-db.js b/server/p2p-db.js new file mode 100644 index 0000000..b9b6738 --- /dev/null +++ b/server/p2p-db.js @@ -0,0 +1,236 @@ +/* ============================================================================ + * 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; } +} + +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'"), + }; +} diff --git a/server/p2p-db.test.js b/server/p2p-db.test.js new file mode 100644 index 0000000..01f4c8d --- /dev/null +++ b/server/p2p-db.test.js @@ -0,0 +1,71 @@ +// P2P tables against a real temp libsql DB (docs/p2p-architecture.md). +import { test, expect, beforeAll } from 'bun:test'; +import { mkdtempSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; + +const root = mkdtempSync(join(tmpdir(), 'ytp-p2p-test-')); +process.env.DB_PATH = join(root, 'test.db'); +const dbmod = await import('./db.js'); +const P = await import('./p2p-db.js'); +const { loadP2pConfig } = await import('./p2p-config.js'); + +const CID = 'a'.repeat(64); +const CID2 = 'b'.repeat(64); +const T0 = Date.UTC(2026, 8, 29, 12); + +beforeAll(async () => { + await dbmod.initDb(); + await P.initP2pSchema(); + await P.initP2pSchema(); // idempotent (ALTER TABLE duplicate column is ignored) +}); + +test('config defaults: P2P on, malware scan off', () => { + const c = loadP2pConfig({}); + expect(c.enabled).toBe(true); + expect(c.malwareScan).toBe(false); + expect(c.staleDays).toBe(7); + expect(loadP2pConfig({ P2P_ENABLED: '0', P2P_MALWARE_SCAN: '1' })).toMatchObject({ enabled: false, malwareScan: true }); +}); + +test('content upsert/get/known/revoke', async () => { + await P.upsertContent({ cid: CID, videoId: 'dQw4w9WgXcQ', size: 1000, height: 720, vcodec: 'h264', acodec: 'aac', duration: 212, meta: { title: 'x' }, origin: 'server', now: T0 }); + const c = await P.getContent(CID); + expect(c.video_id).toBe('dQw4w9WgXcQ'); + expect(c.status).toBe('verified'); + expect([...(await P.knownCids([CID, CID2]))]).toEqual([CID]); + expect((await P.listContentForVideo('dQw4w9WgXcQ')).length).toBe(1); + await P.upsertContent({ cid: CID2, videoId: 'dQw4w9WgXcQ', size: 5, origin: 'intake', now: T0 }); + await P.revokeContent(CID2); + expect([...(await P.knownCids([CID2]))]).toEqual([]); +}); + +test('holders persist, never expire, and a full report removes missing ones', async () => { + await P.createDevice({ deviceId: 'dev_1', secretHash: 'h', fingerprint: 'fp', now: T0 }); + await P.upsertHolder({ cid: CID, deviceId: 'dev_1', now: T0 }); + // 90 days later with no new report: still listed (UI marks it stale). + let hs = await P.listHolders(CID); + expect(hs.length).toBe(1); + expect(hs[0].last_verified_at).toBe(T0); + await P.setHolderTrust({ cid: CID, deviceId: 'dev_1', trust: 'challenged', now: T0 + 1000 }); + await P.upsertHolder({ cid: CID, deviceId: 'dev_1', trust: 'reported', now: T0 + 2000 }); + hs = await P.listHolders(CID); + expect(hs[0].trust).toBe('challenged'); // a later plain report never downgrades trust + expect(hs[0].last_verified_at).toBe(T0 + 2000); + expect(await P.removeHoldersExcept({ deviceId: 'dev_1', keep: [], now: T0 + 3000 })).toBe(1); + expect((await P.listHolders(CID)).length).toBe(0); + expect(await P.activeHoldingsOf('dev_1')).toEqual([]); + await P.upsertHolder({ cid: CID, deviceId: 'dev_1', now: T0 + 4000 }); // re-added + expect(await P.activeHoldingsOf('dev_1')).toEqual([CID]); +}); + +test('views per day and window sums', async () => { + await P.addView('vid00000001', T0); + await P.addView('vid00000001', T0); + await P.addView('vid00000001', T0 - 40 * 86400_000); + expect(await P.viewsSince('vid00000001', T0 - 30 * 86400_000)).toBe(2); + expect(await P.viewsSince('vid00000001', T0 - 50 * 86400_000)).toBe(3); + const s = await P.p2pStats(); + expect(s.content).toBe(1); + expect(s.devices).toBe(1); +}); diff --git a/server/package.json b/server/package.json index cf198f4..0671a1d 100644 --- a/server/package.json +++ b/server/package.json @@ -6,7 +6,7 @@ "scripts": { "start": "bun server.js", "dev": "bun --hot server.js", - "test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js" + "test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js" }, "dependencies": { "@hono/node-server": "^1.14.0", diff --git a/server/server.js b/server/server.js index b7dbb06..266aa7f 100644 --- a/server/server.js +++ b/server/server.js @@ -48,6 +48,7 @@ import { registerNoteRoutes, parseLrc, sanitizeLyrics } from './notes.js'; import { createRemoteHub } from './remote.js'; import { createPartyHub } from './party.js'; import { registerUploadRoutes } from './uploads.js'; +import { initP2pSchema } from './p2p-db.js'; import * as innertube from './innertube.js'; import QRCode from 'qrcode'; import { createYtdlpPool } from './ytdlp-pool.js'; @@ -2021,6 +2022,7 @@ app.get('/*', indexHtml); // ============================================================================ async function main() { await initDb(); + await initP2pSchema(); console.log(`[ytplayer] DB ready`); await media.init(); notes.startBackups();