--- a/frontend/p2p-client.js +++ b/frontend/p2p-client.js @@ -6,6 +6,9 @@ * changed() a save/delete happened → re-report soon * device() { deviceId, secret } or null * authHeaders() { 'X-Device': … } for other P2P calls + * 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) @@ -148,6 +151,74 @@ return j; } + // ---- 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 + + 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 })); + // 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; @@ -158,6 +229,7 @@ await hashPending(recs); const res = await report(recs, dev); lastSync = Date.now(); + if (wantSocket()) connect(); else disconnect(); return res; } @@ -185,5 +257,8 @@ }); } - window.P2PClient = { start, changed, sync, device, authHeaders, config }; + window.P2PClient = { + start, changed, sync, device, authHeaders, config, + onMessage, signal, peer: () => myPeer, isConnected: () => !!(ws && ws.readyState === 1 && myPeer), + }; }());