adapter.py 131 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501250225032504250525062507250825092510251125122513251425152516251725182519252025212522252325242525252625272528252925302531253225332534253525362537253825392540254125422543254425452546254725482549255025512552255325542555255625572558255925602561256225632564256525662567256825692570257125722573257425752576257725782579258025812582258325842585258625872588258925902591259225932594259525962597259825992600260126022603260426052606260726082609261026112612261326142615261626172618261926202621262226232624262526262627262826292630263126322633263426352636263726382639264026412642264326442645264626472648264926502651265226532654265526562657265826592660266126622663266426652666266726682669267026712672267326742675267626772678267926802681268226832684268526862687268826892690269126922693269426952696269726982699270027012702270327042705270627072708270927102711271227132714271527162717271827192720272127222723272427252726272727282729273027312732273327342735273627372738273927402741274227432744274527462747274827492750275127522753275427552756275727582759276027612762276327642765276627672768276927702771277227732774277527762777277827792780278127822783278427852786278727882789279027912792279327942795279627972798279928002801280228032804280528062807280828092810281128122813281428152816281728182819282028212822282328242825282628272828282928302831283228332834283528362837283828392840284128422843284428452846284728482849285028512852285328542855285628572858285928602861286228632864286528662867286828692870287128722873287428752876287728782879288028812882288328842885288628872888288928902891289228932894289528962897289828992900290129022903290429052906290729082909291029112912291329142915291629172918291929202921292229232924292529262927292829292930293129322933293429352936293729382939294029412942294329442945294629472948294929502951295229532954295529562957295829592960296129622963296429652966296729682969297029712972297329742975297629772978297929802981298229832984298529862987298829892990299129922993299429952996299729982999300030013002300330043005300630073008300930103011301230133014301530163017301830193020302130223023302430253026302730283029303030313032303330343035303630373038303930403041304230433044304530463047304830493050305130523053305430553056305730583059306030613062306330643065306630673068306930703071307230733074307530763077307830793080308130823083308430853086308730883089309030913092309330943095309630973098309931003101310231033104310531063107310831093110311131123113311431153116311731183119312031213122312331243125312631273128312931303131313231333134313531363137313831393140314131423143314431453146314731483149315031513152315331543155315631573158315931603161316231633164316531663167316831693170317131723173317431753176317731783179318031813182318331843185318631873188318931903191319231933194319531963197319831993200320132023203320432053206320732083209321032113212321332143215321632173218321932203221322232233224322532263227322832293230323132323233323432353236323732383239324032413242324332443245324632473248324932503251325232533254325532563257325832593260326132623263326432653266326732683269327032713272327332743275327632773278327932803281328232833284328532863287328832893290329132923293329432953296329732983299330033013302330333043305330633073308330933103311
  1. """Chatto Platform Adapter for Hermes Agent.
  2. A plugin-based gateway adapter that connects to a Chatto server
  3. (self-hosted team chat) and relays messages to/from the Hermes agent.
  4. The adapter uses the chattolib library for all Chatto API interactions,
  5. including both outbound messaging and realtime WebSocket connections.
  6. """
  7. from __future__ import annotations
  8. import random
  9. # Put the vendored dependencies for THIS platform on sys.path before importing
  10. # anything from chattolib. Imported relatively as part of the plugin package and
  11. # absolutely when this module is loaded standalone (e.g. by the tests).
  12. try:
  13. from .vendor_path import setup_vendor_path
  14. except ImportError: # pragma: no cover - depends on how the module is loaded
  15. from vendor_path import setup_vendor_path
  16. setup_vendor_path()
  17. import asyncio
  18. import hashlib
  19. import logging
  20. import mimetypes
  21. import os
  22. import re
  23. import tempfile
  24. from collections import deque
  25. from dataclasses import dataclass, field
  26. from datetime import UTC, datetime
  27. from difflib import SequenceMatcher
  28. from enum import StrEnum
  29. from typing import Any
  30. from urllib.parse import unquote, urlsplit
  31. import httpx
  32. logger = logging.getLogger(__name__)
  33. from gateway.config import Platform, PlatformConfig
  34. from gateway.platforms.base import (
  35. BasePlatformAdapter,
  36. MessageEvent,
  37. MessageType,
  38. ProcessingOutcome,
  39. SendResult,
  40. cache_media_bytes,
  41. get_inbound_media_max_bytes,
  42. validate_inbound_media_size,
  43. )
  44. from gateway.session import build_session_key
  45. # Chattolib imports (vendored)
  46. # Using vendored chattolib from vendor/chattolib/
  47. # See vendor_chattolib.sh for how to update the vendored copy
  48. # Absolute imports — vendor/ is on sys.path (see above) and chattolib's own
  49. # modules import each other absolutely. Mixing in relative ".vendor.chattolib"
  50. # imports would load a second, distinct copy of every module, so isinstance()
  51. # checks across the two copies would silently fail.
  52. try:
  53. from chattolib._pb.chatto.realtime.v1 import ( # noqa: F401 — importability probe
  54. realtime_pb2,
  55. )
  56. from chattolib._pb.chatto.realtime.v1.events_pb2 import (
  57. MessagePostedEvent, # type: ignore[import-untyped]
  58. )
  59. from chattolib.client import (
  60. ChattoClient,
  61. )
  62. from chattolib.exceptions import (
  63. ChattoAuthError,
  64. ChattoError,
  65. )
  66. from chattolib.realtime import (
  67. ChattoRealtimeCloseError,
  68. ChattoRealtimeError,
  69. RealtimeEvent,
  70. stream_events,
  71. )
  72. from chattolib.types import (
  73. Message,
  74. MessageAttachment,
  75. PresenceStatus,
  76. Room,
  77. RoomKind,
  78. RoomWithViewerState,
  79. User,
  80. )
  81. except ImportError as e:
  82. # Fail loudly: continuing here only defers the failure to a confusing
  83. # NameError somewhere deep in the adapter.
  84. logger.error("Chatto: failed to import vendored chattolib: %s", e)
  85. raise
  86. try:
  87. from .platform_config import (
  88. ChattoConfiguration,
  89. ChattoConstants,
  90. )
  91. except ImportError: # pragma: no cover - loaded as a top-level module (tests)
  92. from platform_config import (
  93. ChattoConfiguration,
  94. ChattoConstants,
  95. )
  96. # --------------------------------------------------------------------------- #
  97. # Chat types
  98. # --------------------------------------------------------------------------- #
  99. class HermesChatType(StrEnum):
  100. """The ``chat_type`` vocabulary the Hermes gateway understands.
  101. Declared in ``gateway/session.py:161`` as ``"dm", "group", "channel",
  102. "thread"`` and consumed as a bare string all over the gateway:
  103. ``SessionSource.description`` (session.py:239) and the PII-redacting
  104. description in ``build_session_context_prompt`` (session.py:537) both
  105. branch on these exact values and fall back to a nameless generic case for
  106. anything else, and ``build_session_key`` puts the value straight into the
  107. session key. Passing a chattolib ``RoomKind`` (``"ROOM_KIND_CHANNEL"``)
  108. therefore does not fail loudly — it just quietly degrades what the agent is
  109. told about where it is.
  110. A StrEnum so it stays a drop-in ``str`` at every one of those call sites.
  111. GROUP vs CHANNEL
  112. ----------------
  113. There is no strict contract between the two, and the adapters disagree in
  114. practice: Slack labels every non-DM conversation ``"group"`` (including real
  115. channels), Discord uses both, and Telegram reserves ``"channel"`` for actual
  116. broadcast channels. The intended reading is ``group`` = ordinary
  117. multi-participant chat, ``channel`` = broadcast surface.
  118. The distinction only changes behaviour in three places:
  119. 1. Authorization (``gateway/authz_mixin.py``) — the only security-relevant
  120. one. The group-scoped env allowlists apply to ``{"group", "forum"}``
  121. ONLY, never to ``"channel"``: ``{PLATFORM}_GROUP_ALLOWED_USERS`` /
  122. ``_GROUP_ALLOWED_CHATS`` (:616), the chat-id allowlist (:708) and the
  123. Telegram legacy shim (:724). The adapter-delegation paths in turn treat
  124. all three alike (:461, :649, :674, :694), where the value only picks
  125. ``group_allow_from`` over ``allow_from`` from ``config.extra``.
  126. For Chatto both choices are equivalent today: those group env maps hold
  127. Telegram and QQBot only (:535-541), and our own allowlist runs through
  128. ``CHATTO_ALLOWED_USERS``, which is chat_type-independent.
  129. 2. What the agent is told — ``SessionSource.description`` renders
  130. ``"group: Name"`` vs ``"channel: Name"`` (session.py:239-246), likewise
  131. the PII-redacted variant (session.py:537-544).
  132. 3. The session key, which embeds the literal (session.py:1192). Changing
  133. the value for a room re-buckets its existing sessions.
  134. Explicitly NOT affected: ``is_shared_multi_user_session`` (session.py:1063)
  135. only looks at ``"dm"`` and ``thread_id``, so sender prefixes, the multi-user
  136. prompt line and ``group_sessions_per_user`` treat group and channel
  137. identically.
  138. """
  139. DM = "dm"
  140. GROUP = "group"
  141. CHANNEL = "channel"
  142. # Emitted by adapters whose thread events are their own chat type (Slack,
  143. # Discord). We don't: a Chatto thread keeps its room's chat_type and is
  144. # identified by ``thread_id`` on the source instead. Listed for the record,
  145. # because build_session_key rewrites the slot to "thread" itself
  146. # (session.py:1190).
  147. THREAD = "thread"
  148. # Not declared in session.py:161 but real: Telegram forum topics travel as
  149. # "forum", and the authz group allowlists above accept it alongside "group".
  150. # Chatto has no equivalent, so we never emit it.
  151. # Chatto only distinguishes DMs from channels. UNSPECIFIED means the server
  152. # sent a kind this vendored chattolib doesn't know: map it to the generic
  153. # multi-user bucket rather than guessing "channel", and never to "dm" — that
  154. # value drives session isolation (is_shared_multi_user_session, session.py:1063)
  155. # and would silently turn a room into a private conversation.
  156. #
  157. # CHANNEL for RoomKind.CHANNEL is the descriptive choice and carries no
  158. # behavioural cost (see the GROUP vs CHANNEL note above). Switching to GROUP for
  159. # Slack parity would be this one line — plus the re-bucketing of existing
  160. # sessions that point 3 of that note describes.
  161. _ROOM_KIND_TO_CHAT_TYPE: dict[RoomKind, HermesChatType] = {
  162. RoomKind.DM: HermesChatType.DM,
  163. RoomKind.CHANNEL: HermesChatType.CHANNEL,
  164. RoomKind.UNSPECIFIED: HermesChatType.GROUP,
  165. }
  166. def chat_type_for_room_kind(kind: RoomKind | None) -> HermesChatType:
  167. """Map a chattolib RoomKind onto the gateway's chat_type vocabulary.
  168. An unknown or missing kind becomes ``GROUP`` — see ``_ROOM_KIND_TO_CHAT_TYPE``.
  169. """
  170. if kind is None:
  171. return HermesChatType.GROUP
  172. return _ROOM_KIND_TO_CHAT_TYPE.get(kind, HermesChatType.GROUP)
  173. class RoomPolicy(StrEnum):
  174. """How a channel-kind room treats an inbound message.
  175. Derived per room from ``CHATTO_REQUIRE_MENTION_ROOMS`` /
  176. ``CHATTO_OPTIONAL_MENTION_ROOMS`` — see ``_room_policy``. A StrEnum
  177. so the values log readably without a formatting dance.
  178. """
  179. # Listed for unaddressed answers: every message is dispatched, and one
  180. # aimed at a named colleague gets a 🫥 acknowledgement instead of a reply.
  181. OPEN = "open"
  182. # Listed for addressed-only participation: dispatches only messages that
  183. # mention the bot (@name or a broadcast handle); others are dropped.
  184. REQUIRE_MENTION = "require_mention"
  185. # Listed in neither config: the room is silent — not dispatched at all,
  186. # read-only like any unlisted membership.
  187. SILENT = "silent"
  188. # Blocking filesystem/network helpers. The adapter runs on the shared gateway
  189. # event loop, so file reads and HTTP downloads are pushed to a worker thread
  190. # via asyncio.to_thread — a slow disk or dead image URL must not stall every
  191. # platform's message processing.
  192. def _read_file_bytes(path: str) -> bytes:
  193. """Read a whole file synchronously (run via asyncio.to_thread)."""
  194. with open(path, "rb") as f:
  195. return f.read()
  196. def _write_file_bytes(path: str, data: bytes) -> None:
  197. """Write bytes to a file synchronously (run via asyncio.to_thread)."""
  198. with open(path, "wb") as f:
  199. f.write(data)
  200. def _normalise_outbound_text(content: str) -> str:
  201. """Normalise outgoing text for Chatto.
  202. Chatto renders Markdown natively, so there is nothing to escape or
  203. translate — the only transformations here are the ones that measurably
  204. render wrong: CRLF line endings (which show up as stray blank lines)
  205. and runs of more than two blank lines. Shared by the adapter's
  206. ``format_message`` and the standalone cron sender, so both paths render
  207. identically.
  208. """
  209. if not content:
  210. return content
  211. normalised = content.replace("\r\n", "\n").replace("\r", "\n")
  212. while "\n\n\n\n" in normalised:
  213. normalised = normalised.replace("\n\n\n\n", "\n\n\n")
  214. return normalised
  215. # --------------------------------------------------------------------------- #
  216. # Adapter
  217. # --------------------------------------------------------------------------- #
  218. # Presence as the roster block spells it. UNSPECIFIED gets no label: a server
  219. # that never tracked presence should not have the roster claim anything.
  220. _PRESENCE_LABELS = {
  221. PresenceStatus.ONLINE: "online",
  222. PresenceStatus.AWAY: "away",
  223. PresenceStatus.DO_NOT_DISTURB: "do not disturb",
  224. PresenceStatus.OFFLINE: "offline",
  225. }
  226. @dataclass
  227. class _RoomRoster:
  228. """One announced room's roster projection.
  229. A miniature of the server's room membership: which user IDs belong to
  230. the room (the users themselves live in the shared ``_user_cache``) and
  231. how many directory entries were beyond the fetch limit, rendered as
  232. ``… and N more``. Presence changes patch the cached users in place;
  233. only a membership change or reconnect discards this and refetches.
  234. """
  235. member_ids: set[str] = field(default_factory=set)
  236. unfetched: int = 0
  237. def hermes_adapter_factory(config: PlatformConfig):
  238. """Construct a ChattoAdapter from a PlatformConfig."""
  239. return ChattoAdapter(config)
  240. class ChattoAdapter(BasePlatformAdapter):
  241. """Chatto platform adapter.
  242. Receives messages via WebSocket realtime, sends via ConnectRPC.
  243. """
  244. # Read by BasePlatformAdapter.max_message_length_for_chat(), which the
  245. # gateway and the stream consumer use to chunk outgoing messages. Without
  246. # it they fall back to 4096 and split Chatto messages far earlier than
  247. # necessary — send() itself already truncates at SPLIT_THRESHOLD.
  248. MAX_MESSAGE_LENGTH = ChattoConstants.MAX_MESSAGE_LENGTH
  249. splits_long_messages = True
  250. supports_code_blocks: bool = True
  251. supports_status_text: bool = True # client.update_custom_status
  252. def __init__(self, pconfig: PlatformConfig):
  253. """Signature needs to be compatible with BasePlatformAdapter.__init__."""
  254. super().__init__(
  255. config=pconfig, platform=Platform(ChattoConstants.PLATFORM_NAME)
  256. )
  257. # "extra" has been pre-populated by Hermes from config.yaml's extra block.
  258. # --- Configuration from our configuration data class with some logic ---
  259. self.chatto_config: ChattoConfiguration = ChattoConfiguration(pconfig)
  260. # ------ State -------
  261. # Our own user, filled in by connect(). Events arriving before connect()
  262. # completes must not blow up on an undefined attribute.
  263. self.me: User | None = None
  264. # --- Runtime state ---
  265. self._room_names: dict[str, str] = {}
  266. self._room_kinds: dict[str, RoomKind] = {}
  267. # Event IDs already processed — chattolib may redeliver events across
  268. # reconnects, so every inbound event is checked against this list.
  269. # Bounded deques: appending past the cap drops the oldest ID on its own.
  270. self._seen: deque[str] = deque(maxlen=ChattoConstants.SEEN_CAP)
  271. # Message IDs this adapter has handed to the gateway (posted or
  272. # edit-re-dispatched). Edits of anything on this list never start a
  273. # fresh turn — that is the lock against re-answering settled
  274. # conversations by editing old messages.
  275. self._dispatched_ids: deque[str] = deque(maxlen=ChattoConstants.SEEN_CAP)
  276. # Rooms whose member roster has been projected from the directory —
  277. # the roster is room-scoped and rides on every channel turn (threads
  278. # hold isolated sessions, so each turn needs its own copy). Together
  279. # with _rosters this is our miniature projection of server state:
  280. # presence changes patch the cached users in place; membership events
  281. # and reconnects discard a room so its next turn refetches once.
  282. self._roster_announced: deque[str] = deque(maxlen=ChattoConstants.SEEN_CAP)
  283. # Announced rooms' membership projection (see _RoomRoster).
  284. self._rosters: dict[str, _RoomRoster] = {}
  285. # session_key -> message ID currently being processed there. Written
  286. # by on_processing_start, cleared by on_processing_complete; an edit
  287. # landing on the recorded ID is a mid-run correction.
  288. self._processing: dict[str, str] = {}
  289. self._joined_room_ids: list[str] = []
  290. # Rooms the server force-joined everyone into (Room.universal) — used
  291. # only for [universal] tags in the joined-rooms log line, never for
  292. # gating.
  293. self._universal_room_ids: set[str] = set()
  294. # DM room ID -> chat partner's login, resolved once via the member
  295. # directory so the joined-rooms summary can name who a DM is with.
  296. self._dm_partners: dict[str, str] = {}
  297. # One-shot guard for the unjoined-home-channel warning in _refresh_rooms.
  298. self._home_warning_logged = False
  299. self._ws_task: asyncio.Task | None = None
  300. self._presence_task: asyncio.Task | None = None
  301. # Persistent typing indicator loops per room
  302. self._typing_tasks: dict[str, asyncio.Task] = {}
  303. # Member directory cache: user_id -> user info dict
  304. self._user_cache: dict[str, User] = {}
  305. # Handle -> does a user hold it. Cached both ways; see _mentions_someone_else.
  306. self._known_handles: dict[str, bool] = {}
  307. # Chattolib client cache and lock for async access.
  308. self._chatto_client: ChattoClient | None = None
  309. self._chatto_client_lock: asyncio.Lock = asyncio.Lock()
  310. # ------------------------------------------------------------------ #
  311. # Auth
  312. # ------------------------------------------------------------------ #
  313. async def _get_chatto_client(self: ChattoAdapter) -> ChattoClient | None:
  314. """Return the shared ChattoClient, creating and logging in on first use.
  315. The default way to get a client. The fast path is a plain attribute
  316. read; first creation runs under a lock so concurrent callers log in
  317. exactly once. Returns ``None`` when no client exists yet or creation
  318. failed (bad credentials, unreachable server) — a normal state during
  319. startup, shutdown and reconnects, not an exceptional one. Callers
  320. decide what "no client" means for them, as a guard clause:
  321. client = await self._get_chatto_client()
  322. if client is None:
  323. logger.warning("Chatto: dropping X - no client available")
  324. return
  325. """
  326. if self._chatto_client is not None:
  327. return self._chatto_client
  328. async with self._chatto_client_lock:
  329. if self._chatto_client is not None:
  330. return self._chatto_client
  331. try:
  332. client = await self._open_client(
  333. base_url=self.chatto_config.base_url.value,
  334. token=self.chatto_config.token.value,
  335. )
  336. self._chatto_client = client
  337. logger.info(
  338. "Chatto: authenticated via token to '%s'",
  339. self.chatto_config.base_url.value,
  340. )
  341. return client
  342. except ChattoAuthError as e:
  343. logger.error("Chatto: authentication failed: %s", e)
  344. return None
  345. except (ChattoError, ValueError) as e:
  346. logger.error("Chatto: failed to create client: %s", e)
  347. return None
  348. async def _require_chatto_client(self) -> ChattoClient:
  349. """Return a ChattoClient or raise RuntimeError if unavailable.
  350. The exception-flavoured variant of :meth:`_get_chatto_client`, for
  351. callers whose surrounding machinery already routes exceptions.
  352. Currently that is only the realtime event loop, whose except chain
  353. turns the raise into a logged warning plus a backed-off reconnect —
  354. no special "no client" branch needed there. Everywhere else, prefer
  355. the ``is None`` guard shown in :meth:`_get_chatto_client`.
  356. """
  357. client = await self._get_chatto_client()
  358. if client is None:
  359. raise RuntimeError("Chatto client unavailable")
  360. return client
  361. # ------------------------------------------------------------------ #
  362. # Connection
  363. # ------------------------------------------------------------------ #
  364. async def _open_client(
  365. self,
  366. *,
  367. base_url: str,
  368. token: str | None = None,
  369. ) -> ChattoClient:
  370. """Return a connected ``ChattoClient`` using a bearer token.
  371. Raises ValueError when no token is configured — caught upstream as a
  372. normal creation failure, so misconfiguration reads as a logged error
  373. instead of a crash.
  374. """
  375. if not token:
  376. raise ValueError("Chatto: CHATTO_TOKEN not configured")
  377. return ChattoClient(token=token, base_url=base_url)
  378. async def connect(self, *, is_reconnect: bool = False) -> bool:
  379. """Connect to Chatto and start the realtime event stream.
  380. BasePlatformAdapter override
  381. """
  382. logger.info("Chatto: connecting...")
  383. client = await self._get_chatto_client()
  384. if client is None:
  385. self._set_fatal_error(
  386. "connect_failed", "Chatto client not available", retryable=True
  387. )
  388. return False
  389. # Get our first own user info
  390. try:
  391. self.me = await client.me()
  392. except Exception as exc:
  393. logger.error("Chatto: failed to get user info: %s", exc)
  394. self._set_fatal_error(
  395. "chatto_auth_failed",
  396. f"Chatto auth failed: {exc}",
  397. retryable=False,
  398. )
  399. try:
  400. if self._chatto_client is not None:
  401. await self._chatto_client.close()
  402. finally:
  403. self._chatto_client = None
  404. return False
  405. # Announce online presence so the bot appears online in the member list.
  406. # The server treats this as a TTL, so _presence_refresh_loop below has to
  407. # keep re-announcing it — a single call here lapses back to offline.
  408. await self._announce_online()
  409. self._closing = False
  410. # Start background realtime WS event stream loop.
  411. self._ws_task = asyncio.create_task(
  412. self._chattolib_event_loop(),
  413. name="chatto-event-stream",
  414. )
  415. self._presence_task = asyncio.create_task(
  416. self._presence_refresh_loop(),
  417. name="chatto-presence-refresh",
  418. )
  419. self._mark_connected()
  420. logger.info(
  421. "Chatto: connected to %s as '%s' (display_name: '%s' id: '%s')",
  422. self.chatto_config.base_url.value,
  423. self.me.login,
  424. self.me.display_name,
  425. self.me.id,
  426. )
  427. return True
  428. async def _announce_online(self) -> bool:
  429. """Tell the server we are online. Returns whether the call got through.
  430. Logged at warning level on failure: a silently dropped presence call is
  431. indistinguishable from a bot that is simply not running.
  432. """
  433. client = await self._get_chatto_client()
  434. if client is None:
  435. logger.warning(
  436. "Chatto: presence refresh failed, bot may appear offline: no client"
  437. )
  438. return False
  439. try:
  440. await client.set_presence(status=PresenceStatus.ONLINE)
  441. return True
  442. except Exception as exc:
  443. logger.warning(
  444. "Chatto: presence refresh failed, bot may appear offline: %s", exc
  445. )
  446. return False
  447. async def _presence_refresh_loop(self) -> None:
  448. """Re-announce ONLINE until disconnect, since presence expires server-side.
  449. Failures are not fatal — the next tick tries again, so a blip in the
  450. presence endpoint costs at most one interval of visible offline time.
  451. """
  452. while not self._closing:
  453. await self._sleep_interruptible(ChattoConstants.PRESENCE_REFRESH_INTERVAL)
  454. if self._closing:
  455. return
  456. await self._announce_online()
  457. async def disconnect(self) -> None:
  458. """Stop WebSocket, presence refresh, typing tasks, and clear state.
  459. BasePlatformAdapter override
  460. """
  461. # No explicit offline broadcast: chattolib rejects OFFLINE outright
  462. # ("stop refreshing to go offline"), so cancelling the refresh loop
  463. # below is what actually takes the bot offline.
  464. self._closing = True
  465. # Cancel all typing tasks
  466. for chat_id in list(self._typing_tasks.keys()):
  467. await self.stop_typing(chat_id)
  468. if self._ws_task and not self._ws_task.done():
  469. self._ws_task.cancel()
  470. try:
  471. await self._ws_task
  472. except (asyncio.CancelledError, Exception):
  473. logger.debug(
  474. "Chatto: websocket task ended during disconnect", exc_info=True
  475. )
  476. self._ws_task = None
  477. if self._presence_task and not self._presence_task.done():
  478. self._presence_task.cancel()
  479. try:
  480. await self._presence_task
  481. except (asyncio.CancelledError, Exception):
  482. logger.debug(
  483. "Chatto: presence task ended during disconnect", exc_info=True
  484. )
  485. self._presence_task = None
  486. if self._chatto_client:
  487. try:
  488. await self._chatto_client.close()
  489. except Exception:
  490. logger.exception("Chatto: error closing client")
  491. finally:
  492. self._chatto_client = None
  493. logger.info("Chatto: disconnected")
  494. self._mark_disconnected()
  495. async def _seed_room(self, room_id: str) -> None:
  496. """Seed high-water mark from the newest events so a restart doesn't replay history."""
  497. try:
  498. client = await self._get_chatto_client()
  499. if client is None:
  500. logger.debug(
  501. "Chatto: _seed_room aborted - no client available for %s", room_id
  502. )
  503. return
  504. timeline_page = await client.get_room_events(room_id)
  505. for ev in timeline_page.events:
  506. if ev.id:
  507. self._mark_seen(ev.id)
  508. logger.debug(
  509. "Chatto: seeded room %s with %d events",
  510. room_id,
  511. len(timeline_page.events),
  512. )
  513. except Exception as e:
  514. logger.warning("Chatto: get room events failed for %s: %s", room_id, e)
  515. # ------------------------------------------------------------------ #
  516. # Realtime Event List
  517. # ------------------------------------------------------------------ #
  518. def _mark_seen(self, event_id: str) -> None:
  519. # The deque's maxlen evicts the oldest ID — no manual trimming.
  520. self._seen.append(event_id)
  521. def _is_seen(self, event_id: str) -> bool:
  522. return event_id in self._seen
  523. # ------------------------------------------------------------------ #
  524. # WebSocket Realtime Transport
  525. # ------------------------------------------------------------------ #
  526. def _own_handles(self) -> set[str]:
  527. """The handles that address this bot, lowercased for comparison.
  528. Chatto resolves mentions case-insensitively (FDR-006), so ``@Hermes_Bot``
  529. and ``@hermes_bot`` are the same handle everywhere in our gates.
  530. """
  531. if not self.me:
  532. return set()
  533. return {
  534. handle.lower() for handle in (self.me.login, self.me.display_name) if handle
  535. }
  536. def _mention_candidates(self, body: str) -> list[str]:
  537. """Candidate @-handles in a message body, in order of appearance.
  538. Mirrors the Chatto web frontend's extraction (apps/frontend/src/lib/
  539. mentions.ts upstream): candidates come from ``ChattoConstants.
  540. MENTION_RE`` outside code regions. Mentions inside fenced code blocks
  541. and inline code spans do not resolve upstream either, so a ``@bob``
  542. quoted in a snippet must not gate our behaviour.
  543. """
  544. without_fences = re.sub(r"(?s)(```|~~~).*?(\1|$)", " ", body)
  545. without_code = re.sub(r"`[^`\n]*`", " ", without_fences)
  546. return ChattoConstants.MENTION_RE.findall(without_code)
  547. def _mentions_me(self, body: str) -> bool:
  548. """Whether the message addresses this bot.
  549. By login, by display name, or by a broadcast handle — ``@here`` speaks
  550. to everyone present and the bot is one of them, so naming a colleague
  551. alongside it does not take the bot out of the audience. Matching is
  552. case-insensitive, like every mention resolution in Chatto.
  553. """
  554. own = self._own_handles()
  555. for handle in self._mention_candidates(body):
  556. lowered = handle.lower()
  557. if lowered in own or lowered in ChattoConstants.BROADCAST_MENTIONS:
  558. return True
  559. return False
  560. async def _handle_belongs_to_a_user(self, handle: str) -> bool:
  561. """Whether ``handle`` is the login of a real Chatto user.
  562. The API carries no mention entities — ``mention_confirmation_token`` is
  563. reserved in the message descriptor — so an @-token is only a candidate
  564. until the directory confirms it. Results are cached both ways under the
  565. lowercased handle (the mention namespace is case-insensitive per
  566. FDR-006), since the same handles recur and a miss is as reusable as a
  567. hit. Only logins are looked up: matching another user's display name,
  568. as the web frontend does against its room member list, has no directory
  569. equivalent here.
  570. """
  571. cache_key = handle.lower()
  572. known = self._known_handles.get(cache_key)
  573. if known is not None:
  574. return known
  575. client = await self._get_chatto_client()
  576. if client is None:
  577. # Unresolved means "not confirmed", so the message goes through.
  578. return False
  579. try:
  580. member = await client.admin_get_member(login=handle)
  581. except Exception as exc:
  582. logger.debug("Chatto: could not resolve handle @%s: %s", handle, exc)
  583. return False
  584. # admin_get_member returns pb_to_dict(GetMemberResponse):
  585. # {"member": {"user": {...}, ...}, "roles": [...]} when found,
  586. # or {} when not found.
  587. if not member.get("member"):
  588. self._known_handles[cache_key] = False
  589. return False
  590. from chattolib.types import User
  591. user = User.parse(member["member"].get("user"))
  592. self._known_handles[cache_key] = user is not None
  593. return user is not None
  594. async def _mentions_someone_else(self, body: str) -> bool:
  595. """Whether the message @-mentions a person who is not this bot.
  596. Broadcast handles are not a person — they address everyone present,
  597. the bot included, so they do not count as someone else. A handle no
  598. user holds is not a mention at all: someone writing *about* mentioning
  599. ("per @-mention", "@nonexistent") is talking to us, and staying silent
  600. on a false positive is worse than answering one.
  601. """
  602. own = self._own_handles()
  603. for handle in self._mention_candidates(body):
  604. lowered = handle.lower()
  605. if lowered in ChattoConstants.BROADCAST_MENTIONS:
  606. continue
  607. if lowered in own:
  608. continue
  609. if await self._handle_belongs_to_a_user(handle):
  610. return True
  611. return False
  612. def _check_auth(self, user: User) -> bool:
  613. """Whether this Chatto user may talk to the agent.
  614. Deliberately our own gate instead of the gateway's authz_mixin: its
  615. group allowlists key on chat_type and per-platform env vars
  616. (``{PLATFORM}_GROUP_ALLOWED_USERS``), none of which fit Chatto's one
  617. flat member directory. ``CHATTO_ALLOWED_USERS`` matches login and id,
  618. ``CHATTO_ALLOW_ALL_USERS`` overrides both — the gate stays
  619. chat_type-independent by design.
  620. """
  621. if self.chatto_config.allow_all_users.value:
  622. return True
  623. if user.login in self.chatto_config.allowed_users.value:
  624. return True
  625. if user.id in self.chatto_config.allowed_users.value:
  626. return True
  627. logger.warning(
  628. "Chatto: rejecting message from unauthorized user '%s' (%s)",
  629. user.login,
  630. user.id,
  631. )
  632. return False
  633. # ------------------------------------------------------------------ #
  634. # Room management over DM (/join, /leave)
  635. # ------------------------------------------------------------------ #
  636. _DM_COMMANDS = ("/join", "/leave")
  637. async def _handle_dm_command(self, room_id: str, body: str) -> bool:
  638. """Run a ``/join`` or ``/leave`` admin command sent as a direct message.
  639. Returns True when ``body`` is one of the commands — whether it
  640. succeeded or not — so the caller keeps it out of the agent pipeline.
  641. Membership lives on the Chatto server: a joined room reappears in
  642. every future ``list_rooms()`` and therefore survives restarts.
  643. """
  644. verb, _, argument = body.strip().partition(" ")
  645. if verb.lower() not in self._DM_COMMANDS:
  646. return False
  647. client = await self._get_chatto_client()
  648. if client is None:
  649. await self.send(chat_id=room_id, content="Chatto client is not connected.")
  650. return True
  651. argument = argument.strip()
  652. if not argument:
  653. await self.send(
  654. chat_id=room_id,
  655. content="Usage: /join <room-id or #name> | /leave <room-id or #name>",
  656. )
  657. return True
  658. error, target = await self._resolve_room_target(client, argument)
  659. if error or target is None:
  660. await self.send(chat_id=room_id, content=error or "Room lookup failed.")
  661. return True
  662. if verb.lower() == "/join":
  663. reply = await self._run_join(client, target)
  664. else:
  665. reply = await self._run_leave(client, target)
  666. await self.send(chat_id=room_id, content=reply)
  667. return True
  668. async def _resolve_room_target(
  669. self,
  670. client: ChattoClient,
  671. argument: str,
  672. ) -> tuple[str | None, RoomWithViewerState | None]:
  673. """Resolve a ``/join`` or ``/leave`` argument to a room.
  674. ``#name`` is looked up case-insensitively in a fresh directory scan
  675. (which also refreshes our name/kind caches); anything else is treated
  676. as a room ID and verified via GetRoom. An ambiguous name comes back as
  677. an error naming the candidates, so the admin can retry with an ID.
  678. """
  679. if not argument.startswith("#"):
  680. state = await client.get_room(argument)
  681. if state is None or state.room is None:
  682. return f"No room with ID '{argument}'.", None
  683. return None, state
  684. wanted = argument[1:].strip().casefold()
  685. # (state, room) pairs: a listed match's room is already narrowed here,
  686. # so the candidate listing below needs no fresh Optional dance.
  687. matches: list[tuple[RoomWithViewerState, Room]] = []
  688. for state in await client.list_rooms() or []:
  689. room_obj = state.room if state else None
  690. if room_obj and (room_obj.name or "").strip().casefold() == wanted:
  691. matches.append((state, room_obj))
  692. self._room_names[room_obj.id] = room_obj.name
  693. self._room_kinds[room_obj.id] = room_obj.kind
  694. if not matches:
  695. return f"No room named '{argument}'.", None
  696. if len(matches) > 1:
  697. candidates = "\n".join(f"• {room.name} ({room.id})" for _, room in matches)
  698. return (
  699. f"Several rooms are named '{argument}' — pick one by ID:\n{candidates}"
  700. ), None
  701. return None, matches[0][0]
  702. async def _run_join(self, client: ChattoClient, state: RoomWithViewerState) -> str:
  703. """Join a room via RoomService/JoinRoom and track it as joined.
  704. An account that already holds membership (invited natively in Chatto)
  705. needs no JoinRoom call — it only gets seeded and added to the list.
  706. Channel-kind rooms additionally report their mention-list status,
  707. because a channel on neither list stays silent and this reply is
  708. where users copy the room ID from (see ``_room_join_hint``).
  709. """
  710. room_obj = state.room
  711. if room_obj is None:
  712. # Unreachable via _resolve_room_target: both of its paths only
  713. # return states whose room they already inspected.
  714. return "Chatto returned an empty room record — try again."
  715. label = f"'{room_obj.name}' ({room_obj.id})"
  716. joined_room = room_obj
  717. if not state.viewer_state.is_member:
  718. try:
  719. joined_room = await client.join_room(room_obj.id) or room_obj
  720. except ChattoError as exc:
  721. logger.warning("Chatto: /join failed for %s (%s)", room_obj.id, exc)
  722. return f"Could not join {label}: {exc}"
  723. self._room_names[joined_room.id] = joined_room.name
  724. self._room_kinds[joined_room.id] = joined_room.kind
  725. if joined_room.id not in self._joined_room_ids:
  726. # Same rule as _refresh_rooms: silent rooms are joined read-only
  727. # — seeding history nothing will ever answer would be waste.
  728. if self._answers_in_room(joined_room.id):
  729. await self._seed_room(joined_room.id)
  730. else:
  731. logger.info(
  732. "Chatto: %s is on neither mention list - joined read-only",
  733. label,
  734. )
  735. self._joined_room_ids.append(joined_room.id)
  736. if room_obj.kind != RoomKind.DM:
  737. head = (
  738. f"Already a member of {label}."
  739. if state.viewer_state.is_member
  740. else f"Joined {label}."
  741. )
  742. return f"{head}\n{self._room_join_hint(room_obj.id)}"
  743. if state.viewer_state.is_member:
  744. return f"Already a member of {label} — listening there."
  745. return f"Joined {label}."
  746. async def _run_leave(self, client: ChattoClient, state: RoomWithViewerState) -> str:
  747. """Leave a room via RoomService/LeaveRoom and drop it from the joined list.
  748. Two rooms are refused: a DM conversation cannot be left, and leaving
  749. the configured home channel would silently break cron/notification
  750. delivery, which posts there through the standalone sender.
  751. """
  752. room_obj = state.room
  753. if room_obj is None:
  754. # Same invariant as _run_join: _resolve_room_target pre-inspects.
  755. return "Chatto returned an empty room record — try again."
  756. label = f"'{room_obj.name}' ({room_obj.id})"
  757. if room_obj.kind == RoomKind.DM:
  758. return "Direct messages cannot be left."
  759. home_id = (self.chatto_config.home_channel.value or "").strip()
  760. if home_id == room_obj.id:
  761. return (
  762. f"{label} is the configured home channel "
  763. "(CHATTO_HOME_CHANNEL); leaving it would break cron and "
  764. "notification delivery. Point CHATTO_HOME_CHANNEL elsewhere first."
  765. )
  766. try:
  767. left = await client.leave_room(room_obj.id)
  768. except ChattoError as exc:
  769. logger.warning("Chatto: /leave failed for %s (%s)", room_obj.id, exc)
  770. return f"Could not leave {label}: {exc}"
  771. if not left:
  772. return f"Chatto refused to leave {label}."
  773. if room_obj.id in self._joined_room_ids:
  774. self._joined_room_ids.remove(room_obj.id)
  775. # We are no longer part of this audience; keep no roster for it.
  776. self._evict_roster(room_obj.id)
  777. return f"Left {label}."
  778. def _room_join_hint(self, room_id: str) -> str:
  779. """The mention-list status appended to a channel-kind /join reply.
  780. Users do not know their room IDs by heart — this reply is where they
  781. copy them from, so a silent channel names both env vars verbatim,
  782. ready to paste into ~/.hermes/.env. Configuration resolves once at
  783. gateway startup, hence the restart note.
  784. """
  785. policy = self._room_policy(room_id)
  786. if policy == RoomPolicy.REQUIRE_MENTION:
  787. return (
  788. "This channel answers only @mentions "
  789. "(listed in CHATTO_REQUIRE_MENTION_ROOMS)."
  790. )
  791. if policy == RoomPolicy.OPEN:
  792. return (
  793. "This channel answers every message "
  794. "(listed in CHATTO_OPTIONAL_MENTION_ROOMS)."
  795. )
  796. return (
  797. "This channel stays silent until you list its ID in"
  798. " ~/.hermes/.env (then restart the gateway):\n"
  799. f" CHATTO_REQUIRE_MENTION_ROOMS={room_id} <- answer only @mentions\n"
  800. f" CHATTO_OPTIONAL_MENTION_ROOMS={room_id} <- answer every message"
  801. )
  802. # ------------------------------------------------------------------ #
  803. # Inbound attachments
  804. # ------------------------------------------------------------------ #
  805. async def _download_attachment_bytes(self, url: str) -> bytes:
  806. """Download an attachment, refusing to buffer more than the gateway cap.
  807. The Content-Length header is checked first so an oversized asset is
  808. rejected before a single chunk is read; the running total is re-checked
  809. as chunks arrive, because a missing or lying header must not smuggle an
  810. unbounded body past the cap.
  811. """
  812. max_bytes = get_inbound_media_max_bytes()
  813. chunks: list[bytes] = []
  814. total = 0
  815. async with (
  816. httpx.AsyncClient(
  817. timeout=ChattoConstants.HTTP_TIMEOUT,
  818. follow_redirects=True,
  819. ) as http,
  820. http.stream("GET", url) as resp,
  821. ):
  822. resp.raise_for_status()
  823. declared = resp.headers.get("content-length")
  824. if declared:
  825. try:
  826. declared_size = int(declared)
  827. except ValueError:
  828. logger.debug("Chatto: ignoring invalid Content-Length %r", declared)
  829. else:
  830. validate_inbound_media_size(
  831. declared_size,
  832. media_type="attachment",
  833. max_bytes=max_bytes,
  834. )
  835. async for chunk in resp.aiter_bytes():
  836. total += len(chunk)
  837. validate_inbound_media_size(
  838. total,
  839. media_type="attachment",
  840. max_bytes=max_bytes,
  841. )
  842. chunks.append(chunk)
  843. return b"".join(chunks)
  844. async def _cache_attachments(
  845. self,
  846. attachments: list[MessageAttachment],
  847. ) -> tuple[list[str], list[str], list[str]]:
  848. """Download message attachments into the gateway media cache.
  849. Returns ``(media_urls, media_types, media_kinds)`` — the paths are
  850. agent-visible cache paths, exactly what ``cache_media_bytes`` yields for
  851. every other platform. A failing attachment is logged and skipped: the
  852. message itself still reaches the agent.
  853. """
  854. media_urls: list[str] = []
  855. media_types: list[str] = []
  856. media_kinds: list[str] = []
  857. for att in attachments:
  858. url = att.asset_url.url if att.asset_url else ""
  859. filename = att.filename
  860. content_type = att.content_type
  861. if not url:
  862. # Videos are announced before transcoding finishes, so the
  863. # signed URL can legitimately be missing on arrival.
  864. logger.debug(
  865. "Chatto: attachment '%s' has no asset URL yet, skipping",
  866. filename,
  867. )
  868. continue
  869. try:
  870. data = await self._download_attachment_bytes(url)
  871. cached = cache_media_bytes(
  872. data,
  873. filename=filename,
  874. mime_type=content_type,
  875. )
  876. except Exception as e:
  877. logger.warning(
  878. "Chatto: failed to cache attachment '%s' (%s): %s",
  879. filename,
  880. content_type,
  881. e,
  882. )
  883. continue
  884. if cached is None:
  885. logger.warning(
  886. "Chatto: attachment '%s' (%s) could not be cached, skipping",
  887. filename,
  888. content_type,
  889. )
  890. continue
  891. media_urls.append(cached.path)
  892. media_types.append(cached.media_type)
  893. media_kinds.append(cached.kind)
  894. return media_urls, media_types, media_kinds
  895. @staticmethod
  896. def _message_type_for_media_kinds(media_kinds: list[str]) -> MessageType:
  897. """Pick the MessageType for a set of cached attachment kinds."""
  898. if "document" in media_kinds:
  899. return MessageType.DOCUMENT
  900. if "image" in media_kinds:
  901. return MessageType.PHOTO
  902. if "video" in media_kinds:
  903. return MessageType.VIDEO
  904. if "audio" in media_kinds:
  905. return MessageType.AUDIO
  906. return MessageType.TEXT
  907. async def _room_kind_for(
  908. self, client: ChattoClient, room_id: str
  909. ) -> RoomKind | None:
  910. """The room's kind, from cache or a fresh GetRoom lookup.
  911. Returns ``None`` when the room cannot be resolved — the caller treats
  912. that as "not dispatchable" rather than guessing a kind.
  913. """
  914. kind = self._room_kinds.get(room_id)
  915. if kind is not None:
  916. return kind
  917. room_viewer_state = await client.get_room(room_id)
  918. if room_viewer_state is None or room_viewer_state.room is None:
  919. return None
  920. kind = room_viewer_state.room.kind or RoomKind.UNSPECIFIED
  921. self._room_kinds[room_id] = kind
  922. return kind
  923. def _room_policy(self, room_id: str) -> RoomPolicy:
  924. """Which of the mention lists a channel-kind room is on.
  925. The two lists are mutually exclusive (enforced by
  926. ``hermes_validate_config``), so membership decides: optional beats
  927. require in the face of contradictory runtime config, and a room on
  928. neither list stays silent.
  929. """
  930. if room_id in self.chatto_config.optional_mention_rooms.value:
  931. return RoomPolicy.OPEN
  932. if room_id in self.chatto_config.require_mention_rooms.value:
  933. return RoomPolicy.REQUIRE_MENTION
  934. return RoomPolicy.SILENT
  935. def _answers_in_room(self, room_id: str) -> bool:
  936. """Whether inbound messages from this room reach the agent pipeline.
  937. Chatto knows only DMs and channels, so the split is ``kind == DM``:
  938. a DM always answers (``/join`` must stay reachable), every other
  939. room — channel-kind or an unknown kind, which is how servers that
  940. never set ``kind`` show up — opts in through the mention lists. A
  941. room whose policy is SILENT stays read-only (marked as read, never
  942. seeded or answered).
  943. """
  944. if self._room_kinds.get(room_id) == RoomKind.DM:
  945. return True
  946. return self._room_policy(room_id) != RoomPolicy.SILENT
  947. async def _room_roster_context(
  948. self, client: ChattoClient, room_id: str
  949. ) -> _RoomRoster | None:
  950. """Fetch the room's member roster as a fresh projection.
  951. The agent only ever sees its prompt: without this block it cannot know
  952. who else is in a channel, because unaddressed messages are dropped by
  953. mention gating long before they could teach it a name. Fetched users go
  954. into _user_cache so presence patching and mention resolution share one
  955. store. Returns ``None`` on failure or when nothing renderable remains —
  956. best-effort by design, a directory hiccup must never cost the turn.
  957. """
  958. try:
  959. # list_room_members() in 0.5.0b4 returns member IDs only;
  960. # batch_get_room_members() hydrates them into DirectoryMember.
  961. member_id_list, page = await client.list_room_members(
  962. room_id, limit=ChattoConstants.ROSTER_MEMBER_LIMIT
  963. )
  964. if not member_id_list:
  965. return None
  966. members = await client.batch_get_room_members(room_id, member_id_list)
  967. except Exception:
  968. logger.debug(
  969. "Chatto: could not list members of room %s", room_id, exc_info=True
  970. )
  971. return None
  972. logger.debug(
  973. "Chatto: roster for room %s: %d members fetched, total_count=%d",
  974. room_id,
  975. len(members),
  976. page.total_count,
  977. )
  978. member_ids: set[str] = set()
  979. users: list[User] = []
  980. for member in members:
  981. user = member.user
  982. if user is None or user.deleted:
  983. continue
  984. member_ids.add(user.id)
  985. self._user_cache[user.id] = user
  986. users.append(user)
  987. entries = [
  988. entry for entry in (self._roster_entry(user) for user in users) if entry
  989. ]
  990. if not entries:
  991. logger.debug(
  992. "Chatto: roster for room %s empty after filtering "
  993. "(%d fetched, all deleted/self/without user)",
  994. room_id,
  995. len(members),
  996. )
  997. return None
  998. logger.debug(
  999. "Chatto: roster for room %s: %d entries", room_id, len(entries)
  1000. )
  1001. unfetched = max(0, page.total_count - len(members))
  1002. return _RoomRoster(member_ids=member_ids, unfetched=unfetched)
  1003. def _roster_entry(self, user: User) -> str:
  1004. """One roster entry, e.g. ``@bob (Bob Example, online)``.
  1005. Skips deleted users and our own account (the agent knows itself);
  1006. display name and presence are optional parts of the parenthetical.
  1007. """
  1008. if user.deleted:
  1009. return ""
  1010. if self.me is not None and user.id == self.me.id:
  1011. return ""
  1012. entry = f"@{user.login}"
  1013. details = [
  1014. part
  1015. for part in (
  1016. user.display_name,
  1017. _PRESENCE_LABELS.get(user.presence_status, ""),
  1018. )
  1019. if part
  1020. ]
  1021. if details:
  1022. entry += f" ({', '.join(details)})"
  1023. return entry
  1024. def _roster_line(self, room_id: str) -> str:
  1025. """Render the room's roster line from its cached projection."""
  1026. roster = self._rosters.get(room_id)
  1027. if roster is None:
  1028. return ""
  1029. users = [self._user_cache.get(member_id) for member_id in roster.member_ids]
  1030. entries = [
  1031. entry
  1032. for entry in (
  1033. self._roster_entry(user) for user in users if user is not None
  1034. )
  1035. if entry
  1036. ]
  1037. more = f" … and {roster.unfetched} more" if roster.unfetched else ""
  1038. return ", ".join(entries) + more
  1039. def _announce_roster(self, room_id: str, roster: _RoomRoster) -> None:
  1040. """Record a freshly fetched roster as the room's projection.
  1041. Bounded like _dispatched_ids: appending past the cap drops the oldest
  1042. room's projection alongside its deque entry.
  1043. """
  1044. if room_id not in self._roster_announced:
  1045. oldest = (
  1046. self._roster_announced.popleft()
  1047. if len(self._roster_announced) == self._roster_announced.maxlen
  1048. else None
  1049. )
  1050. self._roster_announced.append(room_id)
  1051. if oldest is not None:
  1052. self._rosters.pop(oldest, None)
  1053. self._rosters[room_id] = roster
  1054. def _evict_roster(self, room_id: str) -> None:
  1055. """Drop all roster state for a room.
  1056. Called on membership changes (user_joined/left_room), when we leave a
  1057. room ourselves, and on reconnect — protocol v1 has no presence snapshot
  1058. on subscribe, so a discarded cache is what forces one honest refetch.
  1059. """
  1060. try:
  1061. self._roster_announced.remove(room_id)
  1062. except ValueError:
  1063. pass
  1064. self._rosters.pop(room_id, None)
  1065. async def _roster_for_room(self, client: ChattoClient, *, room_id: str) -> str:
  1066. """The room's roster line for this turn — delivered on every turn.
  1067. Threads hold isolated sessions and the agent only sees
  1068. channel_context per dispatch, so each channel turn carries the current
  1069. audience rather than deduplicating it. The projection exists to spare
  1070. the directory, not the prompt: presence patches land in _user_cache
  1071. in place, membership events or a reconnect evict the room so its next
  1072. turn refetches once. An unprojected room retries on its next turn
  1073. until a lookup succeeds.
  1074. """
  1075. if room_id not in self._roster_announced:
  1076. fetched = await self._room_roster_context(client, room_id)
  1077. if fetched is None:
  1078. return ""
  1079. self._announce_roster(room_id, fetched)
  1080. return self._roster_line(room_id)
  1081. async def _dispatch_message_posted(self, payload: MessagePostedEvent, event_id: str) -> None:
  1082. # Respond-room gate first: read-only memberships must not cost a
  1083. # single API call, so this runs before fetch_message and _get_chatto_client.
  1084. if not self._answers_in_room(payload.room_id):
  1085. logger.debug(
  1086. "Chatto: message from read-only room %s ignored", payload.room_id
  1087. )
  1088. return
  1089. client = await self._get_chatto_client()
  1090. if client is None:
  1091. logger.warning("Chatto: dropping message - no client available")
  1092. return
  1093. logger.debug("Chatto WS: 'message_posted' payload:%s", payload)
  1094. message = await client.get_message(
  1095. room_id=payload.room_id, event_id=event_id
  1096. )
  1097. if message is None or message.deleted_at:
  1098. return
  1099. message_body = message.body or ""
  1100. logger.debug("message: %s", message)
  1101. # A message carrying only an image/PDF has an empty body — dropping it
  1102. # here is what made attachments sent to Hermes disappear silently.
  1103. if not message_body and not message.attachments:
  1104. return
  1105. event = await self._admit_and_build(
  1106. client,
  1107. room_id=payload.room_id,
  1108. message=message,
  1109. source_message_id=event_id,
  1110. thread_root_event_id=payload.thread_root_event_id or None,
  1111. )
  1112. if event is None:
  1113. return
  1114. self._remember_dispatched(event.message_id or "")
  1115. logger.info("Chatto: dispatching message to Hermes")
  1116. await self.handle_message(event)
  1117. return
  1118. def _remember_dispatched(self, message_id: str) -> None:
  1119. """Record a message ID as handed to the gateway, capped like _seen."""
  1120. if not message_id:
  1121. return
  1122. self._dispatched_ids.append(message_id)
  1123. def _edit_is_fresh(self, message: Message) -> bool:
  1124. """Whether this edit is young enough to still be processed.
  1125. Age is measured against the posting time, so an edit to an hours-old
  1126. message cannot resurrect a settled conversation even when it arrives
  1127. right now.
  1128. """
  1129. if message.created_at is None or message.updated_at is None:
  1130. return True
  1131. age_seconds = (message.updated_at - message.created_at).total_seconds()
  1132. return age_seconds <= self.chatto_config.edit_window.value
  1133. def _session_key_for(self, source) -> str:
  1134. """The gateway's own session key for this source.
  1135. Built with exactly the inputs ``handle_message`` uses, so lookups in
  1136. ``_processing`` and calls to ``cancel_session_processing`` hit the
  1137. same session the gateway is running.
  1138. """
  1139. return build_session_key(
  1140. source,
  1141. group_sessions_per_user=self.config.extra.get(
  1142. "group_sessions_per_user", True
  1143. ),
  1144. thread_sessions_per_user=self.config.extra.get(
  1145. "thread_sessions_per_user", False
  1146. ),
  1147. )
  1148. async def _dispatch_message_edited(self, payload: Any) -> None:
  1149. """Route an inbound edit according to the edit-dispatch contract.
  1150. Three outcomes for the edited message:
  1151. - currently being processed → cancel that turn and re-dispatch with
  1152. the corrected text (the cancelled turn reports 🚫 via its
  1153. CANCELLED outcome hook),
  1154. - never dispatched (e.g. a forgotten @mention added later) → re-run
  1155. the admission gates against the new body and answer for real,
  1156. - already answered → stay answered.
  1157. Edits whose text parses as a DM membership command are dropped: the
  1158. command ran when the message was posted and must not run again.
  1159. """
  1160. if not self.chatto_config.edit_dispatch.value:
  1161. return
  1162. # Read-only memberships cost no API call, mirroring the posted path.
  1163. if not self._answers_in_room(payload.room_id):
  1164. logger.debug("Chatto: edit from read-only room %s ignored", payload.room_id)
  1165. return
  1166. client = await self._get_chatto_client()
  1167. if client is None:
  1168. logger.warning("Chatto: dropping edit - no client available")
  1169. return
  1170. logger.debug("Chatto WS: 'message_edited' payload:%s", payload)
  1171. message = await client.get_message(
  1172. room_id=payload.room_id, event_id=payload.message_event_id
  1173. )
  1174. if message is None or message.deleted_at:
  1175. return
  1176. message_body = message.body or ""
  1177. if not message_body and not message.attachments:
  1178. return
  1179. if not self._edit_is_fresh(message):
  1180. logger.debug(
  1181. "Chatto: edit of msg %s outside the edit window",
  1182. payload.message_event_id,
  1183. )
  1184. return
  1185. # Cheap triage before the admission pipeline: an edit to an
  1186. # already-settled message (dispatched, neither running nor queued)
  1187. # must not cost get_room/media calls or re-fire acknowledgements.
  1188. was_dispatched = payload.message_event_id in self._dispatched_ids
  1189. if (
  1190. was_dispatched
  1191. and payload.message_event_id not in self._processing.values()
  1192. and not any(
  1193. pending.message_id == payload.message_event_id
  1194. for pending in self._pending_messages.values()
  1195. )
  1196. ):
  1197. logger.debug(
  1198. "Chatto: edit of already-answered msg %s ignored",
  1199. payload.message_event_id,
  1200. )
  1201. return
  1202. event = await self._admit_and_build(
  1203. client,
  1204. room_id=payload.room_id,
  1205. message=message,
  1206. source_message_id=payload.message_event_id,
  1207. thread_root_event_id=message.thread_root_event_id or None,
  1208. allow_dm_commands=False,
  1209. )
  1210. if event is None:
  1211. return
  1212. session_key = self._session_key_for(event.source)
  1213. if self._processing.get(session_key) == payload.message_event_id:
  1214. logger.info(
  1215. "Chatto: msg %s edited mid-run - restarting the turn",
  1216. payload.message_event_id,
  1217. )
  1218. # The cancelled task's completion hook clears _processing and
  1219. # reports 🚫 before this coroutine moves on, because cancel
  1220. # awaits the task. Queued follow-ups must survive.
  1221. await self.cancel_session_processing(
  1222. session_key, release_guard=True, discard_pending=False
  1223. )
  1224. elif not was_dispatched:
  1225. logger.info(
  1226. "Chatto: msg %s was never dispatched - edit starts a fresh turn",
  1227. payload.message_event_id,
  1228. )
  1229. else:
  1230. pending = self._pending_messages.get(session_key)
  1231. if pending is not None and pending.message_id == payload.message_event_id:
  1232. # Still queued behind the running turn: correct it in place
  1233. # instead of answering stale wording later.
  1234. pending.text = event.text
  1235. logger.info(
  1236. "Chatto: queued msg %s updated to its edited text",
  1237. payload.message_event_id,
  1238. )
  1239. return
  1240. logger.debug(
  1241. "Chatto: edit of already-answered msg %s ignored",
  1242. payload.message_event_id,
  1243. )
  1244. return
  1245. self._remember_dispatched(event.message_id or "")
  1246. logger.info("Chatto: dispatching edited message to Hermes")
  1247. await self.handle_message(event)
  1248. async def _admit_and_build(
  1249. self,
  1250. client: ChattoClient,
  1251. *,
  1252. room_id: str,
  1253. message: Message,
  1254. source_message_id: str,
  1255. thread_root_event_id: str | None,
  1256. allow_dm_commands: bool = True,
  1257. ) -> MessageEvent | None:
  1258. """Run one hydrated inbound message through the admission pipeline.
  1259. Shared by the posted and the edited path: user resolution, auth,
  1260. mention gates, thread anchoring and media caching all behave
  1261. identically for both. Returns ``None`` for anything that must not
  1262. reach the agent. DM membership commands are executed here (side
  1263. effect) unless ``allow_dm_commands`` is False — edits pass False so
  1264. a corrected command line neither runs twice nor leaks to the agent.
  1265. A dispatch that opens a channel thread additionally carries the
  1266. member roster in ``channel_context`` (side effect: once the roster
  1267. was fetched, the thread root is recorded as announced).
  1268. The caller owns dispatching: a non-None result still needs
  1269. ``handle_message()``.
  1270. """
  1271. user = self._user_cache.get(message.actor_id)
  1272. if user is None:
  1273. logger.debug("Chatto: resolving user %s", message.actor_id)
  1274. member = await client.get_room_member(
  1275. room_id=message.room_id, user_id=message.actor_id
  1276. )
  1277. if member is None or member.user is None:
  1278. logger.debug(
  1279. "Chatto: dropping message from unknown user %s",
  1280. message.actor_id,
  1281. )
  1282. return None
  1283. user = member.user
  1284. self._user_cache[user.id] = user
  1285. if not self._check_auth(user):
  1286. logger.debug(
  1287. "Chatto: dropping message from user %s (auth denied)",
  1288. message.actor_id,
  1289. )
  1290. return None
  1291. room_kind = await self._room_kind_for(client, message.room_id)
  1292. if room_kind is None:
  1293. logger.debug(
  1294. "Chatto: dropping message from room %s (room_kind lookup failed)",
  1295. message.room_id,
  1296. )
  1297. return None
  1298. message_body = message.body or ""
  1299. logger.debug("message_body: %s room_kind: %s", message_body, room_kind)
  1300. # Membership commands ride in over DMs only: they change what the bot
  1301. # listens to and must never reach the agent pipeline or the mention
  1302. # gates.
  1303. if room_kind == RoomKind.DM:
  1304. if allow_dm_commands:
  1305. if await self._handle_dm_command(room_id, message_body):
  1306. return None
  1307. elif message_body.startswith("/"):
  1308. logger.debug("Chatto: edited DM command %r not re-run", message_body)
  1309. return None
  1310. # Mention gating deliberately covers everything but DMs: in a channel
  1311. # the bot is one of many listeners and must be addressed, whereas a DM
  1312. # is already addressed at it. Chatto has no group rooms — any
  1313. # multi-participant surface is a channel, and a room whose kind the
  1314. # server never set counts as one too (see _answers_in_room). The
  1315. # shared _mentions_me gate keeps this path and the someone-else check
  1316. # below on one definition of "addressed", broadcast handles included.
  1317. if room_kind != RoomKind.DM:
  1318. policy = self._room_policy(room_id)
  1319. if policy == RoomPolicy.REQUIRE_MENTION and not self._mentions_me(
  1320. message_body
  1321. ):
  1322. logger.debug(
  1323. "Chatto: dropping unaddressed message from %s (policy %s)",
  1324. room_id,
  1325. policy.value,
  1326. )
  1327. return None
  1328. # In an open channel we see every message, including ones plainly
  1329. # aimed at a named colleague. Answering those would be barging in,
  1330. # so acknowledge that we read it and stay quiet. Checked after the
  1331. # bot-mention test above, so a message naming us *and* someone
  1332. # else still counts as ours.
  1333. if (
  1334. policy == RoomPolicy.OPEN
  1335. and not self._mentions_me(message_body)
  1336. and await self._mentions_someone_else(message_body)
  1337. ):
  1338. logger.info(
  1339. "Chatto: message addresses someone else, acknowledging only"
  1340. )
  1341. if self.chatto_config.reactions.value:
  1342. await self.add_reaction(message.room_id, message.id, "🫥")
  1343. return None
  1344. # Thread anchoring — if the incoming message is inside a Chatto thread, we
  1345. # keep that thread by default; otherwise leave thread_id unset so
  1346. # replies land at the root.
  1347. thread_id = (
  1348. thread_root_event_id or None
  1349. ) # we could also take the room id but then, we're in a thread already.
  1350. if not thread_id and room_kind != RoomKind.DM:
  1351. thread_id = message.id
  1352. logger.debug(
  1353. "Chatto: thread anchoring for msg %s: incoming thread_root_event_id=%r, "
  1354. "room_kind=%s, anchored thread_id=%r",
  1355. message.id,
  1356. message.thread_root_event_id,
  1357. room_kind.name,
  1358. thread_id,
  1359. )
  1360. # Every channel turn carries the room's roster: threads hold isolated
  1361. # sessions and the agent only sees channel_context per dispatch, so a
  1362. # repeat is not deduplication but the point (see _roster_for_room).
  1363. # The gateway prepends channel_context above the message text, so the
  1364. # roster never mingles with what the user actually wrote.
  1365. roster = ""
  1366. if room_kind != RoomKind.DM:
  1367. logger.debug(
  1368. "Chatto: fetching roster for room %s",
  1369. room_id,
  1370. )
  1371. roster = await self._roster_for_room(client, room_id=room_id)
  1372. if not roster:
  1373. logger.debug(
  1374. "Chatto: dispatching channel message without roster "
  1375. "(lookup failed or empty)"
  1376. )
  1377. source = self.build_source(
  1378. chat_id=room_id,
  1379. chat_name=self._room_names.get(message.room_id),
  1380. chat_type=chat_type_for_room_kind(room_kind),
  1381. user_id=message.actor_id,
  1382. user_name=user.login, # use login, because display_name is changeable by anyone.
  1383. thread_id=thread_id,
  1384. message_id=source_message_id,
  1385. role_authorized=True,
  1386. )
  1387. # prepare a MessageEvent
  1388. message_event = MessageEvent(
  1389. text=message_body,
  1390. source=source,
  1391. message_id=message.id,
  1392. timestamp=message.created_at or datetime.now(UTC),
  1393. raw_message=message,
  1394. reply_to_message_id=message.in_reply_to,
  1395. channel_context=roster or None,
  1396. )
  1397. if message_event.is_command():
  1398. message_event.message_type = MessageType.COMMAND
  1399. # Attachments — download and hand the local cache paths to the gateway,
  1400. # which runs vision enrichment / document extraction off media_urls.
  1401. (
  1402. message_event.media_urls,
  1403. message_event.media_types,
  1404. media_kinds,
  1405. ) = await self._cache_attachments(list(message.attachments))
  1406. if media_kinds:
  1407. # Same precedence as the Teams/Signal adapters: document-context
  1408. # injection gates strictly on DOCUMENT, image handling keys off the
  1409. # per-path image/* MIME regardless of message_type.
  1410. message_event.message_type = self._message_type_for_media_kinds(media_kinds)
  1411. else:
  1412. message_event.message_type = MessageType.TEXT
  1413. logger.debug("Chatto: MessageEvent: %s", message_event)
  1414. return message_event
  1415. async def _forward_reaction(
  1416. self,
  1417. event: RealtimeEvent,
  1418. payload: Any,
  1419. *,
  1420. removed: bool,
  1421. ) -> None:
  1422. """Forward a human reaction to the gateway's reaction hook surface.
  1423. The handler is registered by the gateway via ``set_reaction_handler``
  1424. and fans out as ``reaction:added`` / ``reaction:removed`` through the
  1425. HookRegistry. The dict shape mirrors the Slack adapter's — hook
  1426. consumers are written against that contract, not against a per-platform
  1427. one. Our own lifecycle reactions (👀/✅/❌) are dropped: forwarding them
  1428. would feed the agent its own markers.
  1429. """
  1430. actor_id = event.actor_id
  1431. if actor_id and self.me and actor_id == self.me.id:
  1432. return
  1433. if not payload.room_id or not payload.message_event_id or not actor_id:
  1434. return
  1435. handler = self._reaction_handler
  1436. if handler is None:
  1437. return
  1438. action = "removed" if removed else "added"
  1439. try:
  1440. await handler(
  1441. {
  1442. "platform": ChattoConstants.PLATFORM_NAME,
  1443. "event_name": f"reaction:{action}",
  1444. "reaction": payload.emoji,
  1445. "user_id": actor_id,
  1446. "item_user_id": None,
  1447. "item_type": "message",
  1448. "channel_id": payload.room_id,
  1449. "message_ts": payload.message_event_id,
  1450. "event_ts": event.id,
  1451. "raw_event": event,
  1452. },
  1453. )
  1454. except Exception: # pragma: no cover - the hook contract is non-blocking
  1455. logger.debug("Chatto: reaction hook forwarding failed", exc_info=True)
  1456. async def _handle_realtime_event(self, event: RealtimeEvent) -> None:
  1457. if self._is_seen(event.id):
  1458. return
  1459. if event.actor_id is None:
  1460. return
  1461. logger.debug("EVENT happened: '%s' from %s", event.kind, event.actor_id)
  1462. # Self-event filter — the actor_id on the envelope is authoritative.
  1463. actor_id = event.actor_id
  1464. if actor_id and self.me and actor_id == self.me.id:
  1465. return
  1466. if event.kind == "message_posted":
  1467. await self._dispatch_message_posted(event.payload, event.id)
  1468. elif event.kind == "message_edited":
  1469. await self._dispatch_message_edited(event.payload)
  1470. elif event.kind in ("reaction_added", "reaction_removed"):
  1471. await self._forward_reaction(
  1472. event,
  1473. event.payload,
  1474. removed=event.kind == "reaction_removed",
  1475. )
  1476. elif event.kind == "presence_changed":
  1477. # A participant's presence changed — patch the cached user in place.
  1478. if (user := self._user_cache.get(event.payload.user_id)) is not None:
  1479. user.presence_status = event.payload.status
  1480. logger.debug(
  1481. "Chatto: patched presence of %s to %s",
  1482. event.payload.user_id,
  1483. event.payload.status.name,
  1484. )
  1485. elif event.kind in ("user_joined_room", "user_left_room"):
  1486. # Membership moved — the room's roster projection is now wrong, so
  1487. # discard it; the next turn there refetches once and re-delivers.
  1488. if event.payload.room_id in self._roster_announced:
  1489. logger.debug(
  1490. "Chatto: %s invalidated the roster of room %s",
  1491. event.kind,
  1492. event.payload.room_id,
  1493. )
  1494. self._evict_roster(event.payload.room_id)
  1495. # confirmed:
  1496. # NOTE: "projection_event" (and "caught_up") belong to realtime protocol
  1497. # v2 on unreleased Chatto main — chattolib speaks v1 and can never
  1498. # deliver them here. Revisit when vendored chattolib gains v2 typing.
  1499. elif event.kind in (
  1500. "mention_notification",
  1501. "notification_dismissed",
  1502. "room_marked_as_read",
  1503. "user_typing",
  1504. "notification_created",
  1505. "new_direct_message_notification",
  1506. "message_retracted",
  1507. "thread_created",
  1508. "thread_follow_changed",
  1509. "room_updated",
  1510. "room_groups_updated",
  1511. ):
  1512. logger.debug(
  1513. "Chatto: '%s' event received. Not yet implemented or not needed.",
  1514. event.kind,
  1515. )
  1516. else:
  1517. logger.warning("Chatto: unknown event kind: '%s'", event.kind)
  1518. async def _chattolib_event_loop(self) -> None:
  1519. """Event loop using chattolib's stream_events.
  1520. This replaces the manual WebSocket loop with chattolib's high-level
  1521. stream_events() which provides pre-decoded RealtimeEvent objects.
  1522. """
  1523. delay = ChattoConstants.WS_RECONNECT_INITIAL_BACKOFF
  1524. while not self._closing:
  1525. try:
  1526. client = await self._require_chatto_client()
  1527. await self._refresh_rooms()
  1528. logger.info(
  1529. "Chatto: starting chattolib event stream with %d rooms",
  1530. len(self._joined_room_ids),
  1531. )
  1532. async for event in stream_events(client):
  1533. if self._closing:
  1534. return
  1535. if isinstance(event, RealtimeEvent):
  1536. await self._handle_realtime_event(event)
  1537. # Iterator exited cleanly — treat as a normal close and reconnect
  1538. # with the local backoff (no server hint available).
  1539. logger.info("Chatto: realtime stream ended, reconnecting")
  1540. except asyncio.CancelledError:
  1541. return
  1542. except ChattoRealtimeCloseError as exc:
  1543. if not exc.reconnect:
  1544. logger.error(
  1545. "Chatto: realtime closed by server (%s: %s), not reconnecting",
  1546. exc.code,
  1547. exc.message,
  1548. )
  1549. return
  1550. wait = max(
  1551. exc.retry_after_ms / 1000.0,
  1552. ChattoConstants.WS_RECONNECT_INITIAL_BACKOFF,
  1553. )
  1554. logger.warning(
  1555. "Chatto: realtime closed by server (%s), reconnecting in %.1fs",
  1556. exc.code,
  1557. wait,
  1558. )
  1559. delay = (
  1560. ChattoConstants.WS_RECONNECT_INITIAL_BACKOFF
  1561. ) # server hint supersedes local backoff
  1562. await self._sleep_interruptible(wait)
  1563. continue
  1564. except ChattoRealtimeError as exc:
  1565. if exc.fatal:
  1566. logger.error(
  1567. "Chatto: fatal realtime error (%s): %s", exc.code, exc.message
  1568. )
  1569. return
  1570. logger.warning(
  1571. "Chatto: realtime error (%s: %s), reconnecting in %.1fs",
  1572. exc.code,
  1573. exc.message,
  1574. delay,
  1575. )
  1576. except Exception as exc:
  1577. cause = exc.__cause__
  1578. if cause is not None and isinstance(cause, TypeError):
  1579. logger.warning(
  1580. "Chatto: realtime wire corruption (%s), reconnecting in %.1fs",
  1581. cause,
  1582. delay,
  1583. )
  1584. else:
  1585. logger.warning(
  1586. "Chatto: unexpected realtime error: %s, reconnecting in %.1fs",
  1587. exc,
  1588. delay,
  1589. )
  1590. if self._closing:
  1591. return
  1592. jitter = delay * 0.2 * random.random()
  1593. await self._sleep_interruptible(delay + jitter)
  1594. delay = min(delay * 2, ChattoConstants.WS_RECONNECT_MAX_BACKOFF)
  1595. async def _sleep_interruptible(self, seconds: float) -> None:
  1596. """Sleep in short slices so disconnect() cancels promptly."""
  1597. end = asyncio.get_running_loop().time() + seconds
  1598. while not self._closing:
  1599. remaining = end - asyncio.get_running_loop().time()
  1600. if remaining <= 0:
  1601. return
  1602. await asyncio.sleep(min(remaining, 0.5))
  1603. def _warn_if_home_channel_unjoined(self, member_ids: set[str]) -> None:
  1604. """Warn once when CHATTO_HOME_CHANNEL names a room the bot is not in.
  1605. Standalone cron delivery posts straight into that room with a fresh
  1606. client and no join logic of its own — without server-side membership
  1607. every proactive send fails there.
  1608. """
  1609. home_id = (self.chatto_config.home_channel.value or "").strip()
  1610. if (
  1611. not home_id
  1612. or home_id in member_ids
  1613. or home_id in self._joined_room_ids
  1614. or self._home_warning_logged
  1615. ):
  1616. return
  1617. self._home_warning_logged = True
  1618. logger.warning(
  1619. "Chatto: CHATTO_HOME_CHANNEL '%s' is not a joined room - cron and "
  1620. "notification delivery will fail until the bot joins it (invite "
  1621. "the account natively in Chatto, or DM it '/join').",
  1622. home_id,
  1623. )
  1624. async def _resolve_dm_partner(self, client: ChattoClient, room_id: str) -> None:
  1625. """Cache the chat partner's login for a DM room, best-effort.
  1626. DM rooms carry no usable name of their own, so the joined-rooms
  1627. summary names the other side instead. Resolved once per room and
  1628. session; without our own user (pre-connect) or on a directory error
  1629. the raw room name stays in place.
  1630. """
  1631. if room_id in self._dm_partners or self.me is None:
  1632. return
  1633. try:
  1634. member_ids, _page = await client.list_room_members(room_id)
  1635. if not member_ids:
  1636. return
  1637. members = await client.batch_get_room_members(room_id, member_ids)
  1638. except Exception:
  1639. logger.debug(
  1640. "Chatto: could not list members of DM %s", room_id, exc_info=True
  1641. )
  1642. return
  1643. for member in members:
  1644. user = member.user
  1645. if user and user.id != self.me.id:
  1646. self._dm_partners[room_id] = user.login or user.display_name
  1647. return
  1648. def _joined_room_label(self, room_id: str) -> str:
  1649. """One entry for the joined-rooms summary, with how-we-answer tags.
  1650. The tag names the room's participation mode at a glance: DMs answer
  1651. unconditionally, listed channels answer on mentions or every message,
  1652. and an unlisted channel is read-only. A DM is labelled with its chat
  1653. partner's login rather than its (empty) room name, and names are
  1654. quoted so empty strings and spaces stay visible.
  1655. """
  1656. name = self._room_names.get(room_id, room_id)
  1657. if self._room_kinds.get(room_id) == RoomKind.DM:
  1658. tags = "[dm]"
  1659. name = self._dm_partners.get(room_id) or name
  1660. else:
  1661. policy = self._room_policy(room_id)
  1662. if policy == RoomPolicy.REQUIRE_MENTION:
  1663. tags = "[on-mention]"
  1664. elif policy == RoomPolicy.OPEN:
  1665. tags = "[every-message]"
  1666. else:
  1667. tags = "[read-only]"
  1668. if room_id in self._universal_room_ids:
  1669. tags += " [universal]"
  1670. return f'"{name}" ({room_id}) {tags}'
  1671. def _log_joined_rooms(self) -> None:
  1672. """Log the joined rooms and how the bot answers in each of them."""
  1673. labels = [self._joined_room_label(rid) for rid in self._joined_room_ids]
  1674. logger.info(
  1675. "Chatto WS: currently joined in %d room(s): %s",
  1676. len(labels),
  1677. ", ".join(labels),
  1678. )
  1679. async def _refresh_rooms(self) -> None:
  1680. """Refresh room list via ConnectRPC, join and seed any newly discovered rooms.
  1681. Runs once per (re)connect, so this is also the bootstrap point for the
  1682. roster projection: protocol v1 sends no presence snapshot on subscribe,
  1683. meaning cached presence could only heal on the next random change.
  1684. Discarding it makes each announced room's first turn refetch fresh
  1685. data — lazily, so quiet rooms stay free of API calls.
  1686. """
  1687. client = await self._get_chatto_client()
  1688. if client is None:
  1689. logger.warning("Chatto WS: _refresh_rooms aborted - no client available")
  1690. return
  1691. for announced in list(self._roster_announced):
  1692. self._evict_roster(announced)
  1693. try:
  1694. rooms_list = await client.list_rooms()
  1695. member_ids: set[str] = set()
  1696. new_room_ids: list[str] = []
  1697. for room_with_state in rooms_list:
  1698. if not room_with_state:
  1699. continue
  1700. room_obj = room_with_state.room or None
  1701. if not room_obj:
  1702. continue
  1703. self._room_names[room_obj.id] = room_obj.name
  1704. self._room_kinds[room_obj.id] = room_obj.kind
  1705. if room_obj.kind == RoomKind.DM:
  1706. await self._resolve_dm_partner(client, room_obj.id)
  1707. if room_obj.universal:
  1708. self._universal_room_ids.add(room_obj.id)
  1709. if not room_with_state.viewer_state.is_member:
  1710. continue
  1711. member_ids.add(room_obj.id)
  1712. if room_obj.id not in self._joined_room_ids:
  1713. new_room_ids.append(room_obj.id)
  1714. # Left via /leave, kicked, deleted: rooms we no longer belong
  1715. # to drop out here — otherwise the next refresh would quietly
  1716. # re-add what /leave just removed.
  1717. stale_room_ids = [
  1718. rid for rid in self._joined_room_ids if rid not in member_ids
  1719. ]
  1720. for rid in stale_room_ids:
  1721. self._joined_room_ids.remove(rid)
  1722. self._evict_roster(rid)
  1723. if stale_room_ids:
  1724. logger.info(
  1725. "Chatto WS: no longer a member of %d room(s): %s",
  1726. len(stale_room_ids),
  1727. stale_room_ids,
  1728. )
  1729. self._warn_if_home_channel_unjoined(member_ids)
  1730. if new_room_ids:
  1731. logger.info(
  1732. "Chatto WS: discovered %d new room(s): %s",
  1733. len(new_room_ids),
  1734. new_room_ids,
  1735. )
  1736. for rid in new_room_ids:
  1737. # Membership came straight from the directory scan
  1738. # (viewer_state.is_member); natively invited rooms need no
  1739. # JoinRoom call — same rule as _run_join.
  1740. if self._answers_in_room(rid):
  1741. await self._seed_room(rid)
  1742. else:
  1743. logger.info(
  1744. "Chatto WS: %s (%s) is on neither mention list -"
  1745. " joined read-only",
  1746. self._room_names.get(rid, rid),
  1747. rid,
  1748. )
  1749. self._joined_room_ids.append(rid)
  1750. self._log_joined_rooms()
  1751. except Exception:
  1752. logger.warning("Chatto WS: _refresh_rooms failed", exc_info=True)
  1753. # ------------------------------------------------------------------ #
  1754. # Read state (best-effort, Chatto-unique)
  1755. # ------------------------------------------------------------------ #
  1756. # Deliberately outside the try block above: this sweep runs on every
  1757. # refresh — including refreshes that discovered no new rooms. The
  1758. # old early return starved it to "only when something changed",
  1759. # leaving silent rooms unmarked for whole connection lifetimes.
  1760. # Best-effort: mark all joined rooms as read (room_id may be undefined here)
  1761. for _rid in list(self._joined_room_ids):
  1762. try:
  1763. await client.mark_room_as_read(room_id=_rid)
  1764. except Exception:
  1765. logger.debug(
  1766. "Chatto: mark_room_as_read failed for %s", _rid, exc_info=True
  1767. )
  1768. # ------------------------------------------------------------------ #
  1769. # Sending (ConnectRPC — unchanged from polling version)
  1770. # ------------------------------------------------------------------ #
  1771. def _resolve_outbound_thread(
  1772. self,
  1773. chat_id: str,
  1774. reply_to: str | None,
  1775. metadata: dict[str, Any] | None,
  1776. ) -> str | None:
  1777. """The thread an outbound message belongs in, or ``None`` for the room.
  1778. One definition for every outbound path (text and attachments):
  1779. ``metadata["thread_id"]`` wins — it names the thread root to stay in;
  1780. otherwise ``reply_to`` anchors a thread under the incoming message,
  1781. but only with ``auto_thread`` enabled. DMs never thread.
  1782. Deliberately does NOT open an auto-thread of its own — only the
  1783. chunked text path in :meth:`send` does that, because only there do
  1784. further chunks follow into the freshly created thread.
  1785. """
  1786. thread_id = (metadata or {}).get("thread_id")
  1787. if reply_to and self.chatto_config.auto_thread.value and not thread_id:
  1788. thread_id = reply_to
  1789. if self._room_kinds.get(chat_id) == RoomKind.DM:
  1790. return None
  1791. return str(thread_id) if thread_id else None
  1792. async def send(
  1793. self,
  1794. chat_id: str,
  1795. content: str,
  1796. reply_to: str | None = None,
  1797. metadata: dict[str, Any] | None = None,
  1798. ) -> SendResult:
  1799. """Send a message to a Chatto room.
  1800. Long messages are split into chunks via ``truncate_message`` and
  1801. each chunk is sent as a separate CreateMessage call. The first
  1802. chunk's message ID is returned as ``message_id``.
  1803. When ``auto_thread`` is enabled and the incoming message was a
  1804. regular room message (not already in a thread), the first chunk is
  1805. sent as a room message and its ID becomes the thread root. Subsequent
  1806. chunks are sent in that thread. This mirrors Discord's auto_thread
  1807. behavior.
  1808. BasePlatformAdapter override
  1809. """
  1810. if not content:
  1811. return SendResult(success=False, error="Empty message")
  1812. formatted = _normalise_outbound_text(content)
  1813. chunks = self.truncate_message(formatted, ChattoConstants.SPLIT_THRESHOLD)
  1814. thread_id = self._resolve_outbound_thread(chat_id, reply_to, metadata)
  1815. room_kind = self._room_kinds.get(chat_id)
  1816. is_dm = room_kind == RoomKind.DM
  1817. # Auto-thread: by default, Chatto creates a thread for replies to room
  1818. # messages (not DMs, not already in a thread). This keeps conversations
  1819. # organized in the room. Can be disabled via extra.auto_thread=false.
  1820. use_auto_thread = (
  1821. self.chatto_config.auto_thread.value and not thread_id and not is_dm
  1822. )
  1823. message_ids: list[str] = []
  1824. last_error: str | None = None
  1825. retryable = False
  1826. client = await self._get_chatto_client()
  1827. if client is None:
  1828. return SendResult(
  1829. success=False, error="Chatto client not available", retryable=True
  1830. )
  1831. for i, chunk in enumerate(chunks):
  1832. try:
  1833. msg_obj = await client.post_message(
  1834. room_id=chat_id,
  1835. body=chunk,
  1836. thread_root_event_id=str(thread_id) if thread_id else "",
  1837. )
  1838. except Exception as e:
  1839. # ChattoError included: both read as "this chunk did not go
  1840. # out" and stop the batch — the SendResult carries the reason.
  1841. # Only server-side failures count as retryable; a bug in our
  1842. # own code must not read as a transient network blip.
  1843. last_error = str(e)
  1844. retryable = isinstance(e, ChattoError)
  1845. break
  1846. self._mark_seen(msg_obj.id)
  1847. message_ids.append(msg_obj.id)
  1848. # Auto-thread: first chunk becomes the thread root,
  1849. # subsequent chunks go in the thread
  1850. if use_auto_thread and i == 0 and not thread_id:
  1851. thread_id = msg_obj.id
  1852. # Nothing got through at all — report the failure instead of a phantom success.
  1853. if not message_ids:
  1854. return SendResult(
  1855. success=False,
  1856. error=last_error or "Chatto: message could not be sent",
  1857. retryable=retryable,
  1858. )
  1859. first_id = message_ids[0]
  1860. # ------------------------------------------------------------------ #
  1861. # Thread following (best-effort, Chatto-unique)
  1862. # ------------------------------------------------------------------ #
  1863. if thread_id:
  1864. try:
  1865. await client.follow_thread(chat_id, thread_id)
  1866. except Exception:
  1867. logger.debug(
  1868. "Chatto: follow_thread failed for %s/%s",
  1869. chat_id,
  1870. thread_id,
  1871. exc_info=True,
  1872. )
  1873. # A later chunk failed after earlier ones went out: partial delivery.
  1874. if last_error:
  1875. logger.warning(
  1876. "Chatto: sent %d/%d chunk(s) to %s before failing: %s",
  1877. len(message_ids),
  1878. len(chunks),
  1879. chat_id,
  1880. last_error,
  1881. )
  1882. # raw_response stays unset (dict-shaped per the SendResult contract):
  1883. # gateway consumers such as the cron scheduler call .get() on it, so a
  1884. # chattolib Message here would crash delivery bookkeeping *after* the
  1885. # send already succeeded — the job then falls back to the standalone
  1886. # path and the room sees the message twice.
  1887. return SendResult(success=True, message_id=first_id)
  1888. def format_message(self, content: str) -> str:
  1889. """Normalise outgoing text for Chatto.
  1890. The transformations live in :func:`_normalise_outbound_text`, shared
  1891. with the standalone cron sender.
  1892. BasePlatformAdapter override
  1893. """
  1894. return _normalise_outbound_text(content)
  1895. async def edit_message(
  1896. self,
  1897. chat_id: str,
  1898. message_id: str,
  1899. content: str,
  1900. *,
  1901. finalize: bool = False,
  1902. ) -> SendResult:
  1903. """Edit a message we previously sent, via MessageService/UpdateMessage.
  1904. The stream consumer drives streaming replies through this: without the
  1905. override the base class reports "Not supported" and every incremental
  1906. update arrives as a *new* message.
  1907. ``finalize`` is a no-op for Chatto — an edit is an edit here, there is
  1908. no in-progress card state to close out (hence no
  1909. ``REQUIRES_EDIT_FINALIZE``).
  1910. Content that exceeds the per-message limit is refused rather than
  1911. silently truncated, so the caller falls back to ``send()``, which
  1912. splits across messages.
  1913. BasePlatformAdapter override
  1914. """
  1915. if not message_id:
  1916. return SendResult(success=False, error="Chatto: no message id to edit")
  1917. if not content:
  1918. return SendResult(success=False, error="Empty message")
  1919. formatted = self.format_message(content)
  1920. if len(formatted) > ChattoConstants.MAX_MESSAGE_LENGTH:
  1921. # Refuse instead of truncating: the caller's fallback path splits.
  1922. return SendResult(
  1923. success=False,
  1924. error=(
  1925. f"Chatto: edit exceeds {ChattoConstants.MAX_MESSAGE_LENGTH} "
  1926. f"chars ({len(formatted)})"
  1927. ),
  1928. )
  1929. client = await self._get_chatto_client()
  1930. if client is None:
  1931. return SendResult(
  1932. success=False, error="Chatto client not available", retryable=True
  1933. )
  1934. try:
  1935. msg = await client.update_message(
  1936. room_id=str(chat_id),
  1937. event_id=str(message_id),
  1938. body=formatted,
  1939. )
  1940. except ChattoError as e:
  1941. logger.warning("Chatto: UpdateMessage failed for %s: %s", message_id, e)
  1942. return SendResult(success=False, error=str(e), retryable=True)
  1943. except Exception as e:
  1944. logger.warning("Chatto: UpdateMessage error for %s: %s", message_id, e)
  1945. return SendResult(success=False, error=str(e), retryable=False)
  1946. # Our own edit comes back as a message_edited event; mark it seen so it
  1947. # is never mistaken for inbound traffic.
  1948. edited_id = msg.id or str(message_id)
  1949. self._mark_seen(edited_id)
  1950. return SendResult(success=True, message_id=edited_id)
  1951. async def delete_message(self, chat_id: str, message_id: str) -> bool:
  1952. """Delete a message via MessageService/DeleteMessage.
  1953. Used by the stream consumer's fresh-final cleanup (removing a preview
  1954. message once the completed reply has been sent) and by the ephemeral
  1955. reply TTL.
  1956. BasePlatformAdapter override
  1957. """
  1958. if not chat_id or not message_id:
  1959. return False
  1960. client = await self._get_chatto_client()
  1961. if client is None:
  1962. logger.warning("Chatto: DeleteMessage — client unavailable")
  1963. return False
  1964. try:
  1965. return bool(
  1966. await client.delete_message(
  1967. room_id=str(chat_id),
  1968. event_id=str(message_id),
  1969. )
  1970. )
  1971. except ChattoError as e:
  1972. logger.warning("Chatto: DeleteMessage failed for %s: %s", message_id, e)
  1973. return False
  1974. except Exception as e:
  1975. logger.warning("Chatto: DeleteMessage error for %s: %s", message_id, e)
  1976. return False
  1977. async def create_handoff_thread(
  1978. self,
  1979. parent_chat_id: str,
  1980. name: str,
  1981. ) -> str | None:
  1982. """Anchor a session handoff in a fresh thread under *parent_chat_id*.
  1983. Chatto threads hang off a message, not off the room, so we post a seed
  1984. message and hand its ID back as the thread root — the same shape the
  1985. Slack adapter uses. DMs don't support threads, so they get ``None``
  1986. and the watcher keeps delivering into the DM itself.
  1987. BasePlatformAdapter override
  1988. """
  1989. if not parent_chat_id:
  1990. return None
  1991. if self._room_kinds.get(parent_chat_id) == RoomKind.DM:
  1992. logger.debug("Chatto: handoff thread skipped — %s is a DM", parent_chat_id)
  1993. return None
  1994. client = await self._get_chatto_client()
  1995. if client is None:
  1996. logger.warning("Chatto: handoff thread — client unavailable")
  1997. return None
  1998. try:
  1999. msg = await client.post_message(
  2000. room_id=str(parent_chat_id),
  2001. body=f"🧵 Hermes handoff — **{(name or 'session').strip()[:80]}**",
  2002. )
  2003. except Exception as e:
  2004. logger.warning(
  2005. "Chatto: handoff thread seed-post failed for room %s: %s",
  2006. parent_chat_id,
  2007. e,
  2008. )
  2009. return None
  2010. seed_id = msg.id
  2011. if not seed_id:
  2012. logger.warning("Chatto: handoff thread seed-post returned no message id")
  2013. return None
  2014. self._mark_seen(seed_id)
  2015. try:
  2016. await client.follow_thread(str(parent_chat_id), seed_id)
  2017. except Exception:
  2018. logger.debug(
  2019. "Chatto: follow_thread failed for handoff %s/%s",
  2020. parent_chat_id,
  2021. seed_id,
  2022. exc_info=True,
  2023. )
  2024. return seed_id
  2025. # Overridden from BaseAdapter:
  2026. async def send_typing(self, chat_id: str, metadata=None) -> None:
  2027. """Start a persistent typing indicator for a room.
  2028. Sends a typing ping every 10 seconds (Chatto's indicator likely
  2029. lasts ~8-10s). The background loop runs until ``stop_typing()``
  2030. is called or the task is cancelled.
  2031. BasePlatformAdapter override
  2032. """
  2033. if chat_id in self._typing_tasks:
  2034. return # already running
  2035. async def _typing_loop() -> None:
  2036. try:
  2037. while True:
  2038. try:
  2039. client = await self._get_chatto_client()
  2040. if client is None:
  2041. return
  2042. await client.refresh_typing_indicator(room_id=str(chat_id))
  2043. except asyncio.CancelledError:
  2044. return
  2045. except Exception:
  2046. logger.debug(
  2047. "Chatto: typing indicator refresh failed for %s",
  2048. chat_id,
  2049. exc_info=True,
  2050. )
  2051. await asyncio.sleep(10)
  2052. except asyncio.CancelledError:
  2053. pass
  2054. finally:
  2055. self._typing_tasks.pop(chat_id, None)
  2056. self._typing_tasks[chat_id] = asyncio.create_task(_typing_loop())
  2057. async def stop_typing(self, chat_id: str) -> None:
  2058. """Stop the persistent typing indicator for a room.
  2059. BasePlatformAdapter override
  2060. """
  2061. task = self._typing_tasks.pop(chat_id, None)
  2062. if task:
  2063. task.cancel()
  2064. try:
  2065. await task
  2066. except (asyncio.CancelledError, Exception):
  2067. logger.debug("Chatto: typing task ended for %s", chat_id, exc_info=True)
  2068. async def get_chat_info(self, chat_id: str) -> dict[str, Any]:
  2069. """Get information about a chat/room.
  2070. BasePlatformAdapter override
  2071. """
  2072. name = self._room_names.get(chat_id, chat_id)
  2073. kind = self._room_kinds.get(chat_id)
  2074. return {
  2075. "name": name,
  2076. "type": chat_type_for_room_kind(kind).value,
  2077. }
  2078. # ------------------------------------------------------------------ #
  2079. # Reactions
  2080. # ------------------------------------------------------------------ #
  2081. @staticmethod
  2082. def _emoji_to_shortcode(emoji: str) -> str:
  2083. """Convert a unicode emoji to a Chatto shortcode name.
  2084. If the emoji is already a shortcode (no unicode mapping found),
  2085. return it as-is.
  2086. """
  2087. shortcode = ChattoConstants.EMOJI_TO_SHORTCODE.get(emoji)
  2088. if shortcode:
  2089. return shortcode
  2090. # Already a shortcode like "thumbsup" — return as-is
  2091. return emoji
  2092. async def add_reaction(self, room_id: str, message_id: str, emoji: str) -> bool:
  2093. """Add a reaction to a message via MessageService/AddReaction.
  2094. BasePlatformAdapter override
  2095. """
  2096. shortcode = self._emoji_to_shortcode(emoji)
  2097. client = await self._get_chatto_client()
  2098. if client is None:
  2099. logger.warning("Chatto: AddReaction — client unavailable")
  2100. return False
  2101. try:
  2102. result = await client.add_reaction(
  2103. room_id=room_id,
  2104. message_event_id=message_id,
  2105. emoji=shortcode,
  2106. )
  2107. return result
  2108. except ChattoError as e:
  2109. logger.warning("Chatto: AddReaction failed: %s", e)
  2110. return False
  2111. except Exception as e:
  2112. logger.warning("Chatto: AddReaction error: %s", e)
  2113. return False
  2114. async def remove_reaction(self, chat_id: str, message_id: str, emoji: str) -> bool:
  2115. """Remove a reaction from a message via MessageService/RemoveReaction.
  2116. BasePlatformAdapter override
  2117. """
  2118. shortcode = self._emoji_to_shortcode(emoji)
  2119. client = await self._get_chatto_client()
  2120. if client is None:
  2121. logger.warning("Chatto: RemoveReaction — client unavailable")
  2122. return False
  2123. try:
  2124. result = await client.remove_reaction(
  2125. room_id=str(chat_id),
  2126. message_event_id=str(message_id),
  2127. emoji=shortcode,
  2128. )
  2129. return result
  2130. except ChattoError as e:
  2131. logger.warning("Chatto: RemoveReaction failed: %s", e)
  2132. return False
  2133. except Exception as e:
  2134. logger.warning("Chatto: RemoveReaction error: %s", e)
  2135. return False
  2136. # ------------------------------------------------------------------ #
  2137. # DM initiation (Chatto-unique)
  2138. # ------------------------------------------------------------------ #
  2139. async def start_dm(self, user_id: str) -> str | None:
  2140. """Start a direct message with a user via RoomService/StartDM.
  2141. Returns the room ID on success, or None on failure.
  2142. BasePlatformAdapter override
  2143. """
  2144. if not user_id:
  2145. return None
  2146. client = await self._get_chatto_client()
  2147. if client is None:
  2148. return None
  2149. try:
  2150. room = await client.start_dm(participant_ids=[str(user_id)])
  2151. self._room_names[room.id] = room.name
  2152. self._room_kinds[room.id] = room.kind
  2153. return room.id
  2154. except ChattoError as e:
  2155. logger.debug("Chatto: StartDM failed: %s", e)
  2156. return None
  2157. except Exception as e:
  2158. logger.debug("Chatto: StartDM error: %s", e)
  2159. return None
  2160. # ------------------------------------------------------------------ #
  2161. # Room creation (Chatto-unique)
  2162. # ------------------------------------------------------------------ #
  2163. async def create_room(
  2164. self,
  2165. name: str,
  2166. description: str = "",
  2167. group_id: str = "",
  2168. universal: bool = True,
  2169. ) -> str | None:
  2170. """Create an ad-hoc room via RoomService/CreateRoom.
  2171. Returns the room ID on success, or None on failure.
  2172. BasePlatformAdapter override
  2173. """
  2174. client = await self._get_chatto_client()
  2175. if client is None:
  2176. return None
  2177. try:
  2178. room = await client.create_room(
  2179. name=name,
  2180. group_id=group_id or "",
  2181. description=description,
  2182. universal=universal,
  2183. )
  2184. rid = str(room.id) if room else ""
  2185. if rid:
  2186. self._room_names[rid] = room.name
  2187. self._room_kinds[rid] = room.kind
  2188. return rid
  2189. logger.debug("Chatto: CreateRoom returned no room id")
  2190. return None
  2191. except ChattoError as e:
  2192. logger.debug("Chatto: CreateRoom failed: %s", e)
  2193. return None
  2194. except Exception as e:
  2195. logger.debug("Chatto: CreateRoom error: %s", e)
  2196. return None
  2197. # ------------------------------------------------------------------ #
  2198. # Processing lifecycle hooks (reactions-based, like Discord)
  2199. # ------------------------------------------------------------------ #
  2200. def _event_room_and_message_id(self, event: MessageEvent) -> tuple[str, str]:
  2201. """Extract room_id and message_id from a MessageEvent."""
  2202. message_id = event.message_id or ""
  2203. return event.source.chat_id, message_id
  2204. async def on_processing_start(self, event: MessageEvent) -> None:
  2205. """Record the turn as open, then add an 👀 (eyes) reaction.
  2206. The record is what lets an edit of this very message be recognized as
  2207. a mid-run correction. It must be maintained even when reactions are
  2208. disabled — tracking and decorating are independent concerns.
  2209. BasePlatformAdapter override
  2210. """
  2211. session_key = (
  2212. self._session_key_for(event.source) if event.source is not None else ""
  2213. )
  2214. if session_key and event.message_id:
  2215. self._processing[session_key] = str(event.message_id)
  2216. if not self.chatto_config.reactions.value:
  2217. return
  2218. chat_id, message_id = self._event_room_and_message_id(event)
  2219. if not chat_id or not message_id:
  2220. # Routine, not a fault: the gateway runs agent-initiated turns
  2221. # (heartbeat polls, goal continuations) through the same pipeline
  2222. # with message_id=None, and there is no inbound message to mark.
  2223. logger.debug(
  2224. "Chatto: nothing to react to (chat_id=%r, message_id=%r)",
  2225. chat_id,
  2226. message_id,
  2227. )
  2228. return
  2229. await self.add_reaction(chat_id, message_id, "👀")
  2230. async def on_processing_complete(
  2231. self,
  2232. event: MessageEvent,
  2233. outcome: ProcessingOutcome,
  2234. ) -> None:
  2235. """Close the turn's record, then swap 👀 for ✅/❌/🚫.
  2236. Fires for every outcome, including CANCELLED (mid-run edit
  2237. correction) — this is what re-arms `_processing` before the corrected
  2238. turn is dispatched.
  2239. BasePlatformAdapter override
  2240. """
  2241. session_key = (
  2242. self._session_key_for(event.source) if event.source is not None else ""
  2243. )
  2244. if session_key and event.message_id:
  2245. self._processing.pop(session_key, None)
  2246. if not self.chatto_config.reactions.value:
  2247. return
  2248. chat_id, message_id = self._event_room_and_message_id(event)
  2249. if not chat_id or not message_id:
  2250. return
  2251. # Remove the processing eyes reaction
  2252. await self.remove_reaction(chat_id, message_id, "👀")
  2253. # Add the outcome reaction
  2254. if outcome == ProcessingOutcome.SUCCESS:
  2255. await self.add_reaction(chat_id, message_id, "✅")
  2256. elif outcome == ProcessingOutcome.FAILURE:
  2257. await self.add_reaction(chat_id, message_id, "❌")
  2258. elif outcome == ProcessingOutcome.CANCELLED:
  2259. await self.add_reaction(chat_id, message_id, "🚫")
  2260. # ------------------------------------------------------------------ #
  2261. # Asset upload (chunked)
  2262. # ------------------------------------------------------------------ #
  2263. async def _upload_asset(self, room_id: str, file_path: str) -> str | None:
  2264. """Upload a file via the chunked AssetUploadService.
  2265. Returns the asset ID on success, or None on failure.
  2266. """
  2267. try:
  2268. file_data = await asyncio.to_thread(_read_file_bytes, file_path)
  2269. except Exception as e:
  2270. logger.error("Chatto: failed to read file %s — %s", file_path, e)
  2271. return None
  2272. if not file_data:
  2273. logger.error("Chatto: file %s is empty", file_path)
  2274. return None
  2275. file_size = len(file_data)
  2276. file_name = os.path.basename(file_path)
  2277. mime_type = mimetypes.guess_type(file_path)[0] or "application/octet-stream"
  2278. sha256_hash = hashlib.sha256(file_data).hexdigest()
  2279. client = await self._get_chatto_client()
  2280. if client is None:
  2281. logger.error("Chatto: upload aborted - no client available")
  2282. return None
  2283. try:
  2284. # Step 1: Create upload session
  2285. upload = await client.create_upload(
  2286. room_id=room_id,
  2287. filename=file_name,
  2288. size=file_size,
  2289. sha256=sha256_hash,
  2290. content_type=mime_type,
  2291. )
  2292. # AssetUpload names this upload_id, not id — reading it through an
  2293. # untyped getattr default is what let the mismatch reach production.
  2294. upload_id = upload.upload_id
  2295. if not upload_id:
  2296. logger.error("Chatto: CreateUpload returned no upload ID")
  2297. return None
  2298. # Step 2: Upload chunks
  2299. offset = 0
  2300. while offset < file_size:
  2301. chunk = file_data[offset : offset + ChattoConstants.UPLOAD_CHUNK_SIZE]
  2302. chunk_sha256 = hashlib.sha256(chunk).hexdigest()
  2303. await client.upload_chunk(
  2304. upload_id=upload_id,
  2305. offset=offset,
  2306. content=chunk,
  2307. chunk_sha256=chunk_sha256,
  2308. )
  2309. offset += len(chunk)
  2310. # Step 3: Complete upload
  2311. upload, asset = await client.complete_upload(upload_id=upload_id)
  2312. if not asset:
  2313. logger.error("Chatto: CompleteUpload returned no asset")
  2314. return None
  2315. logger.info(
  2316. "Chatto: uploaded %s as asset %s (%d bytes)",
  2317. file_name,
  2318. asset.id,
  2319. file_size,
  2320. )
  2321. return asset.id
  2322. except ChattoError as e:
  2323. logger.error("Chatto: upload failed: %s", e)
  2324. return None
  2325. except Exception as e:
  2326. logger.error("Chatto: upload error: %s", e)
  2327. return None
  2328. async def _post_attachment_message(
  2329. self,
  2330. chat_id: str,
  2331. asset_ids: list[str],
  2332. caption: str | None,
  2333. reply_to: str | None,
  2334. metadata: dict[str, Any] | None,
  2335. ) -> SendResult:
  2336. """Post one message carrying already-uploaded assets.
  2337. Threading follows the same rules as a text send
  2338. (:meth:`_resolve_outbound_thread`), so with ``auto_thread`` disabled an
  2339. attachment reply lands in the room like its text counterpart instead of
  2340. quietly opening a thread.
  2341. """
  2342. thread_id = self._resolve_outbound_thread(chat_id, reply_to, metadata)
  2343. client = await self._get_chatto_client()
  2344. if client is None:
  2345. return SendResult(
  2346. success=False, error="Chatto client not available", retryable=True
  2347. )
  2348. try:
  2349. msg = await client.post_message(
  2350. room_id=str(chat_id),
  2351. body=self.format_message(caption) if caption else "",
  2352. attachment_asset_ids=asset_ids,
  2353. thread_root_event_id=thread_id or "",
  2354. )
  2355. self._mark_seen(msg.id)
  2356. return SendResult(success=True, message_id=msg.id)
  2357. except ChattoError as e:
  2358. # Same classification as the chunked text path: server-side
  2359. # failures retry, anything else is ours and must not loop.
  2360. return SendResult(success=False, error=str(e), retryable=True)
  2361. except Exception as e:
  2362. return SendResult(success=False, error=str(e), retryable=False)
  2363. async def _send_local_file_as_attachment(
  2364. self,
  2365. chat_id: str,
  2366. file_path: str,
  2367. caption: str | None,
  2368. reply_to: str | None,
  2369. metadata: dict[str, Any] | None,
  2370. *,
  2371. kind: str,
  2372. ) -> SendResult:
  2373. """Upload a local file and post it as a native Chatto attachment.
  2374. Shared by ``send_image_file``/``send_document``/``send_video``/
  2375. ``send_voice`` — the upload mechanics are identical, only the wording of
  2376. the failure notice differs. On failure we send that notice as text and
  2377. never the host path (it leaks the Hermes home layout).
  2378. """
  2379. notice = f"⚠️ Couldn't deliver the {kind} attachment."
  2380. safe_path = self.validate_media_delivery_path(file_path)
  2381. if not safe_path:
  2382. logger.warning(
  2383. "[%s] send %s: unsafe path %s",
  2384. self.name,
  2385. kind,
  2386. file_path,
  2387. )
  2388. text = f"{caption}\n{notice}" if caption else notice
  2389. return await self.send(chat_id, text, reply_to=reply_to, metadata=metadata)
  2390. asset_id = await self._upload_asset(str(chat_id), safe_path)
  2391. if not asset_id:
  2392. logger.warning(
  2393. "[%s] send %s: upload failed for %s",
  2394. self.name,
  2395. kind,
  2396. safe_path,
  2397. )
  2398. text = f"{caption}\n{notice}" if caption else notice
  2399. return await self.send(chat_id, text, reply_to=reply_to, metadata=metadata)
  2400. return await self._post_attachment_message(
  2401. chat_id,
  2402. [asset_id],
  2403. caption,
  2404. reply_to,
  2405. metadata,
  2406. )
  2407. async def send_image_file(
  2408. self,
  2409. chat_id: str,
  2410. image_path: str,
  2411. caption: str | None = None,
  2412. reply_to: str | None = None,
  2413. metadata: dict[str, Any] | None = None,
  2414. **kwargs,
  2415. ) -> SendResult:
  2416. """Send a local image file via the chunked upload API.
  2417. The parameter is ``image_path``, not ``file_path``: every caller passes
  2418. it by keyword (``gateway/run.py:22354``, ``:22470``, and the base class's
  2419. own ``send_multiple_images`` file:// branch), so a renamed parameter
  2420. makes each of those raise TypeError and silently degrade to a text
  2421. notice.
  2422. BasePlatformAdapter override
  2423. """
  2424. return await self._send_local_file_as_attachment(
  2425. chat_id,
  2426. image_path,
  2427. caption,
  2428. reply_to,
  2429. metadata,
  2430. kind="image",
  2431. )
  2432. async def send_document(
  2433. self,
  2434. chat_id: str,
  2435. file_path: str,
  2436. caption: str | None = None,
  2437. file_name: str | None = None,
  2438. reply_to: str | None = None,
  2439. metadata: dict[str, Any] | None = None,
  2440. **kwargs,
  2441. ) -> SendResult:
  2442. """Send a local file as a native Chatto attachment.
  2443. ``file_name`` exists in the base-class signature and is accepted for
  2444. compatibility, but Chatto takes the recipient-visible filename from
  2445. the upload session (derived from the local path); failures are logged
  2446. and noticed by ``_send_local_file_as_attachment`` itself.
  2447. BasePlatformAdapter override
  2448. """
  2449. return await self._send_local_file_as_attachment(
  2450. chat_id,
  2451. file_path,
  2452. caption,
  2453. reply_to,
  2454. metadata,
  2455. kind="file",
  2456. )
  2457. async def send_video(
  2458. self,
  2459. chat_id: str,
  2460. video_path: str,
  2461. caption: str | None = None,
  2462. reply_to: str | None = None,
  2463. metadata: dict[str, Any] | None = None,
  2464. **kwargs,
  2465. ) -> SendResult:
  2466. """Send a local video as a native Chatto attachment.
  2467. Chatto transcodes and plays it inline.
  2468. BasePlatformAdapter override
  2469. """
  2470. return await self._send_local_file_as_attachment(
  2471. chat_id,
  2472. video_path,
  2473. caption,
  2474. reply_to,
  2475. metadata,
  2476. kind="video",
  2477. )
  2478. async def send_voice(
  2479. self,
  2480. chat_id: str,
  2481. audio_path: str,
  2482. caption: str | None = None,
  2483. reply_to: str | None = None,
  2484. metadata: dict[str, Any] | None = None,
  2485. **kwargs,
  2486. ) -> SendResult:
  2487. """Send a local audio file as a native Chatto attachment.
  2488. Chatto has no dedicated voice-bubble type, so this is an ordinary audio
  2489. attachment — still far better than the base class's text notice.
  2490. BasePlatformAdapter override
  2491. """
  2492. return await self._send_local_file_as_attachment(
  2493. chat_id,
  2494. audio_path,
  2495. caption,
  2496. reply_to,
  2497. metadata,
  2498. kind="audio",
  2499. )
  2500. async def send_image(
  2501. self,
  2502. chat_id: str,
  2503. image_url: str,
  2504. caption: str | None = None,
  2505. reply_to: str | None = None,
  2506. metadata: dict[str, Any] | None = None,
  2507. ) -> SendResult:
  2508. """Send an image to a Chatto room.
  2509. Materialises the URL (size-capped download, like every inbound
  2510. attachment) and uploads it as a native attachment. Falls back to the
  2511. plain URL as a link — Chatto renders link previews — when the URL
  2512. cannot be materialised or the post after a successful upload fails.
  2513. An upload failure needs no fallback on top: the text notice of
  2514. ``_send_local_file_as_attachment`` has already gone out.
  2515. BasePlatformAdapter override
  2516. """
  2517. link_text = f"{caption}\n{image_url}" if caption else image_url
  2518. path, is_temp = await self._materialise_image(image_url)
  2519. if path is not None:
  2520. try:
  2521. result = await self._send_local_file_as_attachment(
  2522. chat_id, path, caption, reply_to, metadata, kind="image"
  2523. )
  2524. finally:
  2525. if is_temp:
  2526. try:
  2527. os.unlink(path)
  2528. except OSError:
  2529. pass
  2530. if result.success:
  2531. return result
  2532. return await self.send(chat_id, link_text, reply_to=reply_to, metadata=metadata)
  2533. async def _materialise_image(self, image_url: str) -> tuple[str | None, bool]:
  2534. """Resolve one ``send_multiple_images`` entry to a local file path.
  2535. Accepts ``http(s)://`` URLs (downloaded to a temp file), ``file://``
  2536. URIs and bare paths. Returns ``(path, is_temp)`` — the caller unlinks
  2537. when ``is_temp``. ``(None, False)`` means the entry is unusable.
  2538. """
  2539. if image_url.startswith(("http://", "https://")):
  2540. ext = os.path.splitext(urlsplit(image_url).path)[1] or ".png"
  2541. tmp_fd, tmp_path = tempfile.mkstemp(suffix=ext, prefix="chatto_img_")
  2542. os.close(tmp_fd)
  2543. try:
  2544. data = await self._download_attachment_bytes(image_url)
  2545. await asyncio.to_thread(_write_file_bytes, tmp_path, data)
  2546. except Exception as e:
  2547. logger.warning("Chatto: image download failed for %s: %s", image_url, e)
  2548. try:
  2549. os.unlink(tmp_path)
  2550. except OSError:
  2551. pass
  2552. return None, False
  2553. return tmp_path, True
  2554. local = image_url
  2555. if local.startswith("file://"):
  2556. local = unquote(urlsplit(local).path)
  2557. return self.validate_media_delivery_path(local), False
  2558. async def send_multiple_images(
  2559. self,
  2560. chat_id: str,
  2561. images: list[tuple[str, str]],
  2562. metadata: dict[str, Any] | None = None,
  2563. human_delay: float = 0.0,
  2564. ) -> None:
  2565. """Send a batch of images as ONE message with several attachments.
  2566. The base implementation posts each image separately; a Chatto message
  2567. carries a list of attachment assets, so a batch belongs in a single
  2568. message (and a single notification).
  2569. ``human_delay`` is ignored deliberately — there is only one outbound
  2570. call to pace. Entries that can't be fetched are dropped with a warning;
  2571. if nothing survives, we fall back to the base class so the user still
  2572. gets the links.
  2573. BasePlatformAdapter override
  2574. """
  2575. if len(images or []) < 2:
  2576. await super().send_multiple_images(
  2577. chat_id,
  2578. images,
  2579. metadata=metadata,
  2580. human_delay=human_delay,
  2581. )
  2582. return
  2583. asset_ids: list[str] = []
  2584. captions: list[str] = []
  2585. for image_url, alt_text in images:
  2586. path, is_temp = await self._materialise_image(image_url)
  2587. if not path:
  2588. logger.warning("Chatto: skipping unusable image %s", image_url)
  2589. continue
  2590. try:
  2591. asset_id = await self._upload_asset(str(chat_id), path)
  2592. finally:
  2593. if is_temp:
  2594. try:
  2595. os.unlink(path)
  2596. except OSError:
  2597. pass
  2598. if not asset_id:
  2599. logger.warning("Chatto: upload failed for image %s", image_url)
  2600. continue
  2601. asset_ids.append(asset_id)
  2602. if alt_text:
  2603. captions.append(alt_text)
  2604. if not asset_ids:
  2605. logger.warning(
  2606. "Chatto: no image survived upload, falling back to per-image delivery",
  2607. )
  2608. await super().send_multiple_images(
  2609. chat_id,
  2610. images,
  2611. metadata=metadata,
  2612. human_delay=human_delay,
  2613. )
  2614. return
  2615. if len(asset_ids) < len(images):
  2616. logger.warning(
  2617. "Chatto: sending %d of %d images — the rest could not be uploaded",
  2618. len(asset_ids),
  2619. len(images),
  2620. )
  2621. await self._post_attachment_message(
  2622. chat_id,
  2623. asset_ids,
  2624. "\n".join(captions) or None,
  2625. None,
  2626. metadata,
  2627. )
  2628. # ---------------------------------------------------------------------------
  2629. # Cron / out-of-process delivery
  2630. # ---------------------------------------------------------------------------
  2631. async def hermes_standalone_sender_fn(
  2632. pconfig: PlatformConfig,
  2633. chat_id: str,
  2634. message: str,
  2635. *,
  2636. thread_id=None,
  2637. media_files=None,
  2638. force_document=False,
  2639. ) -> SendResult:
  2640. """Deliver a message to Chatto without a running gateway adapter. Do not modify signature.
  2641. Used by cron / scheduled routines that run out-of-process. Creates a
  2642. short-lived chattolib client, posts, and closes. Long messages are
  2643. normalised and split like ``send()`` does, so cron output cannot die on
  2644. the server's per-message limit.
  2645. """
  2646. chatto_config: ChattoConfiguration = ChattoConfiguration(pconfig=pconfig)
  2647. # Create a temporary client for standalone sending — needs a base URL and token.
  2648. has_token = bool(chatto_config.token.value)
  2649. if not chatto_config.base_url.value or not has_token:
  2650. return SendResult(success=False, error="Chatto: base URL or token missing")
  2651. client = ChattoClient(
  2652. base_url=chatto_config.base_url.value, token=chatto_config.token.value
  2653. )
  2654. try:
  2655. kwargs: dict[str, Any] = {}
  2656. if chatto_config.auto_thread.value and thread_id:
  2657. kwargs["thread_root_event_id"] = thread_id
  2658. if media_files and media_files.get("attachment_asset_ids"):
  2659. kwargs["attachment_asset_ids"] = list(media_files["attachment_asset_ids"])
  2660. formatted = _normalise_outbound_text(message)
  2661. chunks = BasePlatformAdapter.truncate_message(
  2662. formatted, ChattoConstants.SPLIT_THRESHOLD
  2663. )
  2664. message_ids: list[str] = []
  2665. try:
  2666. for chunk in chunks:
  2667. posted = await client.post_message(chat_id, chunk, **kwargs)
  2668. message_ids.append(posted.id)
  2669. except Exception as exc:
  2670. if message_ids:
  2671. logger.warning(
  2672. "Chatto standalone: sent %d/%d chunk(s) to %s before failing: %s",
  2673. len(message_ids),
  2674. len(chunks),
  2675. chat_id,
  2676. exc,
  2677. )
  2678. # Same classification as the adapter's send paths: only
  2679. # server-side failures read as retryable.
  2680. return SendResult(
  2681. success=False,
  2682. error=str(exc),
  2683. retryable=isinstance(exc, ChattoError),
  2684. )
  2685. return SendResult(success=True, message_id=message_ids[0])
  2686. finally:
  2687. try:
  2688. await client.close()
  2689. except Exception as exc:
  2690. logger.warning(
  2691. "Chatto standalone: error closing short-lived client (perhaps already closed): %s",
  2692. exc,
  2693. )
  2694. # How similar an ``extra`` key must be to a declared config_key before the
  2695. # unknown-key check reads it as a likely typo or renamed field.
  2696. _NEAR_MISS_RATIO = 0.75
  2697. def _warn_on_unknown_extra_keys(config: PlatformConfig) -> None:
  2698. """Warn about ``extra`` keys that look like misspelled config fields.
  2699. A typo'd key silently resolves to the default otherwise — the warning is
  2700. the only thing telling the user their ``require_mention_channles`` never
  2701. reached us. Only near-matches against our declared config_keys are
  2702. flagged (difflib similarity): keys that resemble nothing of ours are
  2703. either placed there by the Hermes gateway itself (the shared-key loop in
  2704. load_gateway_config() bridges reply_in_thread & co. into every platform's
  2705. extra, hardcoded inline upstream — nothing to import) or are deliberate
  2706. pass-throughs, and flagging those would just train users to ignore us.
  2707. Internal markers (``_enabled_explicit``) are skipped outright.
  2708. """
  2709. extra: dict[str, Any] = getattr(config, "extra", None) or {}
  2710. known = [field.config_key for field in ChattoConfiguration.fields()]
  2711. for key in sorted(extra):
  2712. if key.startswith("_") or key in known:
  2713. continue
  2714. best = max(known, key=lambda k: SequenceMatcher(None, key, k).ratio())
  2715. if SequenceMatcher(None, key, best).ratio() >= _NEAR_MISS_RATIO:
  2716. logger.warning(
  2717. "Chatto: 'extra' key '%s' matches no config field — did you mean '%s'?",
  2718. key,
  2719. best,
  2720. )
  2721. def hermes_validate_config(config: PlatformConfig) -> bool:
  2722. """Check whether Chatto Plugin is configured.
  2723. Function name should be the same as register argument name with "hermes_" prefix, so we
  2724. know that it is needed for plugin register(). Do not change signature.
  2725. Takes ``config``. Compare to hermes_is_connected().
  2726. """
  2727. chatto_config = ChattoConfiguration(pconfig=config)
  2728. _warn_on_unknown_extra_keys(config)
  2729. if (
  2730. len(chatto_config.allowed_users.value) > 0
  2731. and chatto_config.allow_all_users.value
  2732. ):
  2733. logger.error(
  2734. "Chatto: Conflicting configuration. Either use 'allowed_users' or 'allow_all_users' but not both."
  2735. )
  2736. return False
  2737. require_rooms = set(chatto_config.require_mention_rooms.value)
  2738. optional_rooms = set(chatto_config.optional_mention_rooms.value)
  2739. overlap = sorted(require_rooms & optional_rooms)
  2740. if overlap:
  2741. logger.error(
  2742. "Chatto: Conflicting configuration. Room(s) %s are on both "
  2743. "'require_mention_rooms' and 'optional_mention_rooms' — "
  2744. "each room must appear in at most one of them.",
  2745. ", ".join(overlap),
  2746. )
  2747. return False
  2748. # base_url always resolves (it defaults to ChattoHQ), so the only real
  2749. # question is whether a token came in.
  2750. if chatto_config.token.value is not None:
  2751. return True
  2752. logger.error("Chatto: CHATTO_TOKEN must be set.")
  2753. return False
  2754. def hermes_check_fn() -> bool:
  2755. """Report whether the vendored chattolib dependency can be imported."""
  2756. try:
  2757. from chattolib import (
  2758. client, # noqa: F401 — vendored dependency probe
  2759. )
  2760. return True
  2761. except ImportError:
  2762. return False
  2763. # ---------------------------------------------------------------------------
  2764. # is_connected probe
  2765. # ---------------------------------------------------------------------------
  2766. def hermes_is_connected(config: PlatformConfig) -> bool:
  2767. """Report whether the Chatto platform is configured and enabled.
  2768. The name is fixed by the register() contract, but despite what it
  2769. suggests this does not open a connection — it validates configuration
  2770. only (see :func:`hermes_validate_config`); the gateway probes real
  2771. connectivity through ``connect()``.
  2772. """
  2773. return bool(hermes_validate_config(config) and config.enabled)
  2774. def hermes_setup_fn() -> None:
  2775. """Interactive setup wizard for Chatto. Is called by and only works in Hermes CLI context.
  2776. Function name should be the same as register argument name with "hermes_" prefix, so we
  2777. know that it is needed for plugin register().
  2778. """
  2779. from hermes_cli.setup import (
  2780. print_success,
  2781. prompt,
  2782. prompt_yes_no,
  2783. save_env_value,
  2784. )
  2785. url = prompt(
  2786. "Chatto server URL (e.g. https://chat.example.com) or leave blank for default ChattoHQ on chat.chatto.run:"
  2787. )
  2788. if url:
  2789. save_env_value(ChattoConfiguration.base_url.env_name, url)
  2790. token = prompt("Chatto bot token (cht_BK_...):", password=True)
  2791. if token:
  2792. save_env_value(ChattoConfiguration.token.env_name, token)
  2793. home = prompt("Home room ID for notifications (or empty):")
  2794. if home:
  2795. save_env_value(ChattoConfiguration.home_channel.env_name, home)
  2796. allow_all = prompt_yes_no("Allow all users to talk? (true/false):")
  2797. if allow_all:
  2798. save_env_value(ChattoConfiguration.allow_all_users.env_name, str(allow_all))
  2799. print_success("\n✓ Chatto configured. Restart the gateway to activate.")
  2800. def hermes_env_enablement_fn() -> dict:
  2801. """Seed PlatformConfig.extra from the CHATTO_* environment variables.
  2802. Returns a dict compatible with the PlatformConfig merge hook, holding
  2803. exactly the values that are actually set via env — never fabricated
  2804. defaults. Called by the platform registry during load_gateway_config(),
  2805. which commits this dict onto the platform's ``extra`` verbatim; seeding a
  2806. value that was not configured would overwrite what the user set in
  2807. config.yaml (e.g. a YAML base_url clobbered by the ChattoHQ default).
  2808. "Nothing set" therefore stays the adapter's business: every ConfigField
  2809. falls back to its declared default (base_url to ChattoHQ) on its own.
  2810. Seed keys are the config.yaml "extra" keys (config_key), NOT the env var
  2811. names — ChattoConfiguration reads extra[config_key]. The special
  2812. 'home_channel' key is extracted by the gateway and becomes a proper
  2813. HomeChannel dataclass on the PlatformConfig; every other key is merged
  2814. into PlatformConfig.extra.
  2815. Function name should be the same as register argument name with "hermes_" prefix, so we
  2816. know that it is needed for plugin register().
  2817. """
  2818. seed: dict[str, Any] = {}
  2819. for config_field in ChattoConfiguration.fields():
  2820. env_value = os.getenv(config_field.env_name)
  2821. if env_value:
  2822. seed[config_field.config_key] = env_value.strip()
  2823. logger.debug("seed: %s", {k: v for k, v in seed.items() if k != "token"})
  2824. return seed
  2825. # ---------------------------------------------------------------------------
  2826. # Plugin registration entry point
  2827. # ---------------------------------------------------------------------------
  2828. # What each capability is called in the startup banner, keyed by the method
  2829. # that implements it. Derived from real overrides rather than hard-coded, so
  2830. # dropping a method drops its claim from the log instead of leaving a lie.
  2831. _CAPABILITY_LABELS = {
  2832. "send": "text",
  2833. "send_image_file": "images",
  2834. "send_multiple_images": "image batches (bundled into one message)",
  2835. "send_video": "video",
  2836. "send_voice": "voice messages",
  2837. "send_document": "documents",
  2838. "add_reaction": "reactions",
  2839. "edit_message": "message editing",
  2840. "delete_message": "message deletion",
  2841. "send_typing": "typing indicators",
  2842. "create_handoff_thread": "threads",
  2843. "start_dm": "direct messages",
  2844. "create_room": "room creation",
  2845. }
  2846. def _capabilities() -> list[str]:
  2847. """Name the things this adapter genuinely implements itself.
  2848. A capability counts only when ChattoAdapter overrides the base method —
  2849. inheriting BasePlatformAdapter's fallback means the feature is not
  2850. natively supported, and announcing it would mislead.
  2851. """
  2852. found = [
  2853. label
  2854. for name, label in _CAPABILITY_LABELS.items()
  2855. if getattr(ChattoAdapter, name, None)
  2856. is not getattr(BasePlatformAdapter, name, None)
  2857. ]
  2858. if ChattoAdapter.supports_code_blocks:
  2859. found.append("code blocks")
  2860. if ChattoAdapter.supports_status_text:
  2861. found.append("custom status text")
  2862. found.append("presence (refreshed while connected)")
  2863. return found
  2864. def register(ctx) -> None:
  2865. """Plugin entry point — called by the Hermes plugin system."""
  2866. logger.info("Registering Chatto platform plugin on Hermes Agent")
  2867. for capability in _capabilities():
  2868. logger.info("Chatto capability: %s", capability)
  2869. ctx.register_platform(
  2870. name=ChattoConstants.PLATFORM_NAME, # this will be the config.yaml key.
  2871. label=ChattoConstants.PLATFORM_LABEL,
  2872. adapter_factory=hermes_adapter_factory,
  2873. check_fn=hermes_check_fn,
  2874. validate_config=hermes_validate_config,
  2875. is_connected=hermes_is_connected,
  2876. install_hint=ChattoConstants.INSTALL_HINT,
  2877. env_enablement_fn=hermes_env_enablement_fn,
  2878. setup_fn=hermes_setup_fn,
  2879. cron_deliver_env_var=ChattoConfiguration.home_channel.env_name,
  2880. standalone_sender_fn=hermes_standalone_sender_fn,
  2881. allowed_users_env=ChattoConfiguration.allowed_users.env_name,
  2882. allow_all_env=ChattoConfiguration.allow_all_users.env_name,
  2883. max_message_length=ChattoConstants.MAX_MESSAGE_LENGTH,
  2884. emoji="😺",
  2885. allow_update_command=True,
  2886. pii_safe=False,
  2887. platform_hint=(
  2888. "Using the 'Hermes Chatto Platform Plugin' you connect to a Chatto server. "
  2889. "Authorized admins and users will contact you and call you his Hermes Agent. They are "
  2890. "natural persons and thus responsible for what they do in terms of rights. You compute "
  2891. "on their behalf. They _may_ address you by @-mentioning your name. If configured, "
  2892. "you also react without a @-mention. Direct messages reach you without a mention."
  2893. "Keep responses conversational. Markdown is supported. "
  2894. "Include MEDIA:/absolute/path/to/file in your response to refer to our local files. Images "
  2895. "(.png, .jpg, .gif, .webp) arrive as inline pictures, videos (.mp4, .mov, .webm) as "
  2896. "video attachments, audio as a voice bubble, anything else as a downloadable document. "
  2897. "Do NOT use markdown image syntax for local files. Local files always go through MEDIA:. "
  2898. "Several images in one response are bundled into a single message. "
  2899. ),
  2900. )