bot.py 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614
  1. """A small, opinionated framework for building bots on Chatto.
  2. Chatto 0.5.0 introduced *bot accounts*: user identities flagged ``is_bot``
  3. that authenticate with a **key** (e.g. ``cht_BK_...``) rather than a
  4. username/password. The key is used **directly as a bearer token** — there is
  5. no ``/auth/login`` round-trip — and it resolves to a ``User`` with a
  6. capability grant set.
  7. This module turns :class:`chattolib.client.ChattoClient` into the standard
  8. library for bots by adding three things on top of the raw client:
  9. * **Key-based login** — :meth:`Bot.login` builds an authenticated client from
  10. a bot key in one call.
  11. * **An event dispatcher** — :meth:`Bot.run` opens the realtime stream and
  12. routes incoming events (messages, reactions, mentions, presence, typing,
  13. room changes) to the async handlers you register.
  14. * **Flattened verbs** — :meth:`Bot.say`, :meth:`Bot.reply`,
  15. :meth:`Bot.react`, :meth:`Bot.set_status`, :meth:`Bot.join_room`,
  16. :meth:`Bot.create_room`, and friends, so a bot's "brain" reads like
  17. natural language instead of protobuf plumbing.
  18. A minimal bot::
  19. import asyncio
  20. from chattolib.bot import Bot
  21. async def on_message(event):
  22. if event.body.startswith("!hello"):
  23. await event.bot.reply(event, "hi!")
  24. async def main():
  25. async with await Bot.login("cht_BK_...") as bot:
  26. bot.on("message", on_message)
  27. await bot.run()
  28. asyncio.run(main())
  29. Requires the ``chattolib[realtime]`` extra for the live stream.
  30. """
  31. from __future__ import annotations
  32. import asyncio
  33. import contextlib
  34. from collections.abc import Awaitable, Callable
  35. from dataclasses import dataclass
  36. from typing import Any
  37. from chattolib import _pb # noqa: F401 — installs the generated pb import path
  38. from chattolib._pb.chatto.api.v1 import presence_pb2
  39. from chattolib.client import ChattoClient
  40. from chattolib.exceptions import ChattoError
  41. from chattolib.realtime import (
  42. ChattoRealtimeCloseError,
  43. ChattoRealtimeError,
  44. RealtimeConnection,
  45. RealtimeEvent,
  46. RealtimeProjectionEvent,
  47. stream_events,
  48. )
  49. from chattolib.types import (
  50. Message,
  51. PresenceStatus,
  52. Room,
  53. RoomGroup,
  54. RoomWithViewerState,
  55. User,
  56. )
  57. def _presence(value: Any) -> PresenceStatus:
  58. """Coerce a protobuf presence value (int or enum) to a PresenceStatus."""
  59. name = value.name if hasattr(value, "name") else presence_pb2.PresenceStatus.Name(int(value))
  60. # The protobuf name is the *value* of our StrEnum (e.g. "PRESENCE_STATUS_ONLINE"),
  61. # so look it up by value, not by member name.
  62. for member in PresenceStatus:
  63. if member.value == name:
  64. return member
  65. return PresenceStatus.UNSPECIFIED
  66. __all__ = [
  67. "Bot",
  68. "BotError",
  69. "BotEvent",
  70. "BotMessageEvent",
  71. "BotPresenceEvent",
  72. "BotReactionEvent",
  73. "BotRoomEvent",
  74. "BotTypingEvent",
  75. "BotUserEvent",
  76. ]
  77. class BotError(ChattoError):
  78. """Raised for bot-framework-level errors (bad handlers, bad keys, ...)."""
  79. # ---------------------------------------------------------------------------
  80. # Event types
  81. # ---------------------------------------------------------------------------
  82. @dataclass
  83. class BotEvent:
  84. """Base class for every event the dispatcher delivers to a handler.
  85. ``bot`` is the owning :class:`Bot`, so a handler can act on the event
  86. (e.g. ``await event.bot.reply(event, "...")``) without closing over
  87. globals.
  88. """
  89. bot: Bot
  90. kind: str
  91. async def _noop(self) -> None: # pragma: no cover - interface
  92. ...
  93. @dataclass
  94. class BotMessageEvent(BotEvent):
  95. """A new or updated message in a room the bot can see.
  96. ``message`` is the full :class:`Message`. ``actor`` is the author's
  97. :class:`User` when the server included it; otherwise ``None`` (hydrate
  98. with ``await bot.get_user(message.actor_id)`` if you need it).
  99. """
  100. message: Message
  101. actor: User | None = None
  102. @property
  103. def room_id(self) -> str:
  104. return self.message.room_id
  105. @property
  106. def body(self) -> str | None:
  107. return self.message.body
  108. @property
  109. def is_mention(self) -> bool:
  110. """True when this message mentions the bot (best-effort)."""
  111. me = self.bot.user
  112. if me is None or not self.message.body:
  113. return False
  114. return f"@{me.login}" in self.message.body or me.display_name in self.message.body
  115. @dataclass
  116. class BotReactionEvent(BotEvent):
  117. """A reaction was added to (or removed from) a message."""
  118. room_id: str
  119. message_event_id: str
  120. emoji: str
  121. user_id: str
  122. added: bool = True
  123. @dataclass
  124. class BotPresenceEvent(BotEvent):
  125. """A user's presence status changed."""
  126. user_id: str
  127. status: PresenceStatus
  128. @dataclass
  129. class BotTypingEvent(BotEvent):
  130. """A user started (or stopped) typing in a room/thread."""
  131. room_id: str
  132. thread_root_event_id: str | None = None
  133. @dataclass
  134. class BotRoomEvent(BotEvent):
  135. """A room lifecycle change (created, updated, archived, member join/...)."""
  136. room: Room | None = None
  137. detail: str = "" # e.g. "created", "updated", "archived", "user_joined"
  138. @dataclass
  139. class BotUserEvent(BotEvent):
  140. """A user profile upsert or removal observed in the projection."""
  141. user: User | None = None
  142. removed: bool = False
  143. # ---------------------------------------------------------------------------
  144. # The Bot
  145. # ---------------------------------------------------------------------------
  146. Handler = Callable[[BotEvent], Awaitable[None]]
  147. class Bot:
  148. """A Chatto bot: an authenticated client plus an event dispatcher and a
  149. set of convenience verbs.
  150. Create one with :meth:`login` (from a bot key) or :meth:`from_client`
  151. (wrapping an already-authenticated :class:`ChattoClient`).
  152. """
  153. def __init__(
  154. self,
  155. client: ChattoClient,
  156. *,
  157. base_url: str | None = None,
  158. ) -> None:
  159. self._client = client
  160. self._base_url = base_url or client.base_url
  161. self._user: User | None = None
  162. self._handlers: dict[str, list[Handler]] = {}
  163. self._connection: RealtimeConnection | None = None
  164. self._running = False
  165. # -- construction ------------------------------------------------------
  166. @classmethod
  167. async def login(
  168. cls,
  169. key: str,
  170. *,
  171. base_url: str | None = None,
  172. ) -> Bot:
  173. """Authenticate a bot with its **key**.
  174. The key (e.g. ``cht_BK_...``) is used directly as the bearer token —
  175. no username/password. ``base_url`` defaults to the public Chatto
  176. server; pass it to target a self-hosted or preview deployment.
  177. """
  178. if base_url is None:
  179. base_url = ChattoClient.DEFAULT_BASE_URL
  180. client = ChattoClient(token=key, base_url=base_url)
  181. bot = cls(client, base_url=base_url)
  182. await bot._probe_identity()
  183. return bot
  184. @classmethod
  185. def from_client(cls, client: ChattoClient) -> Bot:
  186. """Wrap an already-authenticated client (e.g. a human login)."""
  187. return cls(client)
  188. async def _probe_identity(self) -> None:
  189. """Resolve the bot's own identity and warn (not fail) if not a bot."""
  190. try:
  191. self._user = await self._client.me()
  192. except ChattoError:
  193. self._user = None
  194. return
  195. if self._user is not None and not self._user.is_bot:
  196. # Not fatal — a human client can drive the same verbs — but surface
  197. # it so a misconfigured key is obvious.
  198. import warnings
  199. warnings.warn(
  200. f"Bot key authenticated as {self._user.login!r}, which is not "
  201. "flagged is_bot. It will still work, but this is usually a "
  202. "misconfigured key.",
  203. stacklevel=2,
  204. )
  205. # -- identity ----------------------------------------------------------
  206. @property
  207. def client(self) -> ChattoClient:
  208. """The underlying :class:`ChattoClient` for any RPC not wrapped here."""
  209. return self._client
  210. @property
  211. def base_url(self) -> str:
  212. return self._base_url
  213. @property
  214. def user(self) -> User | None:
  215. """The bot's own :class:`User` profile, once known."""
  216. return self._user
  217. @property
  218. def login_name(self) -> str | None:
  219. return self._user.login if self._user else None
  220. @property
  221. def display_name(self) -> str | None:
  222. return self._user.display_name if self._user else None
  223. # -- lifecycle ---------------------------------------------------------
  224. async def __aenter__(self) -> Bot:
  225. return self
  226. async def __aexit__(self, *exc: Any) -> None:
  227. await self.close()
  228. async def close(self) -> None:
  229. if self._connection is not None:
  230. with contextlib.suppress(Exception):
  231. await self._connection.close()
  232. self._connection = None
  233. self._running = False
  234. await self._client.close()
  235. # -- event registration ------------------------------------------------
  236. def on(self, kind: str, handler: Handler) -> Handler:
  237. """Register ``handler`` for events of ``kind``.
  238. ``kind`` is one of: ``message``, ``reaction``, ``presence``,
  239. ``typing``, ``room``, ``user``, ``*`` (all events). ``handler`` is an
  240. ``async def`` taking a single :class:`BotEvent`. Returns the handler
  241. so ``on`` can be used as a decorator.
  242. """
  243. self._handlers.setdefault(kind, []).append(handler)
  244. return handler
  245. def off(self, kind: str, handler: Handler) -> None:
  246. with contextlib.suppress(ValueError):
  247. self._handlers[kind].remove(handler)
  248. async def _dispatch(self, event: BotEvent) -> None:
  249. handlers = list(self._handlers.get(event.kind, [])) + list(self._handlers.get("*", []))
  250. for handler in handlers:
  251. try:
  252. await handler(event)
  253. except Exception: # noqa: BLE001 - one bad handler must not kill the loop
  254. import logging
  255. logging.exception("bot handler for %r raised", event.kind)
  256. # -- the run loop ------------------------------------------------------
  257. # -- verbs: messaging --------------------------------------------------
  258. async def say(self, room_id: str, body: str = "", *, join_if_needed: bool = True) -> Message:
  259. """Post a message to a room. Returns the created :class:`Message`.
  260. With ``join_if_needed`` (the default), the bot joins the room first if
  261. it isn't already a member — so a bot that wants to be present
  262. everywhere can simply ``await bot.say(room_id, ...)`` without a
  263. separate join step. Pass ``join_if_needed=False`` to instead surface
  264. the server's ``permission_denied`` if the bot can't post.
  265. """
  266. try:
  267. return await self._client.post_message(room_id, body)
  268. except ChattoError as exc:
  269. if not join_if_needed or "not a member" not in str(exc):
  270. raise
  271. await self._client.join_room(room_id)
  272. return await self._client.post_message(room_id, body)
  273. async def reply(
  274. self,
  275. target: BotMessageEvent | Message,
  276. body: str,
  277. *,
  278. also_send_to_channel: bool = False,
  279. ) -> Message:
  280. """Reply to a message (or a :class:`BotMessageEvent`) in its room/thread."""
  281. if isinstance(target, BotMessageEvent):
  282. target = target.message
  283. return await self._client.post_message(
  284. target.room_id,
  285. body,
  286. thread_root_event_id=target.thread_root_event_id,
  287. in_reply_to=target.id,
  288. also_send_to_channel=also_send_to_channel,
  289. )
  290. async def react(self, room_id: str, message_event_id: str, emoji: str) -> bool:
  291. """Add an emoji reaction to a message."""
  292. return await self._client.add_reaction(room_id, message_event_id, emoji)
  293. async def unreact(self, room_id: str, message_event_id: str, emoji: str) -> bool:
  294. """Remove one of the bot's emoji reactions from a message."""
  295. return await self._client.remove_reaction(room_id, message_event_id, emoji)
  296. # -- verbs: presence & status -----------------------------------------
  297. async def set_presence(self, status: PresenceStatus) -> PresenceStatus:
  298. """Set the bot's presence (``ONLINE`` / ``IDLE`` / ``DO_NOT_DISTURB``)."""
  299. return await self._client.update_presence(status)
  300. async def set_status(self, emoji: str, text: str) -> dict[str, Any]:
  301. """Set the bot's custom status (e.g. a "working on X" note)."""
  302. return await self._client.update_custom_status(emoji, text)
  303. async def clear_status(self) -> dict[str, Any]:
  304. """Clear the bot's custom status."""
  305. return await self._client.delete_custom_status()
  306. # -- verbs: rooms ------------------------------------------------------
  307. async def join_room(self, room_id: str) -> Room:
  308. """Join a room. Returns the joined :class:`Room`."""
  309. return await self._client.join_room(room_id)
  310. async def join_room_group(self, group_id: str) -> list[str]:
  311. """Join **all** the rooms in a room group in one call.
  312. Mirrors the Chatto UI's one-click "join group" action. Returns the
  313. list of room IDs the bot is now a member of as a result.
  314. """
  315. return await self._client.join_room_group(group_id)
  316. async def list_room_groups(self) -> list[RoomGroup]:
  317. """List the room groups (and the rooms each contains) the bot can see."""
  318. return await self._client.list_room_groups()
  319. async def join_all_rooms(self) -> list[str]:
  320. """Join every room the bot can see, grouped the way the UI does.
  321. Joins each room group in one call (so a bot becomes a member of all
  322. the rooms in a group at once), then joins any ungrouped rooms
  323. individually. Returns the room IDs the bot is now a member of.
  324. """
  325. joined: list[str] = []
  326. groups = await self.list_room_groups()
  327. grouped_room_ids: set[str] = set()
  328. for group in groups:
  329. for rws in group.rooms:
  330. if rws.room is not None:
  331. grouped_room_ids.add(rws.room.id)
  332. try:
  333. joined.extend(await self.join_room_group(group.id))
  334. except ChattoError:
  335. # A group the bot can't join is skipped, not fatal.
  336. continue
  337. # Rooms that are not part of any group.
  338. for rws in await self.list_rooms():
  339. if rws.room is None or rws.room.id in grouped_room_ids:
  340. continue
  341. if rws.viewer_state.is_member:
  342. continue
  343. try:
  344. joined.append((await self.join_room(rws.room.id)).id)
  345. except ChattoError:
  346. continue
  347. return joined
  348. async def leave_room(self, room_id: str) -> bool:
  349. """Leave a room."""
  350. return await self._client.leave_room(room_id)
  351. async def create_room(
  352. self,
  353. name: str,
  354. group_id: str,
  355. *,
  356. description: str = "",
  357. universal: bool = False,
  358. ) -> Room:
  359. """Create a room in a room group."""
  360. return await self._client.create_room(
  361. name, group_id, description=description, universal=universal
  362. )
  363. async def list_rooms(self) -> list[RoomWithViewerState]:
  364. """List the rooms the bot is a member of."""
  365. return await self._client.list_rooms()
  366. async def mark_read(self, room_id: str) -> None:
  367. """Mark a room as read (clears its unread state for the bot)."""
  368. await self._client.mark_room_as_read(room_id)
  369. # -- the run loop ------------------------------------------------------
  370. async def run(
  371. self,
  372. *,
  373. resume_cursor: str | None = None,
  374. retained_room_ids: list[str] | None = None,
  375. until: asyncio.Event | None = None,
  376. ) -> None:
  377. """Connect the realtime stream and dispatch events until it closes.
  378. Reconnects automatically (with the last ``resume_cursor``) when the
  379. server drops the connection, unless a fatal protocol error is raised.
  380. Pass ``until`` to stop the loop on an external signal.
  381. """
  382. self._running = True
  383. cursor = resume_cursor
  384. while self._running:
  385. if until is not None and until.is_set():
  386. break
  387. try:
  388. async for frame in stream_events(
  389. self._client,
  390. resume_cursor=cursor,
  391. retained_room_ids=retained_room_ids,
  392. ):
  393. if until is not None and until.is_set():
  394. break
  395. if isinstance(frame, RealtimeProjectionEvent):
  396. cursor = frame.resume_cursor or cursor
  397. await self._handle_projection(frame)
  398. else:
  399. await self._handle_live(frame)
  400. except ChattoRealtimeCloseError as close:
  401. if not close.reconnect:
  402. raise
  403. # reconnectable close: loop again, resuming from the cursor
  404. continue
  405. except ChattoRealtimeError:
  406. raise
  407. self._running = False
  408. async def _handle_live(self, frame: RealtimeEvent) -> None:
  409. if frame.kind == "presence_changed":
  410. payload = frame.payload
  411. event = BotPresenceEvent(
  412. bot=self,
  413. kind="presence",
  414. user_id=payload.user_id,
  415. status=_presence(payload.status),
  416. )
  417. await self._dispatch(event)
  418. elif frame.kind == "user_typing":
  419. payload = frame.payload
  420. await self._dispatch(
  421. BotTypingEvent(
  422. bot=self,
  423. kind="typing",
  424. room_id=payload.room_id,
  425. thread_root_event_id=payload.thread_root_event_id or None,
  426. )
  427. )
  428. elif frame.kind == "session_terminated":
  429. self._running = False
  430. async def _handle_projection(self, frame: RealtimeProjectionEvent) -> None:
  431. for op in frame.operations:
  432. case = op.operation
  433. if case == "room_timeline_event_upsert":
  434. await self._on_timeline_upsert(op)
  435. elif case == "presences_replace":
  436. for user_id, status in op.payload.statuses.items():
  437. await self._dispatch(
  438. BotPresenceEvent(
  439. bot=self,
  440. kind="presence",
  441. user_id=user_id,
  442. status=_presence(status),
  443. )
  444. )
  445. elif case == "room_upsert":
  446. room = op.payload.room
  447. if room is not None:
  448. await self._dispatch(
  449. BotRoomEvent(
  450. bot=self,
  451. kind="room",
  452. room=Room.parse(_pb_to_dict(room)),
  453. detail="upsert",
  454. )
  455. )
  456. elif case == "room_remove":
  457. await self._dispatch(
  458. BotRoomEvent(bot=self, kind="room", room=None, detail="removed")
  459. )
  460. elif case == "user_upsert":
  461. await self._dispatch(
  462. BotUserEvent(
  463. bot=self,
  464. kind="user",
  465. user=User.parse(_pb_to_dict(op.payload)),
  466. )
  467. )
  468. elif case == "user_remove":
  469. await self._dispatch(BotUserEvent(bot=self, kind="user", user=None, removed=True))
  470. async def _on_timeline_upsert(self, op: Any) -> None:
  471. payload = op.payload
  472. event = payload.event
  473. case = event.WhichOneof("event") if event is not None else None
  474. if case == "message_posted":
  475. posted = event.message_posted
  476. msg = Message.parse(_pb_to_dict(posted.message)) if posted.HasField("message") else None
  477. if msg is None:
  478. return
  479. actor = None
  480. includes = payload.includes
  481. if includes is not None and msg.actor_id:
  482. u = includes.users.get(msg.actor_id)
  483. if u is not None:
  484. actor = User.parse(_pb_to_dict(u))
  485. await self._dispatch(
  486. BotMessageEvent(bot=self, kind="message", message=msg, actor=actor)
  487. )
  488. elif case in (
  489. "room_created",
  490. "room_updated",
  491. "room_deleted",
  492. "room_archived",
  493. "room_unarchived",
  494. "room_threading_mode_changed",
  495. "user_joined_room",
  496. "user_left_room",
  497. ):
  498. room = None
  499. sub = getattr(event, case, None)
  500. if sub is not None and sub.HasField("room"):
  501. room = Room.parse(_pb_to_dict(sub.room))
  502. await self._dispatch(BotRoomEvent(bot=self, kind="room", room=room, detail=case))
  503. def _pb_to_dict(msg: Any) -> dict[str, Any]:
  504. """Best-effort protobuf -> camelCase dict for the dataclass parsers."""
  505. from chattolib._transport import pb_to_dict
  506. return pb_to_dict(msg)