import { AppBskyActorProfile, AppBskyGraphFollow } from "@atcute/bluesky"; import RichtextBuilder from "@atcute/bluesky-richtext-builder"; import { CommitEvent, JetstreamEvent } from "@atcute/jetstream"; import { is, ResourceUri } from "@atcute/lexicons"; import { atprotoClient, botInstances, followerAtUris } from "../main.ts"; import { getDidDocument } from "./atcuteUtils.ts"; import { didResolver } from "./didResolver.ts"; import { OnlineTimtinkersBotCommand, OnlineTimtinkersBotRecurring, OnlineTimtinkersBotShoutout, OnlineTimtinkersBotVariable, PlaceStreamChatMessage, PlaceStreamChatProfile, PlaceStreamLivestream, PlaceStreamLiveTeleport, } from "./lexicons/index.ts"; import StreamplaceBot from "./streamplaceBot.ts"; import { CDN_URL } from "../env.ts"; import { customVariablesAtUris } from "../main.ts"; import { describeRepo, hasCollectionWithWildcard } from "./atcuteUtils.ts"; import { getLatestLivestream } from "./helpers.ts"; // union for everything sent over the WebSocket type WsMessage = | { $type: "place.stream.chat.defs#messageView"; data: MessageView } | { $type: "place.stream.livestream#viewerCount"; data: ViewerCount } | { $type: "online.timtinkers.bot.defs#mediaView"; data: MediaView } | { $type: "place.stream.chat.defs#enrichedMessageView"; data: EnrichedChatMessage; }; // Map of collection name -> handler, async handlers must have their own .catch() type CommitHandler = (event: CommitEvent) => void | Promise; export class EventHandler { private streamerClients: Map>; private commitHandlers: Map; constructor(streamerClients: Map>) { this.streamerClients = streamerClients; this.commitHandlers = new Map([ ["app.bsky.actor.profile", (e) => this.handleActorProfile(e)], [ "app.bsky.graph.follow", (e) => this.handleBskyFollow(e).catch((err) => console.error("Error in handleBskyFollow:", err) ), ], [ "place.stream.chat.message", (e) => this.handleChatMessage(e).catch((err) => console.error("Error in handleChatMessage:", err) ), ], ["place.stream.chat.profile", (e) => this.handleChatProfile(e)], ["place.stream.livestream", (e) => this.handleLivestream(e)], [ "place.stream.live.teleport", (e) => this.handleTeleport(e).catch((err) => console.error("Error in handleTeleport:", err) ), ], [ "online.timtinkers.bot.command", (e) => this.handleNewCommand(e).catch((err) => console.error("Error in handleNewCommand:", err) ), ], [ "online.timtinkers.bot.recurring", (e) => this.handleRecurringCommand(e), ], [ "online.timtinkers.bot.shoutout", (e) => this.handleNewShoutout(e), ], [ "online.timtinkers.bot.variable", (e) => this.handleNewVariable(e), ], ]); } /** All ATProto collection names this bot subscribes to via Jetstream. */ get subscribedCollections(): string[] { return [...this.commitHandlers.keys()]; } handleEvent(event: JetstreamEvent) { if (event.kind === "commit") { try { this.commitHandlers.get(event.commit.collection)?.(event); } catch (err) { console.error( `Error handling commit for ${event.commit.collection}:`, err, ); } } else if (event.kind === "identity") { if (!didResolver.getUser(event.did)) return; const user = didResolver.getUser(event.did)!; user.handle = event.identity.handle; } } private async handleBskyFollow(event: CommitEvent): Promise { if (event.commit.operation === "create") { const record = event.commit.record as AppBskyGraphFollow.Main; // Case: user follows streamer if (botInstances.get(record.subject)) { const follower = await didResolver.resolve(event.did); const bot = botInstances.get(record.subject)!; // Don't send chat message unless streamer is currently live if (!bot.streamerIsLive()) return; // Don't send chat message unless follower has a place.stream.* record const repoDescription = await describeRepo( event.did, follower.pdsEndpoint, ); if ( !repoDescription || !repoDescription.collections || !hasCollectionWithWildcard( repoDescription.collections, "place.stream.*", ) ) return; const { text, facets } = new RichtextBuilder().addText( "Thank you for following the streamer, ", ).addMention(`@${follower.handle}`, event.did).addText("!"); bot.sendMessage(text, facets); } // Case user follows bot else if ( record.subject === atprotoClient.getDid() && !botInstances.get(event.did) ) { const followerDoc = await getDidDocument(event.did); const newBot = new StreamplaceBot(followerDoc); await newBot.init(); botInstances.set(event.did, newBot); } } else if (event.commit.operation === "delete") { // Case user unfollows bot if ( botInstances.get(event.did) ) { botInstances.get(event.did)!.setEnabled(false); followerAtUris.delete(event.did); botInstances.delete(event.did); } } } private handleActorProfile(event: CommitEvent): void { if (!didResolver.getUser(event.did)) return; if (event.commit.rkey !== "self") return; const profile = didResolver.getUser(event.did)!; // Case created or updated profile if ( event.commit.operation === "create" || event.commit.operation === "update" ) { const record = event.commit .record as AppBskyActorProfile.Main; profile.actorProfile = record; } // Case deleted profile if (event.commit.operation === "delete") { profile.actorProfile = undefined; } } private async handleChatMessage(event: CommitEvent): Promise { if (event.commit.operation === "create") { const record = event.commit.record as PlaceStreamChatMessage.Main; if (botInstances.has(record.streamer)) { // Send to bot const streamplaceBot = botInstances.get(record.streamer)!; streamplaceBot.processMessage(event); } // Enrich and pass to subscribed websocket clients const enrichedMessage = await this.enrichMessage( record.streamer, event, ); this.sendMessageToClients( { $type: "place.stream.chat.defs#enrichedMessageView", data: enrichedMessage, }, record.streamer, ); } } private handleChatProfile(event: CommitEvent): void { if (!didResolver.getUser(event.did)) return; if (event.commit.rkey !== "self") return; const profile = didResolver.getUser(event.did)!; // Case created or updated profile if ( event.commit.operation === "create" || event.commit.operation === "update" ) { const record = event.commit .record as PlaceStreamChatProfile.Main; profile.chatProfile = record; } // Case deleted profile else if (event.commit.operation === "delete") { profile.chatProfile = undefined; } } private async handleTeleport(event: CommitEvent): Promise { if (event.commit.operation === "create") { const record = event.commit.record as PlaceStreamLiveTeleport.Main; if (botInstances.get(record.streamer)) { const GRACE_MS = 2000; const SANITY_WINDOW_MS = 60 * 1000; // 1 minute const startsAtMs = Date.parse(record.startsAt); if (Number.isNaN(startsAtMs)) return; const now = Date.now(); const offsetMs = startsAtMs - now; // Discard anything not within one minute of "now", in either direction if (Math.abs(offsetMs) > SANITY_WINDOW_MS) return; const teleporter = await didResolver.resolve(event.did); const livestreamShoutout = await getLatestLivestream(event.did); const richtextBuilder = new RichtextBuilder() .addText("Welcome in teleporters! "); if (livestreamShoutout) { richtextBuilder.addMention( `@${teleporter.handle}`, event.did, ) .addText(livestreamShoutout); } else { richtextBuilder.addText("And thank you for the teleport ") .addMention(`@${teleporter.handle}`, event.did); } const { text, facets } = richtextBuilder; const send = () => { botInstances.get(record.streamer)!.sendMessage( text, facets, ); }; // delay is the time left until startsAt, plus the grace period const delayMs = offsetMs + GRACE_MS; if (delayMs <= 0) { send(); } else { setTimeout(send, delayMs); } } } } // Custom commands and shoutouts private async handleNewCommand(event: CommitEvent): Promise { const bot = botInstances.get(event.did); if (!bot) return; const uri = `at://${event.did}/${event.commit.collection}/${event.commit.rkey}` as ResourceUri; // Case new command if ( event.commit.operation === "create" || event.commit.operation === "update" ) { await bot.getCommandHandler().handleUpsert( uri, event.commit.record as OnlineTimtinkersBotCommand.Main, ); } // Case deleted command else if (event.commit.operation === "delete") { bot.getCommandHandler().handleDelete(uri); } } private handleNewShoutout(event: CommitEvent): void { const uri = `at://${event.did}/${event.commit.collection}/${event.commit.rkey}` as ResourceUri; const bot = botInstances.get(event.did); if (!bot) return; if ( event.commit.operation === "create" || event.commit.operation === "update" ) { const record = event.commit .record as OnlineTimtinkersBotShoutout.Main; bot.getShoutouts().set(uri, record); if (record.shorthand) { bot.getShoutoutShorthands().set( record.shorthand.toLowerCase(), uri, ); } if (event.commit.operation === "update") { bot.getGreeted().delete(record.user as Did); } } else if (event.commit.operation === "delete") { const record = bot.getShoutouts().delete(uri); if (!record) return; bot.getGreeted().delete(record.user as Did); if (record.shorthand) { bot.getShoutoutShorthands().delete( record.shorthand.toLowerCase(), ); } } } private handleLivestream(event: CommitEvent) { if (!botInstances.get(event.did)) return; if (event.commit.operation === "delete") return; if (!is(PlaceStreamLivestream.mainSchema, event.commit.record)) return; const record = event.commit .record as PlaceStreamLivestream.Main; const bot = botInstances.get(event.did); bot!.livestream = record; } private handleNewVariable(event: CommitEvent) { const uri = `at://${event.did}/${event.commit.collection}/${event.commit.rkey}` as ResourceUri; // Case new shoutout if ( event.commit.operation === "create" || event.commit.operation === "update" ) { const record = event.commit .record as OnlineTimtinkersBotVariable.Main; if (!is(OnlineTimtinkersBotVariable.mainSchema, record)) return; customVariablesAtUris.set(uri, record); } // Case deleted shoutout else if (event.commit.operation === "delete") { customVariablesAtUris.delete(uri); } } private handleRecurringCommand(event: CommitEvent): void { const uri = `at://${event.did}/${event.commit.collection}/${event.commit.rkey}` as ResourceUri; if (!botInstances.get(event.did)) return; const bot = botInstances.get(event.did)!; if ( event.commit.operation === "create" || event.commit.operation === "update" ) { const record = event.commit .record as OnlineTimtinkersBotRecurring.Main; bot.getTimedCommandHandler().handleUpsert(uri, record); } else if (event.commit.operation === "delete") { bot.getTimedCommandHandler().handleDelete(uri); } } sendMessageToClients(message: WsMessage, streamerDid: Did) { const subscribedClients = this.streamerClients.get(streamerDid); if (!subscribedClients) return; const messageJson = JSON.stringify(message); for (const client of subscribedClients) { if (client.readyState === WebSocket.OPEN) { client.send(messageJson); } } } private async enrichMessage( streamerDid: Did, commitEvent: CommitEvent, ): Promise { // TODO: Not handling deletes for now if (commitEvent.commit.operation === "delete") throw new Error(); const { handle, pdsEndpoint, actorProfile, chatProfile } = await didResolver.resolve(commitEvent.did); // Get streamer avatar const streamerAvatar = (await didResolver.resolve(streamerDid)) .actorProfile?.avatar; // Get moderators for badge creation const moderators = botInstances.get(streamerDid)?.getModerators(); const badges: BadgeView[] = []; if (streamerDid === commitEvent.did) { badges.push({ badgeType: "place.stream.badge.defs#streamer", issuer: "did:web:stream.place", recipient: commitEvent.did, }); } if (chatProfile?.selfLabels?.includes("bot")) { badges.push({ badgeType: "place.stream.badge.defs#bot", issuer: "did:web:stream.place", recipient: commitEvent.did, }); } if (moderators && moderators.getByKey(commitEvent.did)) { badges.push({ badgeType: "place.stream.badge.defs#mod", issuer: "did:web:stream.place", recipient: commitEvent.did, }); } interface BlobWithRef { $type: string; ref: { $link: string; }; mimeType?: string; size?: number; } return { service: "streamplace", author: { handle: handle, did: commitEvent.did, pdsEndpoint: pdsEndpoint, description: actorProfile?.description, displayName: actorProfile?.displayName, pronouns: actorProfile?.pronouns, website: actorProfile?.website, avatarUrl: actorProfile?.avatar ? this.imageUrl( commitEvent.did, (actorProfile.avatar as BlobWithRef).ref.$link, ) : undefined, bannerUrl: actorProfile?.banner ? this.imageUrl( commitEvent.did, (actorProfile.banner as BlobWithRef).ref.$link, ) : undefined, }, badges: badges, chatProfile: { $type: "place.stream.chat.profile", color: chatProfile?.color, selfLabels: chatProfile?.selfLabels, }, cid: commitEvent.commit.cid, deleted: false, streamerAvatarUrl: streamerAvatar ? this.imageUrl( streamerDid, (streamerAvatar as BlobWithRef).ref.$link, ) : undefined, indexedAt: new Date().toISOString(), // could convert the `time_us` but this is good enough record: commitEvent.commit.record as PlaceStreamChatMessage.Main, uri: `at://${commitEvent.did}/place.stream.chat.message/${commitEvent.commit.rkey}`, }; } // Helper function private imageUrl(did: Did, link: string): `${string}:${string}` { return `${CDN_URL}/feed_thumbnail/plain/${did}/${link}@jpeg`; } }