/* ============================================================================ * 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:, data} relayed as * server → peer {type:'signal', from:, 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 { createDirectRelay } from './direct-relay.js'; 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, authorizeProfile = async () => '', enabled = () => true, now = () => Date.now(), log = console } = {}) { const directRooms = new Map(); 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.profile = await authorizeProfile(m.profile, m.profileSecret, deviceId).catch(() => ''); ws.data.inventory = []; 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 }); directory(ws.data.profile); } function directory(profile) { if (!profile) return; const sockets = [...online.values()].filter(ws => ws.data.profile === profile); const peers = sockets.map(ws => ({ id: ws.data.peer, name: 'Device ' + ws.data.peer.slice(-4), files: ws.data.inventory || [] })); for (const ws of sockets) send(ws, { type: 'direct-peers', peers: peers.filter(p => p.id !== ws.data.peer) }); } async function message(ws, raw) { if (typeof raw !== 'string') return; 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?.type === 'direct-inventory' && ws.data.profile) { ws.data.inventory = Array.isArray(m.ids) ? [...new Set(m.ids.filter(id => typeof id === 'string' && /^[A-Za-z0-9_-]{1,128}$/.test(id)))].slice(0, 1000) : []; directory(ws.data.profile); return; } if (m?.type === 'direct' && ws.data.profile) { const profile = ws.data.profile; const peers = new Map([...online.values()].filter(x => x.data.profile === profile).map(x => [x.data.peer, x])); if (!directRooms.has(profile)) directRooms.set(profile, createDirectRelay()); directRooms.get(profile).handle(ws.data.peer, m, peers, send); return; } if (m?.type === 'signal') send(ws, { type: 'error', error: 'use a confirmed direct-transfer invitation', to: m.to }); } function close(ws) { const { deviceId, peer } = ws.data || {}; if (deviceId && online.get(deviceId) === ws) { online.delete(deviceId); byPeer.delete(peer); directory(ws.data.profile); if (![...online.values()].some(x => x.data.profile === ws.data.profile)) directRooms.delete(ws.data.profile); } } 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 }; } // Flow 8 (plan 018): the source is gone and the server evicted its copy, so // ask ONE online holder that shares to upload it through intake // ({type:'upload-request', videoId, cid}). At most one request per video per // 10 minutes; returns true when a device was asked (or recently was). export function createRehydrator({ p2pDb, hub, hasServerCopy, enabled = () => true, now = () => Date.now() }) { const asked = new Map(); // videoId -> ms return async function rehydrate(videoId) { if (!enabled()) return false; if (await hasServerCopy(videoId)) return false; const last = asked.get(videoId); if (last && now() - last < 10 * 60_000) return true; for (const c of await p2pDb.listContentForVideo(videoId)) { for (const h of await p2pDb.listHolders(c.cid, 50)) { if (Number(h.share) !== 1 || !hub.isOnline(h.device_id)) continue; if (hub.send(h.device_id, { type: 'upload-request', videoId, cid: c.cid })) { asked.set(videoId, now()); if (asked.size > 5000) asked.delete(asked.keys().next().value); return true; } } } return false; }; }