Files
ytplayer/frontend/p2p-client.js

311 lines
13 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-client.js — this device's side of peer-to-peer sharing
* (docs/p2p-architecture.md flows 2–3). window.P2PClient.
*
* start({ getSettings, getProfile }) called once from app.js boot()
* changed() a save/delete happened → re-report soon
* device() { deviceId, secret } or null
* authHeaders() { 'X-Device': … } for other P2P calls
* contribute(videoId) / unknownVideos()
* hand a saved file the server doesn't
* know yet to /api/p2p/intake (plan 016)
* onMessage(type, fn) / signal(to, data) / peer()
* live /ws/p2p socket (plan 014): the
* device is "online" while it is open
*
* What it does, in order, each sync:
* 1. registers the device once (localStorage ytpDevice)
* 2. reconciles DeviceDB with the files really in OPFS (a record without a
* file is dropped; a file without a record is added as 'unhashed')
* 3. hashes 'unhashed' files and files not re-checked for 30 days, one at a
* time in hash-worker.js
* 4. reports the FULL holdings list (empty when sharing is off), answers the
* server's range challenges, and marks accepted files 'verified'
* P2P is ON by default: sharing runs unless settings.p2pShare === false or the
* server says P2P is disabled.
* ========================================================================== */
(function () {
'use strict';
const KEY = 'ytpDevice';
const REHASH_MS = 30 * 24 * 3600_000;
const RESYNC_MS = 6 * 3600_000;
let hooks = { getSettings: () => ({}), getProfile: () => '' };
let serverCfg = null;
let running = null;
let again = false;
let timer = null;
let lastSync = 0;
let lastUnknown = []; // cids the server did not recognise at the last report
function device() {
try {
const d = JSON.parse(localStorage.getItem(KEY) || 'null');
return d && /^dev_[0-9a-f]{16}$/.test(d.deviceId) && /^[0-9a-f]{64}$/.test(d.secret) ? d : null;
} catch { return null; }
}
const authHeaders = () => { const d = device(); return d ? { 'X-Device': d.deviceId + '.' + d.secret } : {}; };
async function config() {
if (serverCfg) return serverCfg;
try {
const j = await (await fetch('/api/p2p/config')).json();
if (j && j.ok) serverCfg = j;
} catch { /* offline — try again next sync */ }
return serverCfg;
}
async function ensureDevice() {
const have = device();
if (have) return have;
const r = await fetch('/api/p2p/device', {
method: 'POST', headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ fingerprint: window.getFingerprint ? window.getFingerprint() : '', profile: hooks.getProfile() || '' }),
});
const j = await r.json();
if (!j || !j.ok) throw new Error((j && j.error) || 'device registration failed');
const d = { deviceId: j.deviceId, secret: j.secret };
try { localStorage.setItem(KEY, JSON.stringify(d)); } catch { /* storage blocked */ }
return d;
}
function hashInWorker(name) {
return new Promise((resolve) => {
let w;
try { w = new Worker('/hash-worker.js'); } catch { resolve(null); return; }
w.onmessage = (e) => { w.terminate(); resolve(e.data && e.data.ok ? e.data : null); };
w.onerror = () => { w.terminate(); resolve(null); };
w.postMessage({ name });
});
}
async function sha256Range(videoId, offset, length) {
const bytes = window.OPFS && window.OPFS.readRange ? await window.OPFS.readRange(videoId, offset, length) : null;
if (!bytes) return null;
return window.Sha256.hex(bytes);
}
// DeviceDB ⇄ OPFS. Returns the records that describe a real file.
async function reconcile() {
const files = await window.OPFS.listVideos({ strict: true });
const byId = new Map(files.map((f) => [f.id, f]));
const recs = await window.DeviceDB.listFiles();
const out = [];
for (const r of recs) {
const f = byId.get(r.videoId);
if (!f) { await window.DeviceDB.deleteFile(r.videoId); continue; }
if (f.size !== r.size && r.size) { r.cid = null; r.state = 'unhashed'; r.size = f.size; await window.DeviceDB.putFile(r); }
r.name = f.name;
out.push(r);
byId.delete(r.videoId);
}
for (const f of byId.values()) {
if (String(f.id).startsWith('edit_')) continue; // edited cuts are never shared
const r = { videoId: f.id, cid: null, size: f.size, savedAt: Date.now(), lastCheckedAt: 0, state: 'unhashed' };
await window.DeviceDB.putFile(r);
out.push({ ...r, name: f.name });
}
return out;
}
async function hashPending(recs) {
const now = Date.now();
for (const r of recs) {
if (r.state !== 'unhashed' && now - (r.lastCheckedAt || 0) < REHASH_MS) continue;
const h = await hashInWorker(r.name);
if (!h) continue;
const changedCid = h.sha256 !== r.cid;
r.cid = h.sha256;
r.size = h.size;
r.lastCheckedAt = Date.now();
if (changedCid || r.state === 'unhashed') r.state = 'unverified';
await window.DeviceDB.putFile(r);
}
}
async function report(recs, dev) {
const share = hooks.getSettings().p2pShare !== false;
const items = recs.filter((r) => r.cid).map((r) => ({ cid: r.cid, size: r.size, videoId: r.videoId }));
const r = await fetch('/api/p2p/holdings', {
method: 'POST',
headers: { 'Content-Type': 'application/json', 'X-Device': dev.deviceId + '.' + dev.secret },
body: JSON.stringify({ share, items: share ? items : [], profile: hooks.getProfile() || '' }),
});
if (r.status === 401) { try { localStorage.removeItem(KEY); } catch { /* ignore */ } return null; }
const j = await r.json();
if (!j || !j.ok) return null;
lastUnknown = Array.isArray(j.unknown) ? j.unknown : [];
const accepted = new Set(j.accepted || []);
for (const rec of recs) {
if (rec.cid && accepted.has(rec.cid) && rec.state !== 'verified') { rec.state = 'verified'; await window.DeviceDB.putFile(rec); }
}
const answers = [];
for (const ch of j.challenges || []) {
const rec = recs.find((x) => x.cid === ch.cid);
if (!rec) continue;
const hex = await sha256Range(rec.videoId, ch.offset, ch.length);
if (hex) answers.push({ ...ch, sha256: hex });
}
if (answers.length) {
await fetch('/api/p2p/challenge', {
method: 'POST',
headers: { 'Content-Type': 'application/json', 'X-Device': dev.deviceId + '.' + dev.secret },
body: JSON.stringify({ answers }),
}).catch(() => {});
}
return j;
}
// ---- intake (plan 016) ---------------------------------------------------------
// Send one saved video to the server so it can hash + validate it and add its
// cid to the catalog. Uploads the whole file: only on the user's request
// ("Verify & share") or when the server asks for a copy (plan 018).
async function contribute(videoId, { cid = null, title = '', channel = '', restore = false } = {}) {
const dev = await ensureDevice();
const rec = await window.DeviceDB.getFile(videoId);
const file = window.OPFS.getFileObject ? await window.OPFS.getFileObject(videoId) : null;
if (!file) return { ok: false, error: 'not saved on this device' };
const h = { 'X-Device': dev.deviceId + '.' + dev.secret };
const t = await (await fetch('/api/p2p/intake', {
method: 'POST', headers: { ...h, 'Content-Type': 'application/json' },
body: JSON.stringify({ videoId, cid: cid || (rec && rec.cid) || undefined, size: file.size, title, channel, restore }),
})).json().catch(() => ({ ok: false, error: 'server unreachable' }));
if (!t.ok || t.known) { if (t.known) changed(); return t; }
const r = await (await fetch(t.url, { method: 'PUT', headers: h, body: file }))
.json().catch(() => ({ ok: false, error: 'upload failed' }));
if (r.ok) changed();
return r;
}
async function unknownVideos() {
if (!lastUnknown.length || !window.DeviceDB) return [];
const recs = await window.DeviceDB.listFiles();
return recs.filter((r) => r.cid && lastUnknown.includes(r.cid)).map((r) => r.videoId);
}
// ---- presence socket (/ws/p2p) — plan 014 ------------------------------------
// Open while sharing OR receiving is on. Credentials go in the first message,
// never the URL. Reconnects with backoff (5 s … 5 min).
let ws = null;
let myPeer = null;
let wsRetry = 0;
let wsTimer = null;
const listeners = new Map(); // type -> Set<fn>
const wantSocket = () => {
const s = hooks.getSettings();
return s.p2pShare !== false || s.p2pReceive !== false;
};
function onMessage(type, fn) {
if (!listeners.has(type)) listeners.set(type, new Set());
listeners.get(type).add(fn);
return () => listeners.get(type).delete(fn);
}
function emit(m) {
for (const fn of listeners.get(m.type) || []) { try { fn(m); } catch { /* a listener's bug is not ours */ } }
}
function connect() {
if (ws || !wantSocket()) return;
const d = device();
if (!d) return;
const proto = location.protocol === 'https:' ? 'wss:' : 'ws:';
let sock;
try { sock = new WebSocket(`${proto}//${location.host}/ws/p2p`); } catch { return; }
ws = sock;
let ping = null;
sock.onopen = () => {
sock.send(JSON.stringify({ type: 'auth', device: d.deviceId, secret: d.secret, profile: hooks.getProfile(), profileSecret: hooks.getProfileSecret?.() || '' }));
// The server drops sockets idle for 120 s (server.js websocketHandler).
ping = setInterval(() => { try { sock.send('{"type":"ping"}'); } catch { /* closing */ } }, 50_000);
};
sock.onmessage = (e) => {
let m;
try { m = JSON.parse(e.data); } catch { return; }
if (m.type === 'hello') { myPeer = m.peer; wsRetry = 0; }
emit(m);
};
sock.onclose = () => {
clearInterval(ping);
if (ws !== sock) return;
ws = null;
myPeer = null;
if (!wantSocket()) return;
clearTimeout(wsTimer);
wsTimer = setTimeout(connect, Math.min(300_000, 5000 * 2 ** wsRetry++));
};
}
function disconnect() {
clearTimeout(wsTimer);
const sock = ws;
ws = null;
myPeer = null;
if (sock) { try { sock.close(); } catch { /* gone */ } }
}
function signal(to, data) {
if (!ws || ws.readyState !== 1) return false;
ws.send(JSON.stringify({ type: 'signal', to, data }));
return true;
}
async function syncOnce() {
if (!window.OPFS || !window.OPFS.isSupported() || !window.DeviceDB || !window.Sha256) return null;
if (navigator.onLine === false) return null;
const cfg = await config();
if (!cfg || !cfg.enabled) return null;
const dev = await ensureDevice();
const recs = await reconcile();
await hashPending(recs);
const res = await report(recs, dev);
lastSync = Date.now();
if (wantSocket()) connect(); else disconnect();
return res;
}
// One sync at a time; a request during a sync schedules exactly one more.
function sync() {
if (running) { again = true; return running; }
running = syncOnce().catch(() => null).finally(() => {
running = null;
if (again) { again = false; sync(); }
});
return running;
}
function changed() {
clearTimeout(timer);
timer = setTimeout(sync, 5000);
}
// Plan 018: the server lost its copy and the source is gone — it asks one
// holder to send the file back. Only while sharing is on, one at a time,
// and only for the exact cid this device holds.
let restoring = false;
onMessage('upload-request', async (m) => {
if (restoring || hooks.getSettings().p2pShare === false || hooks.getSettings().directTransfer !== false) return;
const rec = window.DeviceDB ? await window.DeviceDB.getFile(m.videoId) : null;
if (!rec || rec.cid !== m.cid) return;
restoring = true;
try { await contribute(m.videoId, { cid: m.cid, restore: true }); } catch { /* next request retries */ }
finally { restoring = false; }
});
function start(h) {
hooks = { ...hooks, ...(h || {}) };
const idle = window.requestIdleCallback || ((fn) => setTimeout(fn, 1));
setTimeout(() => idle(() => sync()), 8000);
document.addEventListener('visibilitychange', () => {
if (document.visibilityState === 'visible' && Date.now() - lastSync > RESYNC_MS) sync();
});
}
window.P2PClient = {
start, changed, sync, device, authHeaders, config, contribute, unknownVideos,
send: m => { if (!ws || ws.readyState !== 1) return false; ws.send(JSON.stringify(m)); return true; },
onMessage, signal, peer: () => myPeer, isConnected: () => !!(ws && ws.readyState === 1 && myPeer),
};
}());