Something went wrong. Try again.
Monorepo for Aesthetic.Computer aesthetic.computer
Something went wrong. Try again.
43 kB · 1239 lines
JavaScript
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240// Chat Manager, 25.11.28// Multi-instance chat support for session-server// Adapted from nanos/chat.mjs to run multiple chat instances in one process
import { WebSocket } from "ws";import { promises as fs } from "fs";import fetch from "node-fetch";import https from "https";
import { filter } from "./filter.mjs";import { redact, unredact } from "./redact.mjs";import { ensureIndexes as ensureHeartsIndexes, toggleHeart, countHearts } from "./hearts.mjs";
import { MongoClient, ObjectId } from "mongodb";import { broadcastToTopic } from "../shared/push.mjs"; // Standard push (no Firebase).import { MAX_CHARS, chatCapabilities, profanityFiltered,} from "../shared/chat-capabilities.mjs";
const MAX_MESSAGES = 500;
// A dead or unreachable Mongo should cost a chat connection seconds, not the// driver's default 30s of server selection — several awaits sit on the path// between a socket opening and its `connected` handshake.const mongoClientOptions = { serverSelectionTimeoutMS: 8000 };
// Chat instance configurationsexport const chatInstances = { "chat-system.aesthetic.computer": { name: "chat-system", allowedHost: "chat-system.aesthetic.computer", userInfoEndpoint: "https://aesthetic.us.auth0.com/userinfo", topic: "mood", }, "chat.sotce.net": { name: "chat-sotce", allowedHost: "chat.sotce.net", userInfoEndpoint: "https://sotce.us.auth0.com/userinfo", // Dedicated topic — sotce chat pings must never fan out to the shared // "mood" topic that aesthetic.computer devices subscribe to. topic: "chat-sotce", // And its own face in the notification, not the aesthetic.computer mark. icon: "https://assets.aesthetic.computer/sotce-net/cookie-192.png", }, "chat-clock.aesthetic.computer": { name: "chat-clock", allowedHost: "chat-clock.aesthetic.computer", userInfoEndpoint: "https://aesthetic.us.auth0.com/userinfo", topic: "mood", },};
// Development host mappings (localhost ports)const devHostMappings = { "localhost:8083": "chat-system.aesthetic.computer", "localhost:8084": "chat.sotce.net", "localhost:8085": "chat-clock.aesthetic.computer",};
export class ChatManager { constructor(options = {}) { this.dev = options.dev || false; this.filterDebug = options.filterDebug || false; this.loggerKey = options.loggerKey || process.env.LOGGER_KEY; this.activityEmitter = options.activityEmitter || null; // MongoDB this.mongoClient = null; this.db = null; this.mongoConnectionString = options.mongoConnectionString || process.env.MONGODB_CONNECTION_STRING; this.mongoDbName = options.mongoDbName || process.env.MONGODB_NAME || "aesthetic"; // HTTPS agent for dev mode this.agent = this.dev ? new https.Agent({ rejectUnauthorized: false }) : null; // Per-instance state this.instances = {}; for (const [host, config] of Object.entries(chatInstances)) { this.instances[host] = { config, messages: [], connections: {}, connectionId: 0, authorizedConnections: {}, subsToHandles: {}, subsToSubscribers: {}, }; } console.log("💬 ChatManager initialized with instances:", Object.keys(chatInstances).join(", ")); }
setActivityEmitter(emitter) { this.activityEmitter = typeof emitter === "function" ? emitter : null; }
emitActivity(payload) { if (!this.activityEmitter) return; try { this.activityEmitter(payload); } catch (err) { console.error("💬 Activity emitter failure:", err); } }
async init() { // Connect to MongoDB if (this.mongoConnectionString) { try { console.log("💬 Connecting to MongoDB..."); this.mongoClient = new MongoClient(this.mongoConnectionString, mongoClientOptions); await this.mongoClient.connect(); this.db = this.mongoClient.db(this.mongoDbName); console.log("💬 MongoDB connected!");
// Ensure hearts indexes exist await ensureHeartsIndexes(this.db);
// Load messages for each instance for (const [host, instance] of Object.entries(this.instances)) { await this.loadMessages(instance); } } catch (err) { console.error("💬 MongoDB connection failed:", err); } } else { console.log("💬 No MongoDB connection string, running without persistence"); }
// Chat must outlive its database: messages that fail to persist stay in // memory (`sub`, no `id`) and this flush retries them until Mongo returns. // Without it, recovery only happened at shutdown — a whole outage's worth // of messages rode on the process exiting cleanly. if (!this.dev && this.mongoConnectionString) { this.persistTimer = setInterval(async () => { const stranded = Object.values(this.instances).reduce( (n, instance) => n + instance.messages.filter((m) => m.sub && !m.id).length, 0, ); if (stranded === 0) return; console.log(`💬 Retrying persistence for ${stranded} stranded message(s)...`); try { await this.persistAllMessages(); } catch (err) { console.error("💬 Persistence retry failed:", err.message); } }, 60000); } }
async loadMessages(instance) { if (!this.db) return; const collectionName = instance.config.name; console.log(`💬 Loading messages for ${collectionName}...`); try { const chatCollection = this.db.collection(collectionName); let combinedMessages;
if (collectionName === "chat-sotce") { combinedMessages = (await chatCollection .find({}) .sort({ when: -1 }) .limit(MAX_MESSAGES) .toArray()).reverse(); } else if (collectionName !== "chat-system") { combinedMessages = (await chatCollection .find({}) .sort({ when: -1 }) .limit(MAX_MESSAGES) .toArray()).reverse(); } else { // chat-system includes logs combinedMessages = ( await chatCollection .aggregate([ { $unionWith: { coll: "logs", pipeline: [{ $match: {} }], }, }, { $sort: { when: -1 } }, { $limit: MAX_MESSAGES }, ]) .toArray() ).reverse(); }
// Deduplicate messages (removes duplicates caused by persist-on-shutdown bug) { const seen = new Set(); combinedMessages = combinedMessages.filter((msg) => { const key = `${msg.user || msg.from}:${msg.text}:${msg.when?.getTime?.() ?? msg.when}`; if (seen.has(key)) return false; seen.add(key); return true; }); }
// Check mutes for chat-system if (collectionName === "chat-system") { for (const msg of combinedMessages) { if (await this.isMuted(instance, msg.user)) { redact(msg); } } }
instance.messages = []; for (const message of combinedMessages) { let from; if (message.user) { from = await this.getHandleFromSub(instance, message.user); } else { from = message.from || "deleted"; }
const msg = { from, text: message.deleted ? "[deleted]" : (instance.config.name === "chat-clock" ? message.text : filter(message.text, this.filterDebug)) || "message forgotten", redactedText: message.redactedText, when: message.when, sub: message.user || undefined, font: message.font || "font_1", // 🔤 Include font from DB (default for old messages) }; if (message.via) msg.via = message.via; // 📺 e.g. "youtube" — service relays if (message.link) msg.link = message.link; // 📺 walk-back to the live chat if (message._id) msg.id = message._id.toString(); if (message.deleted) msg.deleted = true; instance.messages.push(msg); } console.log(`💬 Loaded ${instance.messages.length} messages for ${collectionName}`); } catch (err) { console.error(`💬 Failed to load messages for ${collectionName}:`, err); } }
// Check if a host should be handled by chat manager isChatHost(host) { // Direct match if (chatInstances[host]) return true; // Dev mapping if (this.dev && devHostMappings[host]) return true; return false; }
// Get the instance for a host getInstance(host) { // Direct match if (this.instances[host]) return this.instances[host]; // Dev mapping if (this.dev && devHostMappings[host]) { return this.instances[devHostMappings[host]]; } return null; }
// Handle a new WebSocket connection async handleConnection(ws, req) { const host = req.headers.host; const instance = this.getInstance(host); if (!instance) { console.log("💬 Unknown chat host:", host); ws.close(1008, "Unknown host"); return; }
// Validate host in production if (!this.dev && host !== instance.config.allowedHost) { ws.close(1008, "Policy violation"); return; }
const id = instance.connectionId++; instance.connections[id] = ws; const ip = req.socket.remoteAddress || "localhost"; ws.isAlive = true;
ws.on("pong", () => { ws.isAlive = true; });
console.log( `💬 [${instance.config.name}] Connection ${id} from ${ip}, online: ${Object.keys(instance.connections).length}` );
// Set up ping interval for this connection const pingInterval = setInterval(() => { if (ws.isAlive === false) { clearInterval(pingInterval); return ws.terminate(); } ws.isAlive = false; ws.ping(); }, 15000);
ws.on("message", async (data) => { await this.handleMessage(instance, ws, id, data); });
ws.on("close", () => { clearInterval(pingInterval); delete instance.connections[id]; delete instance.authorizedConnections[id];
console.log( `💬 [${instance.config.name}] Connection ${id} closed, online: ${Object.keys(instance.connections).length}` );
// Broadcast updated online handles after removing this connection this.broadcastOnlineHandles(instance); this.broadcast(instance, this.pack("left", { chatters: Object.keys(instance.connections).length, handles: this.getOnlineHandles(instance) }, id)); });
// Load heart counts for current message window. This await sits between // the socket opening and the `connected` packet — if it throws, the client // never finishes its handshake and chat reads as offline (the 2026-08-12 // credential-rotation outage). Hearts are decoration; degrade to none. let heartCounts = {}; if (this.db) { try { const messageIds = instance.messages.filter((m) => m.id).map((m) => m.id); heartCounts = await countHearts(this.db, instance.config.name, messageIds); } catch (err) { console.error(`💬 [${instance.config.name}] Heart counts unavailable, connecting without:`, err.message); } }
// Send welcome message ws.send( this.pack( "connected", { message: `Joined \`${instance.config.name}\` • 🧑🤝🧑 ${Object.keys(instance.connections).length}`, chatters: Object.keys(instance.connections).length, handles: this.getOnlineHandles(instance), messages: instance.messages, heartCounts, // The handshake: what this channel accepts. Clients read their limit // and syntax from here instead of hardcoding a copy of it. capabilities: chatCapabilities(instance.config.name), id, }, id, ), );
// Notify others this.broadcastOthers(instance, ws, this.pack( "joined", { text: `${id} has joined. Connections open: ${Object.keys(instance.connections).length}`, chatters: Object.keys(instance.connections).length, handles: this.getOnlineHandles(instance), }, id, )); }
async handleMessage(instance, ws, id, data) { let msg; try { msg = JSON.parse(data.toString()); } catch (err) { console.log("💬 Failed to parse message:", err); return; }
msg.id = id;
if (msg.type === "logout") { console.log(`💬 [${instance.config.name}] User logged out`); delete instance.authorizedConnections[id]; } else if (msg.type === "chat:message") { await this.handleChatMessage(instance, ws, id, msg); } else if (msg.type === "chat:service-message") { await this.handleServiceMessage(instance, ws, id, msg); } else if (msg.type === "chat:delete") { await this.handleDeleteMessage(instance, ws, id, msg); } else if (msg.type === "chat:edit") { await this.handleEditMessage(instance, ws, id, msg); } else if (msg.type === "chat:heart") { await this.handleChatHeart(instance, ws, id, msg); } }
// 📺 Service relays (the TV chat bridge) post on behalf of off-AC visitors — // e.g. YouTube live-chat viewers of the always-on broadcasts. They authorize // with a shared secret (CHAT_SERVICE_SECRET env), not a user token, and their // messages carry `via` (e.g. "youtube") + the visitor's display name as // `from`, so clients can render them as televised guests rather than handles. // No activity feed, no push notifications — guests are heard, not amplified. async handleServiceMessage(instance, ws, id, msg) { const secret = process.env.CHAT_SERVICE_SECRET; if (!secret || msg.content?.secret !== secret) { console.error(`💬 [${instance.config.name}] Service message rejected (bad secret)`); ws.send(this.pack("unauthorized", { message: "Bad service secret." }, id)); return; }
const via = String(msg.content.via || "youtube").slice(0, 16); // Visitor names are foreign input: strip control chars and our color-code // delimiter, clamp, and fall back rather than ever posting an empty name. const name = String(msg.content.name || "") .replace(/[\u0000-\u001f\\@]/g, "") .trim() .slice(0, 24) || "viewer"; let text = String(msg.content.text || "").trim(); if (!text) return; if (text.length > MAX_CHARS) text = text.slice(0, MAX_CHARS); if (profanityFiltered(instance.config.name)) text = filter(text, this.filterDebug);
// Optional walk-back link (the broadcast's popout chat) — YouTube only. let link; const rawLink = String(msg.content.link || "").slice(0, 160); if (/^https:\/\/(www\.youtube\.com|youtu\.be)\//.test(rawLink)) link = rawLink;
let when = new Date(); if (!this.dev) { try { const clockResponse = await fetch("https://aesthetic.computer/api/clock"); if (clockResponse.ok) when = new Date(await clockResponse.text()); } catch (err) { console.log("💬 Clock fetch failed, using local time"); } }
let insertedId; if (!this.dev && this.db) { try { const collection = this.db.collection(instance.config.name); const result = await collection.insertOne({ from: name, via, text, when, font: "font_1", ...(link ? { link } : {}) }); insertedId = result.insertedId?.toString(); } catch (err) { console.error(`💬 [${instance.config.name}] Service message store failed, broadcasting anyway:`, err.message); } }
const out = { from: name, via, text, when, font: "font_1" }; if (link) out.link = link; if (insertedId) out.id = insertedId; instance.messages.push(out); if (instance.messages.length > MAX_MESSAGES) instance.messages.shift(); this.broadcast(instance, this.pack("message", out)); console.log(`💬 [${instance.config.name}] 📺 via ${via}: ${name} (${text.length} chars)`); }
async handleChatMessage(instance, ws, id, msg) { console.log( `💬 [${instance.config.name}] Message from ${msg.content.sub} (${msg.content.text.length} chars)` );
// Mute check if (await this.isMuted(instance, msg.content.sub)) { ws.send(this.pack("muted", { message: "Your user has been muted." })); return; }
// Length limit. Counted in UTF-16 code units — see shared/chat-capabilities. if (msg.content.text.length > MAX_CHARS) { ws.send( this.pack("too-long", { message: `Please limit to ${MAX_CHARS} characters.`, maxChars: MAX_CHARS, countedAs: "utf16-code-units", was: msg.content.text.length, }), ); return; }
// Authorization (Auth0 token or AC device token) let authorized; if (instance.authorizedConnections[id]?.token === msg.content.token) { authorized = instance.authorizedConnections[id].user; console.log("💬 Pre-authorized"); } else { console.log("💬 Authorizing...", "handle:", msg.content.handle, "token:", msg.content.token?.slice(0, 10) + "...", "content-type:", typeof msg.content); authorized = await this.authorize(instance, msg.content.token); // Fallback: AC device token (from ac-native device-token API) if (!authorized && msg.content.handle && msg.content.token) { console.log("💬 Trying device token auth for @" + msg.content.handle); authorized = await this.authorizeDeviceToken(instance, msg.content.handle, msg.content.token); } if (authorized) { instance.authorizedConnections[id] = { token: msg.content.token, user: authorized, }; } }
// Get handle let handle, subscribed; if (authorized) { handle = await this.getHandleFromSub(instance, authorized.sub); // Store handle in authorizedConnections for online list instance.authorizedConnections[id].handle = handle; // Broadcast updated online handles to all clients this.broadcastOnlineHandles(instance); // Subscription check for sotce if (instance.config.name === "chat-sotce") { subscribed = await this.checkSubscription(instance, authorized.sub, msg.content.token); } else { subscribed = true; } }
if (!authorized || !handle || !subscribed) { console.error("💬 Unauthorized:", { authorized: !!authorized, handle, subscribed }); ws.send(this.pack("unauthorized", { message: "Please login and/or subscribe." }, id)); return; }
// Process and store message try { const message = msg.content; const fromSub = message.sub; let filteredText; const userIsMuted = await this.isMuted(instance, fromSub);
if (userIsMuted) { redact(message); filteredText = message.text; } else { filteredText = profanityFiltered(instance.config.name) ? filter(message.text, this.filterDebug) : message.text; }
// Get server time let when = new Date(); if (!this.dev) { try { const clockResponse = await fetch("https://aesthetic.computer/api/clock"); if (clockResponse.ok) { const serverTimeISO = await clockResponse.text(); when = new Date(serverTimeISO); } } catch (err) { console.log("💬 Clock fetch failed, using local time"); } }
// Store in MongoDB (production only, non-muted users). Persistence is // best-effort: a failure here must not stop the broadcast below — the // message stays in memory with `sub` and no `id`, which is exactly what // the periodic persistAllMessages flush looks for once Mongo recovers. let insertedId; if (!this.dev && !userIsMuted && this.db) { try { const dbmsg = { user: message.sub, text: message.text, when, font: message.font || "font_1", // 🔤 Store user's font preference }; const collection = this.db.collection(instance.config.name); await collection.createIndex({ when: 1 }); const result = await collection.insertOne(dbmsg); insertedId = result.insertedId?.toString(); console.log("💬 Message stored"); } catch (err) { console.error(`💬 [${instance.config.name}] Message store failed, broadcasting anyway:`, err.message); } }
const out = { from: handle, text: filteredText, redactedText: message.redactedText, when, sub: fromSub, font: message.font || "font_1", // 🔤 Include font in broadcast }; if (insertedId) out.id = insertedId;
// Duplicate detection const lastMsg = instance.messages[instance.messages.length - 1]; if (lastMsg && lastMsg.sub === out.sub && lastMsg.text === out.text && !lastMsg.count) { lastMsg.count = (lastMsg.count || 1) + 1; this.broadcast(instance, this.pack("message:update", { index: instance.messages.length - 1, count: lastMsg.count })); } else { instance.messages.push(out); if (instance.messages.length > MAX_MESSAGES) instance.messages.shift(); this.broadcast(instance, this.pack("message", out)); }
if (handle && handle.startsWith("@")) { const compactText = `${filteredText || ""}`.replace(/\s+/g, " ").trim(); const preview = compactText.length > 80 ? `${compactText.slice(0, 77)}...` : compactText;
this.emitActivity({ handle, event: { type: "chat", when: when?.getTime?.() || Date.now(), label: `Chat ${instance.config.name.replace("chat-", "")}: ${preview}`, piece: instance.config.name, ref: insertedId || null, text: compactText, }, countsDelta: { chats: 1 }, }); }
// Push notification (production only, non-muted). chat-sotce rides its // own "chat-sotce" topic, which only sotce.net readers opt into. if (!this.dev && !userIsMuted) { this.notify(instance, handle, filteredText, when); } } catch (err) { console.error("💬 Message handling error:", err); } }
async handleDeleteMessage(instance, ws, id, msg) { const { token, sub, id: messageId } = msg.content; console.log( `💬 [${instance.config.name}] Delete request from ${sub} for message ${messageId}` );
if (!messageId) { ws.send(this.pack("error", { message: "Missing message id." })); return; }
// Authorize the user let authorized; if (instance.authorizedConnections[id]?.token === token) { authorized = instance.authorizedConnections[id].user; } else { authorized = await this.authorize(instance, token); if (authorized) { instance.authorizedConnections[id] = { token, user: authorized }; } }
if (!authorized || authorized.sub !== sub) { ws.send(this.pack("unauthorized", { message: "Please login." }, id)); return; }
// Find the message in memory and verify ownership const message = instance.messages.find((m) => m.id === messageId); if (!message) { ws.send(this.pack("error", { message: "Message not found." })); return; } if (message.sub !== sub) { ws.send(this.pack("error", { message: "You can only delete your own messages." })); return; }
// Already deleted if (message.deleted) return;
// Soft-delete in MongoDB if (this.db) { try { const collection = this.db.collection(instance.config.name); await collection.updateOne( { _id: new ObjectId(messageId) }, { $set: { deleted: true } } ); console.log("💬 Message soft-deleted in DB"); } catch (err) { console.error("💬 Failed to delete message in DB:", err); } }
// Update in-memory message.text = "[deleted]"; message.deleted = true; delete message.redactedText;
// Broadcast deletion to all clients this.broadcast(instance, this.pack("message:delete", { id: messageId })); }
// ✏️ Re-edit your own message in place. Mirrors delete's auth + ownership // checks and message-send's length/profanity rules; DB keeps the raw text // (like handleChatMessage) while the broadcast carries the filtered copy. async handleEditMessage(instance, ws, id, msg) { const { token, sub, id: messageId } = msg.content || {}; let text = msg.content?.text; console.log( `💬 [${instance.config.name}] Edit request from ${sub} for message ${messageId}` );
if (!messageId || typeof text !== "string") { ws.send(this.pack("error", { message: "Missing message id or text." })); return; }
text = text.replace(/\s+$/, ""); if (text.length === 0) { ws.send(this.pack("error", { message: "Nothing to say." })); return; } if (text.length > MAX_CHARS) { ws.send( this.pack("too-long", { message: `Please limit to ${MAX_CHARS} characters.`, maxChars: MAX_CHARS, countedAs: "utf16-code-units", was: text.length, }), ); return; }
if (await this.isMuted(instance, sub)) { ws.send(this.pack("muted", { message: "Your user has been muted." })); return; }
// Authorize the user let authorized; if (instance.authorizedConnections[id]?.token === token) { authorized = instance.authorizedConnections[id].user; } else { authorized = await this.authorize(instance, token); if (authorized) { instance.authorizedConnections[id] = { token, user: authorized }; } }
if (!authorized || authorized.sub !== sub) { ws.send(this.pack("unauthorized", { message: "Please login." }, id)); return; }
// Find the message in memory and verify ownership const message = instance.messages.find((m) => m.id === messageId); if (!message) { ws.send(this.pack("error", { message: "Message not found." })); return; } if (message.sub !== sub) { ws.send(this.pack("error", { message: "You can only edit your own messages." })); return; } if (message.deleted) { ws.send(this.pack("error", { message: "Deleted messages can't be edited." })); return; }
const filteredText = profanityFiltered(instance.config.name) ? filter(text, this.filterDebug) : text; const editedWhen = new Date();
// Update MongoDB (raw text, matching how sends are stored) if (this.db) { try { const collection = this.db.collection(instance.config.name); await collection.updateOne( { _id: new ObjectId(messageId) }, { $set: { text, edited: true, editedWhen } } ); console.log("💬 Message edited in DB"); } catch (err) { console.error("💬 Failed to edit message in DB:", err); } }
// Update in-memory message.text = filteredText; message.edited = true; message.editedWhen = editedWhen; delete message.redactedText;
// Broadcast the edit to all clients this.broadcast( instance, this.pack("message:edit", { id: messageId, text: filteredText, edited: true, editedWhen, }), ); }
async handleChatHeart(instance, ws, id, msg) { const { for: forId, token } = msg.content || {}; if (!forId || !token) return;
// Reuse cached auth or re-authorize let authorized; if (instance.authorizedConnections[id]?.token === token) { authorized = instance.authorizedConnections[id].user; } else { authorized = await this.authorize(instance, token); if (authorized) { instance.authorizedConnections[id] = { token, user: authorized }; } }
if (!authorized) { ws.send(this.pack("unauthorized", { message: "Please login." }, id)); return; }
if (!this.db) return;
try { const { hearted, count } = await toggleHeart(this.db, { user: authorized.sub, type: instance.config.name, for: forId, }); console.log(`💬 [${instance.config.name}] Heart ${hearted ? "+" : "-"} on ${forId} → ${count}`); this.broadcast(instance, this.pack("message:hearts", { for: forId, count })); } catch (err) { console.error("💬 Heart error:", err); } }
async authorizeDeviceToken(instance, handle, token) { // AC device tokens are "hmac.timestamp" — validate by looking up handle in MongoDB if (!token || !token.includes(".") || !handle) return undefined; try { const cleanHandle = handle.replace("@", ""); const doc = await this.db.collection("@handles").findOne({ handle: cleanHandle }); if (doc && doc._id) { console.log("💬 Device token authorized for @" + cleanHandle + " (sub: " + doc._id + ")"); return { sub: doc._id }; } console.log("💬 Device token: handle @" + cleanHandle + " not found in DB"); return undefined; } catch (err) { console.error("💬 Device token auth error:", err); return undefined; } }
async authorize(instance, token) { try { const response = await fetch(instance.config.userInfoEndpoint, { headers: { Authorization: "Bearer " + token, "Content-Type": "application/json", }, });
if (response.status === 200) { return response.json(); } return undefined; } catch (err) { console.error("💬 Authorization error:", err); return undefined; } }
async getHandleFromSub(instance, fromSub) { if (instance.subsToHandles[fromSub]) { return "@" + instance.subsToHandles[fromSub]; }
try { let prefix = instance.config.name === "chat-sotce" ? "sotce-" : ""; let host = this.dev ? "https://localhost:8888" : (instance.config.name === "chat-sotce" ? "https://sotce.net" : "https://aesthetic.computer");
const options = {}; if (this.dev) options.agent = this.agent;
const response = await fetch(`${host}/handle?for=${prefix}${fromSub}`, options); if (response.status === 200) { const data = await response.json(); if (data.handle) { instance.subsToHandles[fromSub] = data.handle; return "@" + data.handle; } } } catch (err) { console.error("💬 Handle lookup error:", err); }
return "@unknown"; }
async checkSubscription(instance, sub, token) { if (instance.subsToSubscribers[sub] !== undefined) { return instance.subsToSubscribers[sub]; }
const host = this.dev ? "https://localhost:8888" : "https://sotce.net"; const options = { method: "POST", body: JSON.stringify({ retrieve: "subscription" }), headers: { Authorization: "Bearer " + token, "Content-Type": "application/json", }, };
if (this.dev) options.agent = this.agent;
try { const response = await fetch(`${host}/sotce-net/subscribed`, options); if (response.status === 200) { const data = await response.json(); instance.subsToSubscribers[sub] = data.subscribed; return data.subscribed; } } catch (err) { console.error("💬 Subscription check error:", err); }
instance.subsToSubscribers[sub] = false; return false; }
async isMuted(instance, sub) { if (!sub || !this.db) return false; try { const mutesCollection = this.db.collection(instance.config.name + "-mutes"); const mute = await mutesCollection.findOne({ user: sub }); return !!mute; } catch (err) { return false; } }
notify(instance, handle, text, when) { let title = handle + " 💬"; if (instance.config.name === "chat-clock") { const getClockEmoji = (date) => { let hour = date.getHours() % 12 || 12; const minutes = date.getMinutes(); const emojiCode = minutes < 30 ? (0x1f550 + hour - 1) : (0x1f55c + hour - 1); return String.fromCodePoint(emojiCode); }; title = handle + " " + getClockEmoji(when); }
if (!this.db) return; broadcastToTopic(this.db, instance.config.topic, { title, body: text, icon: instance.config.icon, // undefined keeps the default mark urgent: true, // time-sensitive on iOS, Urgency: high on web ttl: 0, // don't store undelivered chat pings data: { piece: "chat" }, }).then((summary) => { console.log("💬 Notification sent:", summary); }).catch((err) => { console.log("💬 Notification error:", err); }); }
// Handle HTTP log endpoint async handleLog(instance, body, authHeader) { const token = authHeader?.split(" ")[1]; if (token !== this.loggerKey) { return { status: 403, body: { status: "error", message: "Forbidden" } }; }
try { const parsed = typeof body === "string" ? JSON.parse(body) : body; console.log(`💬 [${instance.config.name}] Log received from: ${parsed.from || 'unknown'}`);
instance.messages.push(parsed); if (instance.messages.length > MAX_MESSAGES) instance.messages.shift();
// Handle actions (mute/unmute, handle updates) if (parsed.action) { await this.handleLogAction(instance, parsed); }
this.broadcast(instance, this.pack("message", parsed));
if (instance.config.name === "chat-system") { this.notify(instance, "log 🪵", parsed.text, new Date()); }
return { status: 200, body: { status: "success", message: "Log received" } }; } catch (err) { return { status: 400, body: { status: "error", message: "Malformed log JSON" } }; } }
async handleLogAction(instance, parsed) { let [object, behavior] = (parsed.action || "").split(":"); if (!behavior) { behavior = object; object = null; }
if (object === "chat-system" && (behavior === "mute" || behavior === "unmute")) { const user = parsed.users[0]; if (behavior === "mute") { instance.messages.forEach((msg) => { if (msg.sub === user) redact(msg); }); } else { instance.messages.forEach((msg) => { if (msg.sub === user) unredact(msg); }); } this.broadcast(instance, this.pack(parsed.action, { user })); }
if (object === "handle") { if (behavior === "colors") { // Broadcast handle color changes to all connected clients. const data = JSON.parse(parsed.value); this.broadcast(instance, this.pack("handle:colors", { user: parsed.users[0], handle: data.handle, colors: data.colors })); } else { instance.subsToHandles[parsed.users[0]] = parsed.value;
if (behavior === "update" || behavior === "strip") { const from = behavior === "update" ? "@" + parsed.value : "nohandle"; instance.messages.forEach((msg) => { if (msg.sub === parsed.users[0]) { msg.from = from; } }); this.broadcast(instance, this.pack(parsed.action, { user: parsed.users[0], handle: from })); } } } }
pack(type, content, id) { if (typeof content === "object") content = JSON.stringify(content); return JSON.stringify({ type, content, id }); }
broadcast(instance, message) { Object.values(instance.connections).forEach((c) => { if (c?.readyState === WebSocket.OPEN) c.send(message); }); }
broadcastOthers(instance, exclude, message) { Object.values(instance.connections).forEach((c) => { if (c !== exclude && c?.readyState === WebSocket.OPEN) c.send(message); }); }
// Set the presence resolver function (called from session.mjs) // This allows us to query which users are on a specific piece setPresenceResolver(resolver) { this.presenceResolver = resolver; }
// Get list of online handles for an instance (all connected & authorized) getOnlineHandles(instance) { const handles = []; for (const [id, auth] of Object.entries(instance.authorizedConnections)) { if (auth.handle && instance.connections[id]) { handles.push(auth.handle); } } return [...new Set(handles)]; // Remove duplicates }
// Get handles of users actually viewing the chat piece right now getHereHandles(instance) { if (!this.presenceResolver) return []; // Get all users on "chat" piece from session server const onChatPiece = this.presenceResolver("chat"); // Intersect with authorized handles for this chat instance const onlineHandles = this.getOnlineHandles(instance); return onChatPiece.filter(h => onlineHandles.includes(h)); }
// Broadcast presence data (online + here) to all clients broadcastOnlineHandles(instance) { const online = this.getOnlineHandles(instance); const here = this.getHereHandles(instance); // Send both for backwards compatibility and new "here" feature this.broadcast(instance, this.pack("presence", { handles: online, // backwards compat with old "online-handles" online, // all authorized chat connections here // only those actually on chat piece })); }
// Get status for a specific instance or all getStatus(host = null) { if (host) { const instance = this.getInstance(host); if (!instance) return null; return { name: instance.config.name, connections: Object.keys(instance.connections).length, messages: instance.messages.length, }; }
return Object.entries(this.instances).map(([host, instance]) => ({ host, name: instance.config.name, connections: Object.keys(instance.connections).length, messages: instance.messages.length, })); }
// The same handshake the `connected` packet carries, for callers with no // socket. `host` may be a chat host or a bare channel name; without one you // get every channel, since the profanity policy differs between them. getCapabilities(host = null) { if (host) { const instance = this.getInstance(host) || Object.values(this.instances).find((i) => i.config.name === host); if (!instance) return null; return { host: instance.config.allowedHost, ...chatCapabilities(instance.config.name), }; }
return Object.values(this.instances).map((instance) => ({ host: instance.config.allowedHost, ...chatCapabilities(instance.config.name), })); }
// Get recent messages for a specific instance (for dashboard) getRecentMessages(host, count = 10) { const instance = this.getInstance(host); if (!instance) return []; // Return the most recent messages, but don't expose sensitive data return instance.messages.slice(-count).map(msg => ({ from: msg.from || 'unknown', text: msg.text || '', when: msg.when })).reverse(); // Most recent first }
// Persist any in-memory messages that weren't saved to MongoDB. // Messages received while this.db was null (broken env) have a `sub` field // but were never inserted. Messages loaded from DB on startup do not have `sub`. async persistAllMessages() { // Establish a connection if we don't have one if (!this.db && this.mongoConnectionString) { try { console.log("💬 Connecting to MongoDB for message persistence..."); this.mongoClient = new MongoClient(this.mongoConnectionString, mongoClientOptions); await this.mongoClient.connect(); this.db = this.mongoClient.db(this.mongoDbName); console.log("💬 MongoDB connected for persistence!"); } catch (err) { console.error("💬 Failed to connect to MongoDB for persistence:", err); return 0; } }
if (!this.db) { console.error("💬 No MongoDB connection available, cannot persist messages"); return 0; }
let totalPersisted = 0;
for (const [, instance] of Object.entries(this.instances)) { const collectionName = instance.config.name; // Only persist messages that were never saved to DB. // DB-loaded and successfully-stored messages have an `id` (from MongoDB _id). // Messages received while DB was unavailable have `sub` but no `id`. const unpersisted = instance.messages.filter(msg => msg.sub && !msg.id);
if (unpersisted.length === 0) { console.log(`💬 [${collectionName}] No unpersisted messages`); continue; }
try { const collection = this.db.collection(collectionName); const docs = unpersisted.map(msg => ({ user: msg.sub, text: msg.text, when: msg.when, font: msg.font || "font_1", }));
const result = await collection.insertMany(docs, { ordered: false }); totalPersisted += result.insertedCount; // Stamp each message with its new id so it stops matching the // unpersisted filter — this runs on a timer now, and without the // stamp every pass would insert the same messages again. for (const [i, insertedId] of Object.entries(result.insertedIds)) { if (unpersisted[i]) unpersisted[i].id = insertedId.toString(); } console.log(`💬 [${collectionName}] Persisted ${result.insertedCount} messages`); } catch (err) { // insertMany with ordered:false continues on duplicate errors if (err.insertedCount) totalPersisted += err.insertedCount; if (err.result?.insertedIds) { for (const [i, insertedId] of Object.entries(err.result.insertedIds)) { if (unpersisted[i]) unpersisted[i].id = insertedId.toString(); } } console.error(`💬 [${collectionName}] Persistence error:`, err.message); } }
return totalPersisted; }
async shutdown() { console.log("💬 ChatManager shutting down..."); if (this.persistTimer) clearInterval(this.persistTimer); const count = await this.persistAllMessages(); console.log(`💬 Persisted ${count} total messages`);
if (this.mongoClient) { await this.mongoClient.close(); console.log("💬 MongoDB connection closed"); } }}