Add the /ws/p2p presence and signalling hub and the holders endpoint

This commit is contained in:
Claude
2026-09-30 07:42:35 +00:00
parent ef7d08f16a
commit 6fa930ad62
7 changed files with 311 additions and 4 deletions

View File

@@ -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<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 }));
// 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),
};
}());

View File

@@ -21,7 +21,7 @@ green, app boots with no JS errors, P2P on by default, offline boot works).
| 011 | 011-browser-sha256-e1793d | Add an incremental SHA-256 library for the browser and node tests | done | Add an incremental SHA-256 library for the browser and node tests | |
| 012 | 012-device-file-registry-288d55 | Add the on-device IndexedDB file registry and hash saves while downloading | done | Add the on-device IndexedDB file registry and hash saves while downloading | browser harness |
| 013 | 013-device-identity-and-holdings-3ba493 | Register devices and report verified holdings to the server | done | Register devices and report verified holdings to the server | persistent holders, no TTL |
| 014 | 014-p2p-presence-hub-ceced8 | Add the /ws/p2p presence and signalling hub and the holders endpoint | in-progress | | stale flag, never hidden |
| 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 | queued | | |
| 016 | 016-intake-and-server-verification-cfe031 | Let a device hand a file to the server for hashing and validation | queued | | needs ffmpeg for tests |
| 017 | 017-peer-transfer-1faaa7 | Download a verified file from another device over WebRTC | queued | | STUN only |

View File

