import { type JetstreamEvent } from "@atcute/jetstream"; import { JETSTREAM_URL } from "../../env.ts"; import { EventHandler } from "./eventHandler.ts"; // Collections we want to filter in Jetstream for const wantedCollections = [ "app.bsky.actor.profile", "app.bsky.graph.follow", "online.timtinkers.bot.command", "online.timtinkers.bot.recurring", "online.timtinkers.bot.shoutout", "online.timtinkers.bot.variable", "place.stream.chat.message", "place.stream.chat.profile", "place.stream.livestream", "place.stream.live.teleport", ]; // Client subscription message interface SubscriptionMessage { type: "subscribe" | "unsubscribe"; streamer: string; // DID } // WebSocket service class class BotWebSocketService { private jetstreamWs: WebSocket | null = null; private clients = new Map>(); // client -> subscribed streamers private streamerClients = new Map>(); // streamer -> clients private _eventHandler = new EventHandler(this.streamerClients); constructor() {} start() { console.log("Starting WebSocket service"); this.connectToJetstream(); } initClient(client: WebSocket) { this.clients.set(client, new Set()); } // Handle new connection with pre-defined streamer DID handleNewConnection(client: WebSocket, streamerDid: Did) { this.initClient(client); // Auto-subscribe to the streamer from the route parameter this.handleClientMessage(client, { type: "subscribe", streamer: streamerDid, }); } private connectToJetstream() { const baseUrl = JETSTREAM_URL; const jetstreamUrl = `${baseUrl}?${ wantedCollections.map((c) => `wantedCollections=${c}`).join("&") }`; this.jetstreamWs = new WebSocket(jetstreamUrl); let pingInterval: number | null = null; this.jetstreamWs.onopen = () => { console.log("Connected to Jetstream"); // Send a ping every 15 seconds to keep the connection alive pingInterval = setInterval(() => { if (this.jetstreamWs?.readyState === WebSocket.OPEN) { this.jetstreamWs.send(JSON.stringify({ type: "ping" })); } }, 15_000); }; this.jetstreamWs.onmessage = (event) => { try { const jetstreamEvent: JetstreamEvent = JSON.parse(event.data); this._eventHandler.handleEvent(jetstreamEvent); } catch (error) { console.error("Error processing jetstream message:", error); } }; this.jetstreamWs.onclose = () => { console.log( "Jetstream connection closed, attempting to reconnect...", ); if (pingInterval !== null) clearInterval(pingInterval); setTimeout(() => this.connectToJetstream(), 5000); }; this.jetstreamWs.onerror = (error) => { console.error("Jetstream WebSocket error:", error); }; } handleClientMessage(client: WebSocket, message: SubscriptionMessage) { const clientStreamers = this.clients.get(client); if (!clientStreamers) return; switch (message.type) { case "subscribe": { clientStreamers.add(message.streamer); if (!this.streamerClients.has(message.streamer)) { this.streamerClients.set(message.streamer, new Set()); } this.streamerClients.get(message.streamer)!.add(client); client.send(JSON.stringify({ type: "subscribed", streamer: message.streamer, })); break; } case "unsubscribe": { clientStreamers.delete(message.streamer); const streamerClientSet = this.streamerClients.get( message.streamer, ); if (streamerClientSet) { streamerClientSet.delete(client); if (streamerClientSet.size === 0) { this.streamerClients.delete(message.streamer); } } console.log( `Client unsubscribed from streamer: ${message.streamer}`, ); client.send(JSON.stringify({ type: "unsubscribed", streamer: message.streamer, })); break; } } } removeClient(client: WebSocket) { if (!this.clients.has(client)) return; const clientStreamers = this.clients.get(client); if (clientStreamers) { for (const streamer of clientStreamers) { const streamerClientSet = this.streamerClients.get(streamer); if (streamerClientSet) { streamerClientSet.delete(client); if (streamerClientSet.size === 0) { this.streamerClients.delete(streamer); } } } } this.clients.delete(client); } get eventHandler(): EventHandler { return this._eventHandler; } } export const botWS = new BotWebSocketService();