214 lines
10 KiB
Diff
214 lines
10 KiB
Diff
--- /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
|
||
+});
|