Files
ytplayer/server/p2p-hub.js

167 lines
8.0 KiB
JavaScript
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.

/* ============================================================================
* 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 { 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;
};
}