@@ -120,3 +120,9 @@ 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-hub.js`, `p2p-hub.test.js`, `p2p-client.js` byte-identical to the pre-tested versions; hub tests 3 pass; all 13 server test files 0 fail; `SERVER_OK`; `FRONT_OK`; two-browser presence check relayed the signal (`{"peers":[true,true],"distinct":true,"sent":true,"got":[{"from":true,"data":{"hi":1}}]}`). No leftover processes.
- Executor Findings (verbatim): Both git apply patches applied successfully without errors. All 13 test files (24 total test groups) passed with 0 failures. Server builds successfully. Frontend syntax check passes. Browser presence test matches expected output exactly. All modifications applied according to plan steps 1-7. No files touched outside the specified changes.

120
server/p2p-hub.js Normal file
View File

@@ -0,0 +1,120 @@
/* ============================================================================
* p2p-hub.js — live presence + WebRTC signalling for P2P devices, and the
* availability payload (docs/p2p-architecture.md flows 4–5).
*
* /ws/p2p (credentials never go in the URL — proxies log URLs)
* you → server {type:'auth', device, secret} first message, within 10 s
* server → you {type:'hello', peer} after auth
* you → peer {type:'signal', to:<peer>, data} relayed as
* server → peer {type:'signal', from:<your peer>, data}
* server → you {type:'error', error}
* Server-initiated messages (plan 018) go through hub.send(deviceId, msg).
*
* "Online" lives only in memory: a server restart forgets it, a closed
* socket drops it. Holder ROWS are persistent (p2p_holders) — the payload
* reports both: `online` (now) and `lastVerifiedAt` / `stale` (history).
* Peers are addressed by an opaque id (peerIdOf), never by device id.
* ========================================================================== */
import { createHash, timingSafeEqual } from 'node:crypto';
const sha = (s) => createHash('sha256').update(String(s)).digest('hex');
const same = (a, b) => { const x = Buffer.from(String(a)), y = Buffer.from(String(b)); return x.length === y.length && timingSafeEqual(x, y); };
export const peerIdOf = (deviceId) => sha('peer:' + deviceId).slice(0, 12);
const MAX_MSG = 64 * 1024;
const BUDGET_PER_MIN = 300;
export function createP2pHub({ getDevice, enabled = () => true, now = () => Date.now(), log = console } = {}) {
const online = new Map(); // deviceId -> ws
const byPeer = new Map(); // peer -> deviceId (online only)
const send = (ws, m) => { try { ws.send(JSON.stringify(m)); } catch { /* gone */ } };
function upgrade(req, server) {
if (!enabled()) return new Response('p2p disabled', { status: 404 });
const ok = server.upgrade(req, { data: { hub: 'p2p', authed: false, budget: [] } });
return ok ? undefined : new Response('upgrade failed', { status: 400 });
}
function open(ws) {
const t = setTimeout(() => { if (!ws.data.authed) { try { ws.close(4401, 'auth timeout'); } catch { /* gone */ } } }, 10_000);
t.unref?.();
}
async function auth(ws, m) {
const deviceId = String(m.device || '');
const secret = String(m.secret || '');
const d = /^dev_[0-9a-f]{16}$/.test(deviceId) ? await getDevice(deviceId).catch(() => null) : null;
if (!d || !same(d.secret_hash, sha(secret))) { try { ws.close(4401, 'unknown device'); } catch { /* gone */ } return; }
ws.data.authed = true;
ws.data.deviceId = deviceId;
ws.data.peer = peerIdOf(deviceId);
const prev = online.get(deviceId);
if (prev && prev !== ws) { prev.data.deviceId = null; try { prev.close(4409, 'replaced'); } catch { /* gone */ } }
online.set(deviceId, ws);
byPeer.set(ws.data.peer, deviceId);
send(ws, { type: 'hello', peer: ws.data.peer });
}
async function message(ws, raw) {
const t = now();
ws.data.budget = ws.data.budget.filter((x) => t - x < 60_000);
if (ws.data.budget.length >= BUDGET_PER_MIN) { send(ws, { type: 'error', error: 'slow down' }); return; }
ws.data.budget.push(t);
const s = typeof raw === 'string' ? raw : Buffer.from(raw).toString('utf8');
if (s.length > MAX_MSG) return;
let m;
try { m = JSON.parse(s); } catch { return; }
if (!ws.data.authed) { if (m && m.type === 'auth' && !ws.data.authing) { ws.data.authing = true; await auth(ws, m); } return; }
if (m && m.type === 'signal' && typeof m.to === 'string') {
const target = online.get(byPeer.get(m.to));
if (!target || target === ws) { send(ws, { type: 'error', error: 'peer offline', to: m.to }); return; }
send(target, { type: 'signal', from: ws.data.peer, data: m.data });
}
}
function close(ws) {
const { deviceId, peer } = ws.data || {};
if (deviceId && online.get(deviceId) === ws) {
online.delete(deviceId);
byPeer.delete(peer);
}
}
return {
upgrade,
websocket: { open, message, close },
isOnline: (deviceId) => online.has(deviceId),
send: (deviceId, msg) => { const ws = online.get(deviceId); if (!ws) return false; send(ws, msg); return true; },
onlineCount: () => online.size,
};
}
// GET /api/p2p/holders payload. Only devices that share are listed; holders
// are never dropped for age — `stale` just flags an old last check.
export async function holdersPayload({ videoId, cid, p2pDb, isOnline, staleDays, serverHas, now = Date.now() }) {
const contents = cid
? [await p2pDb.getContent(cid)].filter((c) => c && c.status === 'verified')
: await p2pDb.listContentForVideo(videoId);
const staleMs = staleDays * 86400_000;
const out = [];
for (const c of contents) {
const holders = (await p2pDb.listHolders(c.cid, 50))
.filter((h) => Number(h.share) === 1)
.map((h) => ({
peer: peerIdOf(h.device_id),
online: isOnline(h.device_id),
lastVerifiedAt: Number(h.last_verified_at),
stale: now - Number(h.last_verified_at) > staleMs,
trust: h.trust,
}))
.sort((a, b) => (b.online - a.online) || (b.lastVerifiedAt - a.lastVerifiedAt));
out.push({
cid: c.cid, videoId: c.video_id, size: Number(c.size), height: Number(c.height), vcodec: c.vcodec,
serverHas: !!(await serverHas(c.cid)),
counts: { holders: holders.length, online: holders.filter((h) => h.online).length, fresh: holders.filter((h) => !h.stale).length },
holders,
});
}
return { ok: true, staleDays, cids: out };
}

87
server/p2p-hub.test.js Normal file
View File

@@ -0,0 +1,87 @@
// Presence/signalling hub and the holders payload (plan 014).
import { test, expect, beforeAll } from 'bun:test';
import { mkdtempSync } from 'node:fs';
import { tmpdir } from 'node:os';
import { join } from 'node:path';
import { createHash } from 'node:crypto';
const root = mkdtempSync(join(tmpdir(), 'ytp-p2p-hub-'));
process.env.DB_PATH = join(root, 'test.db');
const dbmod = await import('./db.js');
const p2pDb = await import('./p2p-db.js');
const { createP2pHub, holdersPayload, peerIdOf } = await import('./p2p-hub.js');
const sha = (s) => createHash('sha256').update(s).digest('hex');
const SECRET_A = '1'.repeat(64);
const SECRET_B = '2'.repeat(64);
const CID = 'a'.repeat(64);
const DAY = 86400_000;
const NOW = Date.UTC(2026, 8, 29, 12);
function fakeWs(data) {
return { data, sent: [], closed: null, send(s) { this.sent.push(JSON.parse(s)); }, close(code) { this.closed = code; } };
}
beforeAll(async () => {
await dbmod.initDb();
await p2pDb.initP2pSchema();
await p2pDb.createDevice({ deviceId: 'dev_000000000000000a', secretHash: sha(SECRET_A), now: NOW });
await p2pDb.createDevice({ deviceId: 'dev_000000000000000b', secretHash: sha(SECRET_B), now: NOW });
await p2pDb.createDevice({ deviceId: 'dev_000000000000000c', secretHash: sha('x'), now: NOW });
await p2pDb.touchDevice('dev_000000000000000c', { now: NOW, share: false });
await p2pDb.upsertContent({ cid: CID, videoId: 'dQw4w9WgXcQ', size: 1234, height: 720, vcodec: 'h264', origin: 'server', now: NOW });
await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000a', now: NOW - 1 * DAY });
await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000b', now: NOW - 30 * DAY }); // stale, still listed
await p2pDb.upsertHolder({ cid: CID, deviceId: 'dev_000000000000000c', now: NOW }); // share off → hidden
});
test('auth on open, hello with opaque peer id, bad secret is closed', async () => {
const hub = createP2pHub({ getDevice: p2pDb.getDevice });
const a = fakeWs({ budget: [] });
hub.websocket.open(a);
await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: 'x', data: {} })); // ignored before auth
expect(a.sent.length).toBe(0);
await hub.websocket.message(a, JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }));
expect(a.sent[0]).toEqual({ type: 'hello', peer: peerIdOf('dev_000000000000000a') });
expect(hub.isOnline('dev_000000000000000a')).toBe(true);
const bad = fakeWs({ budget: [] });
hub.websocket.open(bad);
await hub.websocket.message(bad, JSON.stringify({ type: 'auth', device: 'dev_000000000000000b', secret: SECRET_A }));
expect(bad.closed).toBe(4401);
expect(hub.isOnline('dev_000000000000000b')).toBe(false);
});
test('signals are relayed between online peers only; close drops presence', async () => {
const hub = createP2pHub({ getDevice: p2pDb.getDevice });
const a = fakeWs({ budget: [] });
const b = fakeWs({ budget: [] });
hub.websocket.open(a);
hub.websocket.open(b);
await hub.websocket.message(a, JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }));
await hub.websocket.message(b, JSON.stringify({ type: 'auth', device: 'dev_000000000000000b', secret: SECRET_B }));
await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: { sdp: 'x' } }));
expect(b.sent.at(-1)).toEqual({ type: 'signal', from: a.data.peer, data: { sdp: 'x' } });
hub.websocket.close(b);
expect(hub.isOnline('dev_000000000000000b')).toBe(false);
await hub.websocket.message(a, JSON.stringify({ type: 'signal', to: b.data.peer, data: {} }));
expect(a.sent.at(-1)).toMatchObject({ type: 'error', error: 'peer offline' });
expect(hub.send('dev_000000000000000a', { type: 'ping' })).toBe(true);
expect(hub.send('dev_000000000000000b', { type: 'ping' })).toBe(false);
});
test('holders payload: persistent rows, online flag, stale marker, share-off hidden', async () => {
const online = new Set(['dev_000000000000000b']);
const p = await holdersPayload({
videoId: 'dQw4w9WgXcQ', p2pDb, isOnline: (d) => online.has(d), staleDays: 7,
serverHas: async () => false, now: NOW,
});
expect(p.staleDays).toBe(7);
expect(p.cids.length).toBe(1);
const c = p.cids[0];
expect(c).toMatchObject({ cid: CID, size: 1234, height: 720, serverHas: false, counts: { holders: 2, online: 1, fresh: 1 } });
expect(c.holders.map((h) => [h.peer, h.online, h.stale])).toEqual([
[peerIdOf('dev_000000000000000b'), true, true], // online first even though stale
[peerIdOf('dev_000000000000000a'), false, false],
]);
expect(JSON.stringify(p)).not.toContain('dev_'); // never leak device ids
});

View File

@@ -6,7 +6,7 @@
"scripts": {
"start": "bun server.js",
"dev": "bun --hot server.js",
"test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js && bun test ./p2p-admit.test.js && bun test ./p2p-retention.test.js && bun test ./p2p-routes.test.js"
"test": "bun test --timeout 60000 ./media-cache.test.js && bun test ./notes.test.js && bun test ./remote.test.js && bun test ./party.test.js && bun test ./uploads.test.js && bun test ./innertube.test.js && bun test ./ytdlp-pool.test.js && bun test ./p2p-db.test.js && bun test ./p2p-admit.test.js && bun test ./p2p-retention.test.js && bun test ./p2p-routes.test.js && bun test ./p2p-hub.test.js"
},
"dependencies": {
"@hono/node-server": "^1.14.0",

View File

@@ -52,6 +52,7 @@ import { admitFile } from './p2p-admit.js';
import { P2P } from './p2p-config.js';
import * as p2pDb from './p2p-db.js';
import { registerP2pRoutes } from './p2p-routes.js';
import { createP2pHub, holdersPayload } from './p2p-hub.js';
import { sha256Range } from './hash.js';
import * as innertube from './innertube.js';
import QRCode from 'qrcode';
@@ -1884,7 +1885,9 @@ const party = createPartyHub();
// Bun allows ONE websocket handler per server: party sockets are tagged
// (ws.data.hub === 'party'), everything else belongs to the remote relay.
const pickHub = (ws) => (ws.data && ws.data.hub === 'party' ? party.websocket : remote.websocket);
// …and P2P sockets are tagged ws.data.hub === 'p2p' (p2p-hub.js).
const pickHub = (ws) => (ws.data && ws.data.hub === 'party' ? party.websocket
: ws.data && ws.data.hub === 'p2p' ? p2pHub.websocket : remote.websocket);
const websocketHandler = {
maxPayloadLength: 256 * 1024,
idleTimeout: 120,
@@ -1962,6 +1965,21 @@ async function fileForCid(cid) {
return (await f.exists()) ? { path, size: f.size } : null;
}
const p2p = registerP2pRoutes(app, { cfg: P2P, p2pDb, fileForCid, sha256Range });
const p2pHub = createP2pHub({ getDevice: p2pDb.getDevice, enabled: () => P2P.enabled });
// GET /api/p2p/holders?v=<videoId>|cid=<sha256> — who holds a copy. Rows are
// persistent; `stale` flags a holder not re-verified for P2P_STALE_DAYS.
app.get('/api/p2p/holders', p2p.gate, async (c) => {
const v = (c.req.query('v') || '').trim();
const cid = (c.req.query('cid') || '').trim().toLowerCase();
const hasCid = /^[0-9a-f]{64}$/.test(cid);
if (!hasCid && !/^[A-Za-z0-9_-]{6,64}$/.test(v)) return c.json({ ok: false, error: 'missing v or cid' }, 400);
const payload = await holdersPayload({
videoId: v, cid: hasCid ? cid : null, p2pDb, isOnline: p2pHub.isOnline,
staleDays: P2P.staleDays, serverHas: fileForCid,
});
return c.json(payload, 200, { 'Cache-Control': 'no-store' });
});
// ============================================================================
// GET /sw.js — serve the service worker with BUILD_TAG injected
@@ -2081,6 +2099,7 @@ async function main() {
const path = new URL(req.url).pathname;
if (path === '/ws/remote') return remote.upgrade(req, server, clientIpOf(req, server));
if (path === '/ws/party') return party.upgrade(req, server, clientIpOf(req, server));
if (path === '/ws/p2p') return p2pHub.upgrade(req, server);
return app.fetch(req, server);
},
websocket: websocketHandler,