import { AppBskyActorProfile, AppBskyRichtextFacet } from "@atcute/bluesky"; import { CommitEvent } from "@atcute/jetstream"; import { is, ResourceUri } from "@atcute/lexicons"; import { atprotoClient } from "../main.ts"; import { AtprotoClient, type RepoWalkHandlers, walkRepo, } from "./atcuteUtils.ts"; import { CommandHandler, MessageCreateEvent } from "./commandHandler.ts"; import { didResolver } from "./didResolver.ts"; import { OnlineTimtinkersBotCommand, OnlineTimtinkersBotRecurring, OnlineTimtinkersBotShoutout, OnlineTimtinkersBotVariable, PlaceStreamChatMessage, PlaceStreamChatProfile, PlaceStreamLivestream, PlaceStreamRichtextFacet, } from "./lexicons/index.ts"; import { buildRichtext } from "./richtextUtils.ts"; import { RecurringCommandHandler } from "./recurringCommandHandler.ts"; import { type DidDocument, getAtprotoHandle, getPdsEndpoint, } from "@atcute/identity"; import { Cid } from "@atcute/lexicons/syntax"; /** * A collection of records keyed canonically by AT-URI, with a secondary index * keyed by an arbitrary semantic key derived from each record. */ export class IndexedCollection { private store: Map = new Map(); private index: Map = new Map(); constructor(private readonly keyFn: (record: V) => K) {} set(uri: ResourceUri, record: V): void { const existing = this.store.get(uri); if (existing !== undefined) { this.index.delete(this.keyFn(existing)); } this.store.set(uri, record); this.index.set(this.keyFn(record), uri); } delete(uri: ResourceUri): V | undefined { const record = this.store.get(uri); if (record === undefined) return undefined; this.index.delete(this.keyFn(record)); this.store.delete(uri); return record; } getByKey(key: K): V | undefined { const uri = this.index.get(key); return uri !== undefined ? this.store.get(uri) : undefined; } getByUri(uri: ResourceUri): V | undefined { return this.store.get(uri); } getUriByKey(key: K): ResourceUri | undefined { return this.index.get(key); } clear(): void { this.store.clear(); this.index.clear(); } get size(): number { return this.store.size; } values(): IterableIterator { return this.store.values(); } entries(): IterableIterator<[ResourceUri, V]> { return this.store.entries(); } } export default class StreamplaceBot { private streamerDid: Did; private streamerProfile: UserProfile; private commandPrefix: string; private enabled: boolean; private client: AtprotoClient; private commandHandler: CommandHandler = new CommandHandler(); private recurringCommandHandler: RecurringCommandHandler = new RecurringCommandHandler(this); // Caching private repoLoaded: boolean = false; private moderators = new IndexedCollection( (r) => r.moderator, ); private shoutouts = new IndexedCollection< Did, OnlineTimtinkersBotShoutout.Main >((r) => r.user); private shoutoutShorthands: Map = new Map(); private hasBeenGreeted: Map = new Map(); private greetingCooldown: number = 12 * 60 * 60 * 1000; // 12 hours private livestreamRecord: PlaceStreamLivestream.Main | undefined; /** * Create a new StreamplaceBot * @param streamerDid The DID of the streamer whose chat to respond in * @param commandPrefix The prefix that triggers bot commands (default: "!") */ constructor( streamerDocument: DidDocument, commandPrefix = "!", ) { const pdsHostUrl = getPdsEndpoint(streamerDocument); const handle = getAtprotoHandle(streamerDocument); if (!pdsHostUrl || !handle) { throw new Error( `Error fetching PDS host or handle for ${streamerDocument.id}`, ); } this.streamerDid = streamerDocument.id; this.streamerProfile = { pdsEndpoint: pdsHostUrl as `${string}:${string}`, handle: handle, }; this.commandPrefix = commandPrefix; this.enabled = true; this.client = atprotoClient; } // Initialize the bot - must be called before use async init(): Promise { await this.client.init(); // Walk the repo first — populates all caches before any handler initializes. await this.loadStreamerRepo(); // commandHandler.init() registers defaults and logs the command count; // custom commands were already ingested by the walk above. this.commandHandler.init(); this.recurringCommandHandler.init(); console.log( `StreamplaceBot initialized for streamer: ${this.streamerProfile.handle}`, ); } /** * Walk the streamer's full repo as a single CAR download, populating all * in-memory caches in one pass. To load additional collections at startup * (emotes, badges, …), add a handler entry here — no other changes needed. */ private async loadStreamerRepo(): Promise { if (this.repoLoaded) return; const handlers: RepoWalkHandlers = { // Profiles "app.bsky.actor.profile": ({ record }) => { const actorProfile = record as AppBskyActorProfile.Main; this.streamerProfile.actorProfile = actorProfile; return Promise.resolve(); }, "place.stream.chat.profile": ({ record }) => { const chatProfile = record as PlaceStreamChatProfile.Main; this.streamerProfile.chatProfile = chatProfile; return Promise.resolve(); }, // Moderators "place.stream.moderation.permission": ({ record, uri }) => { const mod = record as { moderator: Did }; this.moderators.set(uri as ResourceUri, mod); return Promise.resolve(); }, // Shoutouts "online.timtinkers.bot.shoutout": ({ record, uri }) => { if (!is(OnlineTimtinkersBotShoutout.mainSchema, record)) { return Promise.resolve(); } this.shoutouts.set(uri as ResourceUri, record); if (record.shorthand) { this.shoutoutShorthands.set( record.shorthand.toLowerCase(), uri as ResourceUri, ); } return Promise.resolve(); }, // Custom variables (must be indexed before commands) // CAR ordering is MST-lexicographic, so "variable" < "command" and // this handler is guaranteed to run before any command record. "online.timtinkers.bot.variable": ({ record, uri }) => { if (!is(OnlineTimtinkersBotVariable.mainSchema, record)) { return Promise.resolve(); } this.commandHandler.ingestVariable(uri as ResourceUri, record); return Promise.resolve(); }, // Custom commands "online.timtinkers.bot.command": async ({ record, uri }) => { if (!is(OnlineTimtinkersBotCommand.mainSchema, record)) return; await this.commandHandler.ingestRecord( uri as ResourceUri, record, ); }, // Recurring commands "online.timtinkers.bot.recurring": ({ record, uri }) => { if (!is(OnlineTimtinkersBotRecurring.mainSchema, record)) { return Promise.resolve(); } this.recurringCommandHandler.ingestRecord( uri as ResourceUri, record, ); return Promise.resolve(); }, }; try { await walkRepo( this.streamerProfile.pdsEndpoint, this.streamerDid, handlers, ); this.repoLoaded = true; console.log( `Loaded streamer repo: ${this.moderators.size} moderators, ` + `${this.shoutouts.size} shoutouts`, // Custom command / recurring counts logged by their own handlers ); } catch (error) { console.error("Error walking streamer repo:", error); } } // Process an incoming chat message and respond if it's a command async processMessage(event: CommitEvent): Promise { if (!this.enabled) return; if (event.commit.operation !== "create") return; const record = event.commit.record as PlaceStreamChatMessage.Main; const text = record.text.trim(); // Get or cache chatter information const chatter = await didResolver.resolve(event.did); // Increment timed command counters on every real message this.recurringCommandHandler.incrementMessageCount(); // Auto-greet first-time chatters with shoutout if they have one const lastGreeting = this.hasBeenGreeted.get(event.did); if ( this.streamerIsLive() && (!lastGreeting || (Date.now() - lastGreeting.getTime()) >= this.greetingCooldown) ) { this.hasBeenGreeted.set(event.did, new Date()); const shoutout = this.shoutouts.getByKey(event.did); if (shoutout) { const { text, facets } = await buildRichtext(shoutout.text); await this.sendMessage(text, facets); } } // Check if message starts with command prefix if (!text.startsWith(this.commandPrefix)) return; // Disallow bot self-prompts if (event.did === this.getBotDid()) return; // Extract command name and arguments const parts = text.slice(this.commandPrefix.length).split(/\s+/); const commandName = parts[0].toLowerCase(); const params = parts.slice(1); // Look up and execute command const command = this.commandHandler.getCommand(commandName); if (!command) return; console.log( `Executing command: ${commandName} in: ${this.streamerProfile.handle} from: ${chatter.handle}`, ); try { await command.run(event as MessageCreateEvent, params, this); } catch (error) { console.error(`Error executing command ${commandName}:`, error); } } // Send a message to the chat async sendMessage( text: string, facets?: PlaceStreamRichtextFacet.Main[] | AppBskyRichtextFacet.Main[], ): Promise<{ cid: Cid; uri: ResourceUri } | undefined> { try { const result = await this.client.createMessage( text, this.streamerDid, facets as PlaceStreamRichtextFacet.Main[], ); return { cid: result.cid, uri: result.uri }; } catch (error) { console.error("Error sending bot message:", error); } return; } // Write attestation on successful command execution async writeAttestiation( command: ResourceUri, chatter: Did, call: ResourceUri, reply?: ResourceUri, ): Promise<{ cid: Cid; uri: ResourceUri } | undefined> { try { const result = await this.client.createAttestation( command, this.streamerDid, chatter, call, reply, ); return { cid: result.cid, uri: result.uri }; } catch (error) { console.error("Error creating bot attestation:", error); } return; } // Check if streamer is live streamerIsLive(): boolean { if (this.livestreamRecord) { const record = this.livestreamRecord; if (!record.lastSeenAt) return false; const lastSeen = new Date(record.lastSeenAt).getTime(); const now = Date.now(); const idleTimeoutMs = record.idleTimeoutSeconds && record.idleTimeoutSeconds > 0 ? record.idleTimeoutSeconds * 1000 : 300 * 1000; return (now - lastSeen) <= idleTimeoutMs; } return false; } // Getters, Setters (alphabetically) getBotDid(): Did | undefined { return this.client.getDid(); } getCommandPrefix(): string { return this.commandPrefix; } getCommandHandler(): CommandHandler { return this.commandHandler; } getGreeted(): Map { return this.hasBeenGreeted; } getModerators(): IndexedCollection { return this.moderators; } getTimedCommandHandler(): RecurringCommandHandler { return this.recurringCommandHandler; } getShoutout(did: Did): OnlineTimtinkersBotShoutout.Main | undefined { return this.shoutouts.getByKey(did); } getShoutoutByUri( uri: ResourceUri, ): OnlineTimtinkersBotShoutout.Main | undefined { return this.shoutouts.getByUri(uri); } getShoutouts(): IndexedCollection { return this.shoutouts; } getShoutoutShorthands(): Map { return this.shoutoutShorthands; } getStreamerDid(): Did { return this.streamerDid; } setEnabled(enabled: boolean): void { if (!enabled) this.recurringCommandHandler.destroy(); this.enabled = enabled; } get livestream(): PlaceStreamLivestream.Main | undefined { return this.livestreamRecord; } set livestream(record: PlaceStreamLivestream.Main) { this.livestreamRecord = record; } }