146 lines
6.6 KiB
JavaScript
146 lines
6.6 KiB
JavaScript
/* ============================================================================
|
||
* 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 };
|
||
}
|
||
|
||
// 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;
|
||
};
|
||
}
|