/* ============================================================================ * p2p-transfer.js — download a verified file from another device over a * WebRTC data channel (docs/p2p-architecture.md flow 6). window.P2PTransfer. * * start({ canShare }) serve requests from other devices (called once) * download({ videoId, cid, size, peers, onProgress }) → { ok, error? } * * Signalling rides P2PClient.signal()/onMessage('signal') over /ws/p2p: * {k:'offer', xid, cid, sdp} · {k:'answer', xid, sdp} · {k:'ice', xid, cand} * {k:'deny', xid, reason} * Data channel "file" (ordered): requester sends {t:'get'}; sender answers * {t:'meta', size}, 64 KiB binary frames (back-pressured on bufferedAmount), * then {t:'end'}. The receiver writes through p2p-recv-worker.js, which only * commits the file when its SHA-256 equals the cid. STUN only — devices * behind strict NATs won't connect (same as watch-party voice); the next * holder is tried. * ========================================================================== */ (function () { 'use strict'; const ICE = [{ urls: 'stun:stun.l.google.com:19302' }, { urls: 'stun:stun1.l.google.com:19302' }]; const FRAME = 64 * 1024; const HIGH_WATER = 4 * 1024 * 1024; const OPEN_TIMEOUT = 20_000; const IDLE_TIMEOUT = 30_000; const MAX_SERVING = 1; let iceServers = ICE; let canShare = () => true; const sessions = new Map(); // xid -> { pc, onSignal } let serving = 0; const newXid = () => Math.random().toString(36).slice(2) + Date.now().toString(36); function onSignal(m) { const d = m && m.data; if (!d || !d.xid) return; const s = sessions.get(d.xid); if (s) { s.onSignal(d, m.from); return; } if (d.k === 'offer') serve(d, m.from).catch(() => {}); } // ---- sender ------------------------------------------------------------------- async function serve(offer, from) { const deny = (reason) => window.P2PClient.signal(from, { k: 'deny', xid: offer.xid, reason }); if (!canShare()) return deny('sharing off'); if (serving >= MAX_SERVING) return deny('busy'); const rec = await window.DeviceDB.getByCid(offer.cid); const file = rec && window.OPFS.getFileObject ? await window.OPFS.getFileObject(rec.videoId) : null; if (!rec || !file || file.size !== rec.size) return deny('not here'); serving++; const pc = new RTCPeerConnection({ iceServers }); let done = false; const finish = () => { if (done) return; done = true; serving--; sessions.delete(offer.xid); try { pc.close(); } catch { /* closed */ } }; sessions.set(offer.xid, { pc, onSignal: (d) => { if (d.k === 'ice' && d.cand) pc.addIceCandidate(d.cand).catch(() => {}); }, }); pc.onicecandidate = (e) => { if (e.candidate) window.P2PClient.signal(from, { k: 'ice', xid: offer.xid, cand: e.candidate.toJSON() }); }; pc.onconnectionstatechange = () => { if (['failed', 'closed', 'disconnected'].includes(pc.connectionState)) finish(); }; setTimeout(() => { if (pc.connectionState !== 'connected') finish(); }, OPEN_TIMEOUT); pc.ondatachannel = (e) => { const dc = e.channel; dc.bufferedAmountLowThreshold = 1024 * 1024; dc.onmessage = async (ev) => { let msg; try { msg = JSON.parse(ev.data); } catch { return; } if (msg.t !== 'get') return; try { dc.send(JSON.stringify({ t: 'meta', size: file.size })); for (let pos = 0; pos < file.size && !done; pos += 4 * FRAME) { const buf = await file.slice(pos, pos + 4 * FRAME).arrayBuffer(); for (let i = 0; i < buf.byteLength; i += FRAME) { if (dc.bufferedAmount > HIGH_WATER) { await new Promise((r) => { dc.onbufferedamountlow = () => { dc.onbufferedamountlow = null; r(); }; }); } dc.send(buf.slice(i, i + FRAME)); } } dc.send(JSON.stringify({ t: 'end' })); } catch { finish(); } }; dc.onclose = finish; }; await pc.setRemoteDescription({ type: 'offer', sdp: offer.sdp }); const answer = await pc.createAnswer(); await pc.setLocalDescription(answer); window.P2PClient.signal(from, { k: 'answer', xid: offer.xid, sdp: answer.sdp }); } // ---- receiver ----------------------------------------------------------------- function tryPeer({ videoId, cid, size, peer, onProgress }) { return new Promise((resolve) => { const xid = newXid(); const pc = new RTCPeerConnection({ iceServers }); const worker = new Worker('/p2p-recv-worker.js'); let settled = false; let idle = null; let expected = size; const end = (res) => { if (settled) return; settled = true; clearTimeout(idle); clearTimeout(openTimer); sessions.delete(xid); try { pc.close(); } catch { /* closed */ } if (!res.ok) worker.postMessage({ op: 'abort' }); setTimeout(() => worker.terminate(), res.ok ? 0 : 2000); resolve(res); }; const bump = () => { clearTimeout(idle); idle = setTimeout(() => end({ ok: false, error: 'peer went quiet' }), IDLE_TIMEOUT); }; const openTimer = setTimeout(() => end({ ok: false, error: 'could not connect' }), OPEN_TIMEOUT); worker.onmessage = (e) => { const m = e.data || {}; if (m.op === 'progress' && onProgress) onProgress(m.received, expected); if (m.op === 'done') end(m.ok ? { ok: true, sha256: m.sha256, size: m.size } : { ok: false, error: m.error }); }; worker.postMessage({ op: 'open', videoId }); sessions.set(xid, { pc, onSignal: (d) => { if (d.k === 'answer') pc.setRemoteDescription({ type: 'answer', sdp: d.sdp }).catch(() => end({ ok: false, error: 'bad answer' })); else if (d.k === 'ice' && d.cand) pc.addIceCandidate(d.cand).catch(() => {}); else if (d.k === 'deny') end({ ok: false, error: 'peer declined: ' + d.reason }); }, }); pc.onicecandidate = (e) => { if (e.candidate) window.P2PClient.signal(peer, { k: 'ice', xid, cand: e.candidate.toJSON() }); }; const dc = pc.createDataChannel('file', { ordered: true }); dc.binaryType = 'arraybuffer'; dc.onopen = () => { clearTimeout(openTimer); bump(); dc.send(JSON.stringify({ t: 'get' })); }; dc.onmessage = (e) => { bump(); if (typeof e.data === 'string') { let m; try { m = JSON.parse(e.data); } catch { return; } if (m.t === 'meta') { if (m.size !== size) end({ ok: false, error: 'peer has a different file size' }); expected = m.size; } else if (m.t === 'end') { clearTimeout(idle); worker.postMessage({ op: 'finish', cid, size }); } return; } worker.postMessage({ op: 'chunk', buf: e.data }, [e.data]); }; dc.onclose = () => { if (!settled) setTimeout(() => end({ ok: false, error: 'channel closed' }), 5000); }; (async () => { try { const offer = await pc.createOffer(); await pc.setLocalDescription(offer); if (!window.P2PClient.signal(peer, { k: 'offer', xid, cid, sdp: offer.sdp })) end({ ok: false, error: 'not connected to the P2P hub' }); } catch (err) { end({ ok: false, error: err.message }); } })(); }); } // Try each online holder in turn until one delivers a verified file. async function download({ videoId, cid, size, peers, onProgress }) { if (!window.P2PClient || !window.P2PClient.isConnected()) return { ok: false, error: 'not connected to the P2P hub' }; let last = { ok: false, error: 'no online device has this video' }; for (const p of peers || []) { last = await tryPeer({ videoId, cid, size, peer: p.peer, onProgress }); if (last.ok) return last; } return last; } function start(opts = {}) { if (opts.canShare) canShare = opts.canShare; if (opts.iceServers) iceServers = opts.iceServers; window.P2PClient.onMessage('signal', onSignal); } window.P2PTransfer = { start, download }; }());