Files
ytplayer/plans/patches/017-p2p-transfer-new.diff

281 lines
12 KiB
Diff

--- /dev/null
+++ b/frontend/p2p-transfer.js
@@ -0,0 +1,184 @@
+/* ============================================================================
+ * 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 };
+}());
--- /dev/null
+++ b/frontend/p2p-recv-worker.js
@@ -0,0 +1,90 @@
+/* ============================================================================
+ * p2p-recv-worker.js — writes a file arriving from another device into OPFS
+ * while hashing it; commits it ONLY when the SHA-256 equals the expected
+ * content id (docs/p2p-architecture.md flow 6).
+ *
+ * In: { op:'open', videoId } start videos/<videoId>.p2p.part
+ * { op:'chunk', buf } ArrayBuffer (transferred)
+ * { op:'finish', cid, size } verify + rename to <videoId>.mp4
+ * { op:'abort' } drop the partial file
+ * Out: { op:'opened' } | { op:'progress', received } |
+ * { op:'done', ok:true, sha256, size } | { op:'done', ok:false, error }
+ * The ".part" suffix keeps OPFS.listVideos() from ever listing a partial file.
+ * ========================================================================== */
+'use strict';
+importScripts('/sha256.js');
+
+let dir = null;
+let handle = null;
+let access = null;
+let partName = '';
+let finalName = '';
+let hasher = null;
+let offset = 0;
+let lastProgress = 0;
+
+async function cleanup() {
+ try { if (access) access.close(); } catch { /* closed */ }
+ access = null;
+ try { if (dir && partName) await dir.removeEntry(partName); } catch { /* gone */ }
+}
+
+// Messages are handled strictly one after another: an async handler would
+// otherwise let 'chunk' or 'abort' run while 'open' is still awaiting OPFS.
+let chain = Promise.resolve();
+self.onmessage = (e) => { chain = chain.then(() => onOp(e.data || {})); };
+
+async function onOp(m) {
+ try {
+ if (m.op === 'open') {
+ const root = await navigator.storage.getDirectory();
+ dir = await root.getDirectoryHandle('videos', { create: true });
+ partName = `${m.videoId}.p2p.part`;
+ finalName = `${m.videoId}.mp4`;
+ handle = await dir.getFileHandle(partName, { create: true });
+ access = await handle.createSyncAccessHandle();
+ access.truncate(0);
+ hasher = self.Sha256.create();
+ offset = 0;
+ self.postMessage({ op: 'opened' });
+ } else if (m.op === 'chunk') {
+ const u8 = new Uint8Array(m.buf);
+ access.write(u8, { at: offset });
+ hasher.update(u8);
+ offset += u8.byteLength;
+ if (offset - lastProgress > 1024 * 1024) { lastProgress = offset; self.postMessage({ op: 'progress', received: offset }); }
+ } else if (m.op === 'finish') {
+ access.truncate(offset);
+ access.flush();
+ access.close();
+ access = null;
+ const sha256 = hasher.hex();
+ if (offset !== m.size) throw new Error(`size mismatch (${offset} of ${m.size})`);
+ if (sha256 !== m.cid) throw new Error('content hash mismatch');
+ try { await dir.removeEntry(finalName); } catch { /* no previous copy */ }
+ let renamed = false;
+ if (typeof handle.move === 'function') { try { await handle.move(finalName); renamed = true; } catch { /* copy below */ } }
+ if (!renamed) {
+ const out = await (await dir.getFileHandle(finalName, { create: true })).createSyncAccessHandle();
+ try {
+ const file = await handle.getFile();
+ for (let pos = 0; pos < file.size; pos += 8 * 1024 * 1024) {
+ const buf = new Uint8Array(await file.slice(pos, pos + 8 * 1024 * 1024).arrayBuffer());
+ out.write(buf, { at: pos });
+ }
+ out.truncate(file.size);
+ out.flush();
+ } finally { out.close(); }
+ await dir.removeEntry(partName);
+ }
+ partName = '';
+ self.postMessage({ op: 'done', ok: true, sha256, size: offset });
+ } else if (m.op === 'abort') {
+ await cleanup();
+ self.postMessage({ op: 'done', ok: false, error: 'aborted' });
+ }
+ } catch (err) {
+ await cleanup();
+ self.postMessage({ op: 'done', ok: false, error: err && err.message ? err.message : String(err) });
+ }
+}