"use strict"; const { WebSocketServer, WebSocket } = require("ws"); const { isSecureRequest, requestOrigin } = require("../../../src/services/proxy-security"); const { insecureDeviceAllowed } = require("../../lumi_transcription/backend/companion/device_store"); const { onOverlayChatMessage } = require("../../../src/services/overlay-chat"); const { onLumiEvent } = require("../../../src/services/lumi-events"); const { parseClientMessage, sanitizeChat, sanitizeEvent, MAX_MESSAGE_BYTES } = require("./protocol"); const { listSources, listStatuses, sourceIdForChat, sourceIdForEvent, listEventTypes } = require("./sources"); const MAX_QUEUE = 300; const MAX_DEDUPE = 2000; class OverlayFeed { constructor({ devices, logger }) { this.devices = devices; this.logger = logger || { info() {}, warn() {}, error() {} }; this.wss = new WebSocketServer({ noServer: true, maxPayload: MAX_MESSAGE_BYTES, perMessageDeflate: false }); this.sessions = new Set(); this.publishedSeen = new Map(); this.sequence = 0; this.stopChat = onOverlayChatMessage((message) => this.publishChat(message)); this.stopEvent = onLumiEvent((event) => this.publishEvent(event)); this.wss.on("connection", (socket, request, device) => this.connected(socket, request, device)); } upgrade(request, socket, head) { const auth = this.devices.authenticate(request.headers.authorization, "overlay.read.v1"); if (!auth.allowed) return reject(socket, auth.reason === "capability_revoked" ? 403 : 401, auth.reason); let origin; try { origin = requestOrigin(request); } catch { return reject(socket, 400, "invalid_host"); } if (!isSecureRequest(request) && !insecureDeviceAllowed(auth.device, origin, request.socket.remoteAddress)) return reject(socket, 426, "tls_required"); this.wss.handleUpgrade(request, socket, head, (ws) => this.wss.emit("connection", ws, request, auth.device)); } connected(socket, _request, device) { const session = { socket, device, sources: new Set(), events: new Set(), subscribed: false, queue: [], draining: false, seen: new Map(), lastPong: Date.now() }; this.sessions.add(session); this.send(session, "hello", { protocol_version: 1, capabilities: ["overlay.feed.v1", "sources.catalog.v1", "live-state.v1"], sources: listSources(), event_types: listEventTypes(), statuses: listStatuses(), heartbeat_ms: 15000 }); const heartbeat = setInterval(() => { if (Date.now() - session.lastPong > 45000) return socket.close(4408, "heartbeat_timeout"); this.send(session, "ping", { at: Date.now() }); this.send(session, "catalog", { sources: listSources(), statuses: listStatuses() }); }, 15000); heartbeat.unref?.(); socket.on("message", (raw, binary) => { if (binary) return socket.close(4400, "text_messages_only"); try { const message = parseClientMessage(raw); if (message.type === "subscribe") { session.sources = new Set(message.sources); session.events = new Set(message.events); session.subscribed = true; this.send(session, "subscribed", { sources: [...session.sources], events: [...session.events] }); } else if (message.type === "pong") session.lastPong = Date.now(); else if (message.type === "ping") this.send(session, "pong", { at: Date.now() }); } catch (error) { this.send(session, "error", { code: error.code || "INVALID_MESSAGE", message: error.message }); } }); socket.on("close", () => { clearInterval(heartbeat); this.sessions.delete(session); }); socket.on("error", (error) => this.logger.warn?.("Native overlay feed connection failed", { device_id: device.id, error: error.message }, { event: "overlay_feed_socket_error" })); } publishChat(message) { const sourceId = sourceIdForChat(message); if (!this.acceptPublished(`chat:${message.id}`)) return; for (const session of this.sessions) { if (!session.subscribed || !session.sources.has(sourceId)) continue; this.enqueue(session, "chat", { source_id: sourceId, message: sanitizeChat(message) }, `chat:${message.id}`); } } publishEvent(event) { const sourceId = sourceIdForEvent(event); if (!this.acceptPublished(`event:${event.id}`)) return; for (const session of this.sessions) { if (!session.subscribed || !session.events.has(event.type) || !matchesSelectedSource(session.sources, sourceId)) continue; this.enqueue(session, "event", { source_id: sourceId, event: sanitizeEvent(event) }, `event:${event.id}`); } } enqueue(session, type, payload, key) { const now = Date.now(); if (session.seen.has(key)) return; session.seen.set(key, now); while (session.seen.size > MAX_DEDUPE) session.seen.delete(session.seen.keys().next().value); const envelope = { type, sequence: ++this.sequence, sent_at: now, ...payload }; if (session.queue.length >= MAX_QUEUE) session.queue.shift(); session.queue.push(envelope); this.drain(session); } acceptPublished(key) { this.publishedSeen ||= new Map(); if (this.publishedSeen.has(key)) return false; this.publishedSeen.set(key, Date.now()); while (this.publishedSeen.size > MAX_DEDUPE) this.publishedSeen.delete(this.publishedSeen.keys().next().value); return true; } drain(session) { if (session.draining) return; session.draining = true; try { while (session.queue.length && session.socket.readyState === WebSocket.OPEN) { session.socket.send(JSON.stringify(session.queue.shift())); } } finally { session.draining = false; } } send(session, type, payload) { if (session.socket.readyState === WebSocket.OPEN) session.socket.send(JSON.stringify({ type, sequence: ++this.sequence, sent_at: Date.now(), ...payload })); } async close() { this.stopChat(); this.stopEvent(); for (const session of this.sessions) session.socket.close(1001, "plugin_shutdown"); await new Promise((resolve) => this.wss.close(resolve)); } } function matchesSelectedSource(selected, sourceId) { if (selected.has(sourceId)) return true; if (sourceId.startsWith("discord:") && sourceId.endsWith(":*")) return [...selected].some((id) => id.startsWith(sourceId.slice(0, -1))); if (sourceId.endsWith(":current")) return [...selected].some((id) => id.startsWith(sourceId.split(":")[0] + ":")); return false; } function reject(socket, status, reason) { const labels = { 400: "Bad Request", 401: "Unauthorized", 403: "Forbidden", 426: "Upgrade Required" }; socket.write(`HTTP/1.1 ${status} ${labels[status] || "Rejected"}\r\nConnection: close\r\nContent-Type: application/json\r\n\r\n${JSON.stringify({ error: reason })}`); socket.destroy(); } module.exports = { OverlayFeed, MAX_QUEUE, matchesSelectedSource };