151 lines
6.7 KiB
JavaScript
151 lines
6.7 KiB
JavaScript
"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 };
|