// 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('legacy unconfirmed signalling is refused; 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(a.sent.at(-1)).toMatchObject({ type: 'error', error: 'use a confirmed direct-transfer invitation' }); expect(b.sent.some(m => m.type === 'signal')).toBe(false); 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: 'use a confirmed direct-transfer invitation' }); 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 }); test('rehydrator asks one online sharing holder, once per 10 min, only when the server has no copy', async () => { const { createRehydrator } = await import('./p2p-hub.js'); const sent = []; let t = NOW; let serverHas = false; const hub = { isOnline: (d) => d === 'dev_000000000000000a' || d === 'dev_000000000000000c', send: (d, m) => { sent.push([d, m]); return true; } }; const rehydrate = createRehydrator({ p2pDb, hub, hasServerCopy: async () => serverHas, now: () => t }); expect(await rehydrate('dQw4w9WgXcQ')).toBe(true); // c is online but has sharing off; a is online and shares. expect(sent).toEqual([['dev_000000000000000a', { type: 'upload-request', videoId: 'dQw4w9WgXcQ', cid: CID }]]); expect(await rehydrate('dQw4w9WgXcQ')).toBe(true); // recently asked — no second message expect(sent.length).toBe(1); t += 11 * 60_000; serverHas = true; expect(await rehydrate('dQw4w9WgXcQ')).toBe(false); // server copy is back expect(await rehydrate('unknownVid1')).toBe(false); // nobody holds it }); test('direct directory and invitations are restricted to authorized profile sockets', async () => { const hub = createP2pHub({ getDevice: p2pDb.getDevice, authorizeProfile: async (name, secret) => secret === 'proof' ? name : '' }); const a = fakeWs({ budget: [] }), b = fakeWs({ budget: [] }), c = fakeWs({ budget: [] }); for (const [ws, id, secret, profile] of [[a, 'a', SECRET_A, 'shared'], [b, 'b', SECRET_B, 'shared'], [c, 'c', 'x', 'different']]) { await hub.websocket.message(ws, JSON.stringify({ type: 'auth', device: 'dev_000000000000000' + id, secret, profile, profileSecret: 'proof' })); } await hub.websocket.message(a, JSON.stringify({ type: 'direct-inventory', ids: ['localSong'] })); expect(b.sent.at(-1).peers[0].files).toEqual(['localSong']); expect(c.sent.at(-1).peers).toEqual([]); const before = c.sent.length; await hub.websocket.message(a, JSON.stringify({ type: 'direct', action: 'request', to: peerIdOf('dev_000000000000000c'), id: 'localSong' })); expect(c.sent.length).toBe(before); await hub.websocket.message(b, JSON.stringify({ type: 'direct', action: 'request', to: peerIdOf('dev_000000000000000a'), id: 'localSong' })); expect(a.sent.at(-1).action).toBe('request'); }); test('binary websocket payloads cannot serve as signalling envelopes', async () => { const hub = createP2pHub({ getDevice: p2pDb.getDevice }); const ws = fakeWs({ budget: [] }); await hub.websocket.message(ws, Buffer.from(JSON.stringify({ type: 'auth', device: 'dev_000000000000000a', secret: SECRET_A }))); expect(ws.sent).toEqual([]); expect(ws.data.authed).not.toBe(true); });