/* ============================================================================ * party.js — watch party relay: synced playback, text chat, voice signalling * * One host drives playback; guests follow. The server only relays: * /ws/party?role=host[&code=…&secret=…]&pid=…&name=… create or resume a party * /ws/party?role=guest&code=…&pid=…&name=… join one * * Messages (JSON): * host → all {type:'state', state} (server stamps state.ts = server ms) * host → all {type:'settings', allowControl} * any → all {type:'chat', text} (server adds from/pid/at; ≤ 500 chars) * any → all {type:'voice', on} (voice-chat presence) * any → one {type:'rtc', to, data} (WebRTC offer/answer/ICE, relayed) * guest → host {type:'cmd', cmd, …} (only when the host allows control) * server → you {type:'hello', code, secret?, you, now, members, state, chat, allowControl} * server → all {type:'members', members} · {type:'ended'} * A party survives its host dropping for PARTY_GRACE_MS (a reload, a flaky * link); the host resumes it with the secret it was given. * ========================================================================== */ import { createDirectRelay } from './direct-relay.js'; import { createDJQueue, changeDJQueue, queueSnapshot, takeDJVideo } from './party-dj.js'; import { randomBytes, randomInt, timingSafeEqual } from 'node:crypto'; const ALPHABET = 'ABCDEFGHJKLMNPQRSTUVWXYZ23456789'; // no 0/O, 1/I const PARTY_GRACE_MS = 5 * 60_000; const CHAT_KEEP = 60; const MAX_MSG = 64 * 1024; export const PARTY_COMMANDS = new Set(['toggle', 'play', 'pause', 'seek', 'next', 'prev']); const clean = (v, max, fallback = '') => String(v || '').replace(/[\u0000-\u001f\u007f<>]+/g, ' ').trim().slice(0, max) || fallback; const same = (a, b) => { const x = Buffer.from(String(a)), y = Buffer.from(String(b)); return x.length === y.length && timingSafeEqual(x, y); }; export function createPartyHub({ now = () => Date.now() } = {}) { const parties = new Map(); // code → party const joinFails = new Map(); const send = (ws, m) => { try { ws.send(JSON.stringify(m)); } catch { /* gone */ } }; const members = (p) => [...p.members.values()].map((m) => ({ pid: m.pid, name: m.name, host: m.pid === p.hostPid, voice: !!m.voice })); const broadcast = (p, m, except) => { for (const x of p.members.values()) if (x.ws !== except) send(x.ws, m); }; function newCode() { let c; do { c = Array.from({ length: 6 }, () => ALPHABET[randomInt(ALPHABET.length)]).join(''); } while (parties.has(c)); return c; } const sweeper = setInterval(() => { const t = now(); for (const [code, p] of parties) { if (!p.hostOnline && t - p.hostLeftAt > PARTY_GRACE_MS) { broadcast(p, { type: 'ended', reason: 'The host left' }); for (const m of p.members.values()) { try { m.ws.close(4010, 'party ended'); } catch { /* gone */ } } parties.delete(code); } } }, 30_000); sweeper.unref?.(); function upgrade(req, server, ip) { const u = new URL(req.url); const q = (k) => u.searchParams.get(k) || ''; const role = q('role'); const pid = q('pid'); if (!/^[A-Za-z0-9_-]{8,40}$/.test(pid)) return new Response('bad pid', { status: 400 }); const name = clean(q('name'), 40, 'Guest'); const code = q('code').toUpperCase(); if (role === 'guest') { const t = now(); const fails = (joinFails.get(ip) || []).filter((x) => t - x < 10 * 60_000); if (fails.length >= 20) return new Response('too many attempts', { status: 429 }); if (!parties.has(code)) { fails.push(t); joinFails.set(ip, fails); if (joinFails.size > 5000) joinFails.clear(); return new Response('no such party', { status: 404 }); } } else if (role !== 'host') { return new Response('unknown role', { status: 400 }); } const ok = server.upgrade(req, { data: { hub: 'party', role, pid, name, code, secret: q('secret') } }); return ok ? undefined : new Response('upgrade failed', { status: 400 }); } function open(ws) { const d = ws.data; let p; if (d.role === 'host') { p = d.code && parties.get(d.code); if (p && !same(p.secret, d.secret)) { ws.close(4003, 'not the host of this party'); return; } if (!p) { const code = newCode(); p = { direct: createDirectRelay(), code, secret: randomBytes(18).toString('base64url'), hostPid: d.pid, members: new Map(), state: null, chat: [], dj: createDJQueue(), allowControl: false, hostOnline: true, hostLeftAt: 0 }; parties.set(code, p); } p.hostPid = d.pid; p.hostOnline = true; d.code = p.code; } else { p = parties.get(d.code); if (!p) { ws.close(4004, 'no such party'); return; } } const prev = p.members.get(d.pid); if (prev && prev.ws !== ws) { try { prev.ws.close(4000, 'replaced'); } catch { /* gone */ } } p.members.set(d.pid, { ws, pid: d.pid, name: d.name, voice: false, lastChat: 0 }); send(ws, { type: 'hello', code: p.code, secret: d.role === 'host' ? p.secret : undefined, you: d.pid, now: now(), members: members(p), state: p.state, chat: p.chat, dj: queueSnapshot(p.dj), allowControl: p.allowControl, }); broadcast(p, { type: 'members', members: members(p), joined: d.name }, ws); } function message(ws, raw) { if (typeof raw !== 'string' || raw.length > MAX_MSG) return; let m; try { m = JSON.parse(raw); } catch { return; } const d = ws.data; const p = parties.get(d.code); const me = p && p.members.get(d.pid); if (!me || me.ws !== ws || !m || typeof m !== 'object') return; const isHost = d.pid === p.hostPid; if (m.type === 'direct') { p.direct.handle(me.pid, m, new Map([...p.members].map(([pid, member]) => [pid, member.ws])), send); } else if (m.type === 'dj') { if (m.action === 'take' && isHost) { const video = takeDJVideo(p.dj); send(ws, { type: 'dj-play', video }); if (video) broadcast(p, { type: 'dj', ...queueSnapshot(p.dj) }); } else if (changeDJQueue(p.dj, m, me.pid, isHost)) broadcast(p, { type: 'dj', ...queueSnapshot(p.dj) }); } else if (m.type === 'state' && isHost) { p.state = { ...(m.state || {}), ts: now() }; broadcast(p, { type: 'state', state: p.state }, ws); } else if (m.type === 'settings' && isHost) { p.allowControl = !!m.allowControl; broadcast(p, { type: 'settings', allowControl: p.allowControl }); } else if (m.type === 'chat') { const t = now(); if (t - me.lastChat < 400) return; // flood guard me.lastChat = t; const text = clean(m.text, 500); if (!text) return; const msg = { from: me.name, pid: me.pid, text, at: t }; p.chat.push(msg); if (p.chat.length > CHAT_KEEP) p.chat.shift(); broadcast(p, { type: 'chat', msg }); } else if (m.type === 'voice') { me.voice = !!m.on; broadcast(p, { type: 'members', members: members(p) }); } else if (m.type === 'rtc' && typeof m.to === 'string') { const target = p.members.get(m.to); if (target) send(target.ws, { type: 'rtc', from: me.pid, data: m.data }); } else if (m.type === 'cmd' && !isHost && p.allowControl && PARTY_COMMANDS.has(m.cmd)) { const host = p.members.get(p.hostPid); const { type, ...args } = m; if (host) send(host.ws, { type: 'cmd', from: me.name, ...args }); } else if (m.type === 'end' && isHost) { broadcast(p, { type: 'ended', reason: 'The host ended the party' }); for (const x of p.members.values()) { try { x.ws.close(4010, 'party ended'); } catch { /* gone */ } } parties.delete(p.code); } } function close(ws) { const d = ws.data; const p = parties.get(d.code); if (!p) return; const me = p.members.get(d.pid); if (!me || me.ws !== ws) return; p.members.delete(d.pid); if (d.pid === p.hostPid) { p.hostOnline = false; p.hostLeftAt = now(); } broadcast(p, { type: 'members', members: members(p), left: d.name, hostAway: !p.hostOnline }); } return { upgrade, websocket: { maxPayloadLength: MAX_MSG, idleTimeout: 120, open, message: (ws, msg) => message(ws, typeof msg === 'string' ? msg : null), close }, _parties: parties, stop() { clearInterval(sweeper); }, }; }