Something went wrong. Try again.
A work-in-progress chat bot for Streamplace with chat overlay functionality
Something went wrong. Try again.
12 kB · 418 lines
TypeScript
at main
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419import { 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<K, V> { private store: Map<ResourceUri, V> = new Map(); private index: Map<K, ResourceUri> = 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<V> { 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<Did, { moderator: Did }>( (r) => r.moderator, ); private shoutouts = new IndexedCollection< Did, OnlineTimtinkersBotShoutout.Main >((r) => r.user); private shoutoutShorthands: Map<string, ResourceUri> = new Map();
private hasBeenGreeted: Map<Did, Date> = 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<void> { 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<void> { 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<void> { 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<Did, Date> { return this.hasBeenGreeted; }
getModerators(): IndexedCollection<Did, { moderator: Did }> { 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<Did, OnlineTimtinkersBotShoutout.Main> { return this.shoutouts; }
getShoutoutShorthands(): Map<string, ResourceUri> { 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; }}