Files
ytplayer/plans/patches/014-p2p-hub-new.diff

214 lines
10 KiB
Diff
Raw Permalink Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

--- /dev/null
+++ b/server/p2p-hub.js
@@ -0,0 +1,120 @@
+/* ============================================================================
+ * p2p-hub.js — live presence + WebRTC signalling for P2P devices, and the
+ * availability payload (docs/p2p-architecture.md flows 4–5).
+ *
+ * /ws/p2p (credentials never go in the URL — proxies log URLs)
+ * you → server {type:'auth', device, secret} first message, within 10 s
+ * server → you {type:'hello', peer} after auth
+ * you → peer {type:'signal', to:<peer>, data} relayed as
+ * server → peer {type:'signal', from:<your peer>, data}
+ * server → you {type:'error', error}
+ * Server-initiated messages (plan 018) go through hub.send(deviceId, msg).
+ *
+ * "Online" lives only in memory: a server restart forgets it, a closed
+ * socket drops it. Holder ROWS are persistent (p2p_holders) — the payload
+ * reports both: `online` (now) and `lastVerifiedAt` / `stale` (history).
+ * Peers are addressed by an opaque id (peerIdOf), never by device id.
+ * ========================================================================== */
+import { createHash, timingSafeEqual } from 'node:crypto';
+
+const sha = (s) => createHash('sha256').update(String(s)).digest('hex');
+const same = (a, b) => { const x = Buffer.from(String(a)), y = Buffer.from(String(b)); return x.length === y.length && timingSafeEqual(x, y); };
+export const peerIdOf = (deviceId) => sha('peer:' + deviceId).slice(0, 12);
+
+const MAX_MSG = 64 * 1024;
+const BUDGET_PER_MIN = 300;
+
+export function createP2pHub({ getDevice, enabled = () => true, now = () => Date.now(), log = console } = {}) {
+ const online = new Map(); // deviceId -> ws
+ const byPeer = new Map(); // peer -> deviceId (online only)
+
+ const send = (ws, m) => { try { ws.send(JSON.stringify(m)); } catch { /* gone */ } };
+
+ function upgrade(req, server) {
+ if (!enabled()) return new Response('p2p disabled', { status: 404 });
+ const ok = server.upgrade(req, { data: { hub: 'p2p', authed: false, budget: [] } });
+ return ok ? undefined : new Response('upgrade failed', { status: 400 });
+ }
+
+ function open(ws) {
+ const t = setTimeout(() => { if (!ws.data.authed) { try { ws.close(4401, 'auth timeout'); } catch { /* gone */ } } }, 10_000);
+ t.unref?.();
+ }
+
+ async function auth(ws, m) {
+ const deviceId = String(m.device || '');
+ const secret = String(m.secret || '');
+ const d = /^dev_[0-9a-f]{16}$/.test(deviceId) ? await getDevice(deviceId).catch(() => null) : null;
+ if (!d || !same(d.secret_hash, sha(secret))) { try { ws.close(4401, 'unknown device'); } catch { /* gone */ } return; }
+ ws.data.authed = true;
+ ws.data.deviceId = deviceId;
+ ws.data.peer = peerIdOf(deviceId);
+ const prev = online.get(deviceId);
+ if (prev && prev !== ws) { prev.data.deviceId = null; try { prev.close(4409, 'replaced'); } catch { /* gone */ } }
+ online.set(deviceId, ws);
+ byPeer.set(ws.data.peer, deviceId);
+ send(ws, { type: 'hello', peer: ws.data.peer });
+ }
+
+ async function message(ws, raw) {
+ const t = now();
+ ws.data.budget = ws.data.budget.filter((x) => t - x < 60_000);
+ if (ws.data.budget.length >= BUDGET_PER_MIN) { send(ws, { type: 'error', error: 'slow down' }); return; }
+ ws.data.budget.push(t);
+ const s = typeof raw === 'string' ? raw : Buffer.from(raw).toString('utf8');
+ if (s.length > MAX_MSG) return;
+ let m;
+ try { m = JSON.parse(s); } catch { return; }
+ if (!ws.data.authed) { if (m && m.type === 'auth' && !ws.data.authing) { ws.data.authing = true; await auth(ws, m); } return; }
+ if (m && m.type === 'signal' && typeof m.to === 'string') {
+ const target = online.get(byPeer.get(m.to));
+ if (!target || target === ws) { send(ws, { type: 'error', error: 'peer offline', to: m.to }); return; }
+ send(target, { type: 'signal', from: ws.data.peer, data: m.data });
+ }
+ }
+
+ function close(ws) {
+ const { deviceId, peer } = ws.data || {};
+ if (deviceId && online.get(deviceId) === ws) {
+ online.delete(deviceId);
+ byPeer.delete(peer);
+ }
+ }
+
+ return {
+ upgrade,
+ websocket: { open, message, close },
+ isOnline: (deviceId) => online.has(deviceId),
+ send: (deviceId, msg) => { const ws = online.get(deviceId); if (!ws) return false; send(ws, msg); return true; },
+ onlineCount: () => online.size,
+ };
+}
+
+// GET /api/p2p/holders payload. Only devices that share are listed; holders
+// are never dropped for age — `stale` just flags an old last check.
+export async function holdersPayload({ videoId, cid, p2pDb, isOnline, staleDays, serverHas, now = Date.now() }) {
+ const contents = cid
+ ? [await p2pDb.getContent(cid)].filter((c) => c && c.status === 'verified')
+ : await p2pDb.listContentForVideo(videoId);
+ const staleMs = staleDays * 86400_000;
+ const out = [];
+ for (const c of contents) {
+ const holders = (await p2pDb.listHolders(c.cid, 50))
+ .filter((h) => Number(h.share) === 1)
+ .map((h) => ({
+ peer: peerIdOf(h.device_id),
+ online: isOnline(h.device_id),
+ lastVerifiedAt: Number(h.last_verified_at),
+ stale: now - Number(h.last_verified_at) > staleMs,
+ trust: h.trust,
+ }))
+ .sort((a, b) => (b.online - a.online) || (b.lastVerifiedAt - a.lastVerifiedAt));
+ out.push({
+ cid: c.cid, videoId: c.video_id, size: Number(c.size), height: Number(c.height), vcodec: c.vcodec,
+ serverHas: !!(await serverHas(c.cid)),
+ counts: { holders: holders.length, online: holders.filter((h) => h.online).length, fresh: holders.filter((h) => !h.stale).length },
+ holders,
+ });
+ }
+ return { ok: true, staleDays, cids: out };
+}
--- /dev/null
+++ b/server/p2p-hub.test.js
@@ -0,0 +1,87 @@
+// Presence/signalling hub and the holders payload (plan 014).
+import { test, expect, beforeAll } from 'bun:test';
+import { mkdtempSync } from 'node:fs';
+import { tmpdir } from 'node:os';
+import { join } from 'node:path';
+import { createHash } from 'node:crypto';
+
+const root = mkdtempSync(join(tmpdir(), 'ytp-p2p-hub-'));
+process.env.DB_PATH = join(root, 'test.db');
+const dbmod = await import('./db.js');
+const p2pDb = await import('./p2p-db.js');
+const { createP2pHub, holdersPayload, peerIdOf } = await import('./p2p-hub.js');
+
+const sha = (s) => createHash('sha256').update(s).digest('hex');
+const SECRET_A = '1'.repeat(64);
+const SECRET_B = '2'.repeat(64);
+const CID = 'a'.repeat(64);
+const DAY = 86400_000;
+const NOW = Date.UTC(2026, 8, 29, 12);
+
+function fakeWs(data) {
+ return { data, sent: [], closed: null, send(s) { this.sent.push(JSON.parse(s)); }, close(code) { this.closed = code; } };
+}
+
+beforeAll(async () => {
+ await dbmod.initDb();
+ await p2pDb.initP2pSchema();
+ await p2pDb.createDevice({ deviceId: 'dev_000000000000000a', secretHash: sha(SECRET_A), now: NOW });
+ await p2pDb.createDevice({ deviceId: 'dev_000000000000000b', secretHash: sha(SECRET_B), now: NOW });
+ await p2pDb.createDevice({ deviceId: 'dev_000000000000000c', secretHash: sha('x'), now: NOW });
+ await p2pDb.touchDevice('dev_000000000000000c', { now: NOW, share: false });
+ await p2pDb.upsertContent({ cid: CID, videoId: 'dQw4w9WgXcQ', size: 1234, height: 720, vcodec: 'h264', origin: 'server', now: NOW });
+ await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000a', now: NOW - 1 * DAY });
+ await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000b', now: NOW - 30 * DAY }); // stale, still listed
+ await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000c', now: NOW }); // share off → hidden
+});
+
+test('auth on open, hello with opaque peer id, bad secret is closed', async () => {
+ const hub = createP2pHub({ getDevice: p2pDb.getDevice });
+ const a = fakeWs({ budget: [] });
+ hub.websocket.open(a);
+ await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: 'x', data: {} })); // ignored before auth
+ expect(a.sent.length).toBe(0);
+ await hub.websocket.message(a, JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }));
+ expect(a.sent[0]).toEqual({ type: 'hello', peer: peerIdOf('dev_000000000000000a') });
+ expect(hub.isOnline('dev_000000000000000a')).toBe(true);
+ const bad = fakeWs({ budget: [] });
+ hub.websocket.open(bad);
+ await hub.websocket.message(bad, JSON.stringify({ type: 'auth', device: 'dev_000000000000000b', secret: SECRET_A }));
+ expect(bad.closed).toBe(4401);
+ expect(hub.isOnline('dev_000000000000000b')).toBe(false);
+});
+
+test('signals are relayed between online peers only; close drops presence', async () => {
+ const hub = createP2pHub({ getDevice: p2pDb.getDevice });
+ const a = fakeWs({ budget: [] });
+ const b = fakeWs({ budget: [] });
+ hub.websocket.open(a);
+ hub.websocket.open(b);
+ await hub.websocket.message(a, JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }));
+ await hub.websocket.message(b, JSON.stringify({ type: 'auth', device: 'dev_000000000000000b', secret: SECRET_B }));
+ await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: { sdp: 'x' } }));
+ expect(b.sent.at(-1)).toEqual({ type: 'signal', from: a.data.peer, data: { sdp: 'x' } });
+ hub.websocket.close(b);
+ expect(hub.isOnline('dev_000000000000000b')).toBe(false);
+ await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: {} }));
+ expect(a.sent.at(-1)).toMatchObject({ type: 'error', error: 'peer offline' });
+ expect(hub.send('dev_000000000000000a', { type: 'ping' })).toBe(true);
+ expect(hub.send('dev_000000000000000b', { type: 'ping' })).toBe(false);
+});
+
+test('holders payload: persistent rows, online flag, stale marker, share-off hidden', async () => {
+ const online = new Set(['dev_000000000000000b']);
+ const p = await holdersPayload({
+ videoId: 'dQw4w9WgXcQ', p2pDb, isOnline: (d) => online.has(d), staleDays: 7,
+ serverHas: async () => false, now: NOW,
+ });
+ expect(p.staleDays).toBe(7);
+ expect(p.cids.length).toBe(1);
+ const c = p.cids[0];
+ expect(c).toMatchObject({ cid: CID, size: 1234, height: 720, serverHas: false, counts: { holders: 2, online: 1, fresh: 1 } });
+ expect(c.holders.map((h) => [h.peer, h.online, h.stale])).toEqual([
+ [peerIdOf('dev_000000000000000b'), true, true], // online first even though stale
+ [peerIdOf('dev_000000000000000a'), false, false],
+ ]);
+ expect(JSON.stringify(p)).not.toContain('dev_'); // never leak device ids
+});