311 lines
13 KiB
JavaScript
311 lines
13 KiB
JavaScript
/* ============================================================================
|
||
* 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 = window.Lazy ? window.Lazy.worker('/hash-worker.js') : 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()), h?.initialDelay ?? 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),
|
||
};
|
||
}());
|