From c8837befed1aeb63e09baca480fd4b7a19c41790 Mon Sep 17 00:00:00 2001 From: Claude Date: Wed, 30 Sep 2026 07:54:34 +0000 Subject: [PATCH] Download a verified file from another device over WebRTC --- frontend/app.js | 63 +++++- frontend/index.html | 1 + frontend/p2p-recv-worker.js | 90 +++++++++ frontend/p2p-transfer.js | 184 ++++++++++++++++++ frontend/styles.css | 3 + frontend/sw.js | 2 + plans/INDEX.md | 2 +- .../017-peer-transfer-1faaa7.md | 7 + 8 files changed, 350 insertions(+), 2 deletions(-) create mode 100644 frontend/p2p-recv-worker.js create mode 100644 frontend/p2p-transfer.js rename plans/{active => done}/017-peer-transfer-1faaa7.md (86%) diff --git a/frontend/app.js b/frontend/app.js index 1f78c49..0431fd0 100755 --- a/frontend/app.js +++ b/frontend/app.js @@ -1786,10 +1786,11 @@ const Player = { btn.className = 'retry-btn'; btn.textContent = 'โ†ป Retry'; btn.addEventListener('click', () => { - btn.remove(); + els.playerPane.querySelectorAll('.retry-btn').forEach((b) => b.remove()); Player.loadVideo(videoObj, { preferStream, resume, reveal, nocache }); }); els.playerPane.appendChild(btn); + offerPeerDownload(videoObj); } }, @@ -2465,6 +2466,65 @@ async function renderAvailability(id) { el.classList.remove('hidden'); } +// Fetch a verified copy from another device (docs/p2p-architecture.md flow 6). +// Download-then-play: on success the file is saved offline like any save. +async function getFromPeers(videoObj, c, peers, onProgress) { + const id = videoObj.id; + downloading.add(id); + markCardCacheState(id, 'downloading'); + updateDownloadBadge(); + try { + const res = await window.P2PTransfer.download({ videoId: id, cid: c.cid, size: c.size, peers, onProgress }); + if (res.ok) { + cachedIds.add(id); + cacheMutations++; + recordDeviceFile(id, { sha256: res.sha256, expectedSha: c.cid, size: res.size }); + toast(`Saved "${videoObj.title || id}" from another device โœ“`); + } + return res; + } finally { + downloading.delete(id); + markCardCacheState(id, cachedIds.has(id) ? 'cached' : 'none'); + updateDownloadBadge(); + renderSidebar(); + } +} + +// Offered when a video won't load: only if some device holding it is online now. +async function offerPeerDownload(videoObj) { + if (!WEB || !window.P2PTransfer || !window.P2PCore || data.settings.p2pReceive === false) return; + const id = videoObj && videoObj.id; + if (!id || videoObj.custom || cachedIds.has(id)) return; + let payload = null; + try { + const r = await fetch(`/api/p2p/holders?v=${encodeURIComponent(id)}`); + payload = r.ok ? await r.json() : null; + } catch { return; } + const c = payload && Array.isArray(payload.cids) ? payload.cids[0] : null; + const peers = c ? window.P2PCore.pickPeers(payload, c.cid) : []; + if (!peers.length || !els.playerPane.querySelector('.retry-btn')) return; // nobody online / user moved on + const old = els.playerPane.querySelector('.peer-btn'); + if (old) old.remove(); + const btn = document.createElement('button'); + btn.className = 'retry-btn peer-btn'; + btn.textContent = `๐Ÿ“ก Get it from a device (${peers.length} online)`; + btn.addEventListener('click', async () => { + btn.disabled = true; + const res = await getFromPeers(videoObj, c, peers, (got, total) => { + btn.textContent = `๐Ÿ“ก ${Math.min(99, Math.round((got / total) * 100))}%โ€ฆ`; + }); + if (res.ok) { + els.playerPane.querySelectorAll('.retry-btn').forEach((b) => b.remove()); + Player.loadVideo(videoObj); + } else { + btn.disabled = false; + btn.textContent = '๐Ÿ“ก Try again'; + toast('โš  ' + res.error); + } + }); + els.playerPane.appendChild(btn); +} + // Reflect cache state on the now-playing Save button. function updateNowPlayingActions() { if (!current || !current.meta) return; @@ -9981,6 +10041,7 @@ async function boot() { // Peer-to-peer: report what this device holds (on by default; settings.p2pShare). if (WEB && window.P2PClient) { window.P2PClient.start({ getSettings: () => data.settings, getProfile: () => (data.profile && data.profile.name) || '' }); + if (window.P2PTransfer) window.P2PTransfer.start({ canShare: () => data.settings.p2pShare !== false }); } data.playlists.forEach(preloadPlaylist); preloadPinnedPlaylists(); diff --git a/frontend/index.html b/frontend/index.html index cc1a29c..e1d6581 100755 --- a/frontend/index.html +++ b/frontend/index.html @@ -566,6 +566,7 @@ + diff --git a/frontend/p2p-recv-worker.js b/frontend/p2p-recv-worker.js new file mode 100644 index 0000000..66f76d8 --- /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/.p2p.part + * { op:'chunk', buf } ArrayBuffer (transferred) + * { op:'finish', cid, size } verify + rename to .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) }); + } +} diff --git a/frontend/p2p-transfer.js b/frontend/p2p-transfer.js new file mode 100644 index 0000000..74c8c21 --- /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 }; +}()); diff --git a/frontend/styles.css b/frontend/styles.css index 08dc941..8564f91 100755 --- a/frontend/styles.css +++ b/frontend/styles.css @@ -4343,3 +4343,6 @@ body.landscape-fs .player-stage { touch-action: none; } /* fullscreen (real or t transform: none; } } + +/* P2P: "Get it from a device" next to Retry (plan 017) */ +.peer-btn { margin-left: 8px; } diff --git a/frontend/sw.js b/frontend/sw.js index eb718bd..aeabe0a 100644 --- a/frontend/sw.js +++ b/frontend/sw.js @@ -69,6 +69,8 @@ const SHELL = [ '/hash-worker.js', '/p2p-client.js', '/p2p-core.js', + '/p2p-transfer.js', + '/p2p-recv-worker.js', '/app.js', '/manifest.webmanifest', '/icons/icon-192.png', diff --git a/plans/INDEX.md b/plans/INDEX.md index c3a6e42..60aaf5c 100644 --- a/plans/INDEX.md +++ b/plans/INDEX.md @@ -24,6 +24,6 @@ green, app boots with no JS errors, P2P on by default, offline boot works). | 014 | 014-p2p-presence-hub-ceced8 | Add the /ws/p2p presence and signalling hub and the holders endpoint | done | Add the /ws/p2p presence and signalling hub and the holders endpoint | stale flag, never hidden | | 015 | 015-availability-ui-and-settings-3b9397 | Show peer availability with stale markers and add Sharing settings | done | Show peer availability with stale markers and add Sharing settings | | | 016 | 016-intake-and-server-verification-cfe031 | Let a device hand a file to the server for hashing and validation | done | Let a device hand a file to the server for hashing and validation | needs ffmpeg for tests | -| 017 | 017-peer-transfer-1faaa7 | Download a verified file from another device over WebRTC | in-progress | | STUN only | +| 017 | 017-peer-transfer-1faaa7 | Download a verified file from another device over WebRTC | done | Download a verified file from another device over WebRTC | STUN only | | 018 | 018-server-rehydrate-from-peer-4fb8bd | Restore an evicted server copy from an online holder | queued | | | | 019 | 019-admin-p2p-panel-4dc623 | Add a P2P panel to the admin page | queued | | | diff --git a/plans/active/017-peer-transfer-1faaa7.md b/plans/done/017-peer-transfer-1faaa7.md similarity index 86% rename from plans/active/017-peer-transfer-1faaa7.md rename to plans/done/017-peer-transfer-1faaa7.md index 6399c79..cf48ddf 100644 --- a/plans/active/017-peer-transfer-1faaa7.md +++ b/plans/done/017-peer-transfer-1faaa7.md @@ -177,3 +177,10 @@ Output ONLY the following, no other prose: 3. `Findings:` โ€” max 10 lines. Do not commit. Do not push. Do not touch files outside the Steps. + +## Execution log + +- Executor: in-session Agent (haiku). Attempts: 1. Fix rounds: 0. +- Orchestrator re-ran Verification: `p2p-transfer.js` and `p2p-recv-worker.js` byte-identical to the pre-tested versions; all frontend files pass `node --check`; 61 frontend tests pass; index.html references the new script once and sw.js twice; two-browser WebRTC check reproduced the expected JSON (5 MiB delivered and hash-verified, unknown cid `peer declined: not here`, corrupted holder `content hash mismatch` with no partial file left); full app load in Chromium: no JS errors, all globals present including `P2PTransfer`, share/receive on by default, device online. No leftover processes. +- Orchestrator note: the executor's `npm`/`bun install` left an untracked root `bun.lock`; deleted before commit (not part of the plan). +- Executor Findings (verbatim): All syntax checks passed (FRONT_OK). All unit tests passed (61 pass, 0 fail). Script inclusions correct: index.html has 1 entry, sw.js has 2 entries. P2P transfer test executed successfully with expected JSON output containing good/bad/corrupt test cases. New functions getFromPeers and offerPeerDownload added correctly to app.js. P2P transfer and recv worker files created by git patch successfully.