178 lines
8.2 KiB
JavaScript
178 lines
8.2 KiB
JavaScript
/* ============================================================================
|
|
* 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); },
|
|
};
|
|
}
|