adapter.py 97 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437
  1. """
  2. Chatto Platform Adapter for Hermes Agent.
  3. A plugin-based gateway adapter that connects to a Chatto server
  4. (self-hosted team chat) and relays messages to/from the Hermes agent.
  5. The adapter uses the chattolib library for all Chatto API interactions,
  6. including both outbound messaging and realtime WebSocket connections.
  7. Configuration in config.yaml::
  8. gateway:
  9. platforms:
  10. chatto:
  11. enabled: true
  12. extra:
  13. url: https://chat.lacy.casa
  14. channels: # room IDs to watch (empty = all joined)
  15. - REljMv5Pgolo6Y9
  16. home_channel: REljMv5Pgolo6Y9
  17. require_mention: true # only respond to @mentions in rooms
  18. allowed_users: [] # empty = allow all
  19. allow_all_users: true
  20. Or via environment variables (overrides config.yaml):
  21. CHATTO_URL, CHATTO_LOGIN, CHATTO_PASSWORD (secrets in ~/.hermes/.env),
  22. CHATTO_CHANNELS, CHATTO_HOME_CHANNEL,
  23. CHATTO_REQUIRE_MENTION, CHATTO_ALLOWED_USERS, CHATTO_ALLOW_ALL_USERS
  24. """
  25. from __future__ import annotations
  26. import asyncio
  27. import hashlib
  28. import logging
  29. import mimetypes
  30. import os
  31. import threading
  32. from collections import OrderedDict
  33. from datetime import datetime, timezone
  34. from typing import Any, Dict, List, Optional, Tuple, cast
  35. from urllib.parse import urlsplit, urlunsplit
  36. from plugins.plugin_utils import lazy_singleton
  37. import utils
  38. logger = logging.getLogger(__name__)
  39. from gateway.platforms.base import (
  40. BasePlatformAdapter,
  41. SendResult,
  42. MessageEvent,
  43. MessageType,
  44. ProcessingOutcome,
  45. )
  46. from gateway.config import Platform, PlatformConfig
  47. # Chattolib imports (vendored)
  48. # Using vendored chattolib from vendor/chattolib/
  49. # See vendor_chattolib.sh for how to update the vendored copy
  50. try:
  51. # Try vendored chattolib first
  52. from vendor.chattolib.client import (
  53. ChattoClient,
  54. ChattoError,
  55. ChattoAuthError,
  56. )
  57. from vendor.chattolib.exceptions import (
  58. ChattoConnectError,
  59. )
  60. from vendor.chattolib.realtime import (
  61. ChattoRealtimeError,
  62. ChattoRealtimeCloseError,
  63. RealtimeConnection,
  64. RealtimeEvent,
  65. ServerHello,
  66. stream_events,
  67. )
  68. from vendor.chattolib.types import (
  69. RoomKind,
  70. PresenceStatus,
  71. RoomWithViewerState,
  72. User,
  73. Message,
  74. )
  75. from vendor.chattolib._pb.chatto.api.v1 import (
  76. rooms_pb2 as room_service_pb2,
  77. threads_pb2 as thread_service_pb2
  78. )
  79. from vendor.chattolib._transport import pb_to_dict
  80. except ImportError as e:
  81. logger.error("Chatto: failed to import vendored chattolib: %s", e)
  82. from .platform_config import ChattoConfig
  83. class AsyncSingletonSlot:
  84. """Thread-safe async singleton helper (extends SingletonSlot pattern for async factories)."""
  85. def __init__(self):
  86. self._instance = None
  87. self._lock = threading.Lock()
  88. self._future = None
  89. async def get(self, factory):
  90. """Get or create instance using async factory. Thread-safe with double-checked locking."""
  91. if self._instance is not None:
  92. return self._instance
  93. with self._lock:
  94. if self._instance is None:
  95. if self._future is None:
  96. self._future = asyncio.ensure_future(factory())
  97. return await self._future
  98. return self._instance
  99. def reset(self):
  100. """Reset the singleton instance."""
  101. with self._lock:
  102. self._instance = None
  103. self._future = None
  104. # --------------------------------------------------------------------------- #
  105. # Constants
  106. # --------------------------------------------------------------------------- #
  107. _MAX_MESSAGE_LENGTH = 10000
  108. _SEEN_CAP = 500
  109. # WebSocket / realtime protocol
  110. _WS_PATH = "/api/realtime"
  111. _WS_AUTH_TIMEOUT = 20.0
  112. _WS_MAX_MESSAGE_BYTES = 4_000_000
  113. _WS_PING_INTERVAL = 30.0
  114. _WS_RECONNECT_INITIAL_BACKOFF = 1.0
  115. _WS_RECONNECT_MAX_BACKOFF = 30.0
  116. # HTTP timeout used for outbound URL fetches
  117. _HTTP_TIMEOUT = 30
  118. # Emoji shortcode mapping (Chatto uses shortcode names, not unicode emoji)
  119. # Emoji shortcode mapping (Chatto uses shortcode names, not unicode emoji)
  120. _EMOJI_TO_SHORTCODE: Dict[str, str] = {
  121. "👍": "thumbsup",
  122. "👎": "thumbsdown",
  123. "❤️": "heart",
  124. "❤": "heart",
  125. "✅": "white_check_mark",
  126. "❌": "x",
  127. "👀": "eyes",
  128. "🎉": "tada",
  129. "😂": "joy",
  130. "🚀": "rocket",
  131. "🔥": "fire",
  132. "💯": "100",
  133. "🤔": "thinking",
  134. "👏": "clap",
  135. "🙏": "pray",
  136. "😅": "sweat_smile",
  137. "😴": "sleeping",
  138. "⏳": "hourglass",
  139. }
  140. # Chunk size for asset uploads (256 KB)
  141. _UPLOAD_CHUNK_SIZE = 256 * 1024
  142. def _decode_event_envelope(data: bytes) -> dict:
  143. """Decode a transient event envelope bytes into a dict.
  144. Tries protobuf parsing via the vendored realtime pb, falls back to JSON.
  145. """
  146. try:
  147. from vendor.chattolib._pb.chatto.realtime.v1 import realtime_pb2
  148. frame = cast(Any, realtime_pb2).RealtimeServerFrame()
  149. frame.ParseFromString(data)
  150. if frame.WhichOneof("frame") == "event":
  151. ev = frame.event
  152. try:
  153. return pb_to_dict(ev)
  154. except Exception:
  155. # Fall through to manual conversion
  156. out = {}
  157. # best-effort: copy simple fields
  158. if hasattr(ev, "room_id"):
  159. out["roomId"] = getattr(ev, "room_id")
  160. return out
  161. except Exception:
  162. pass
  163. try:
  164. import json
  165. return json.loads(data.decode("utf-8")) if data else {}
  166. except Exception:
  167. return {}
  168. # --------------------------------------------------------------------------- #
  169. # Adapter
  170. # --------------------------------------------------------------------------- #
  171. class ChattoAdapter(BasePlatformAdapter):
  172. """Chatto platform adapter — receives messages via WebSocket realtime,
  173. sends via ConnectRPC REST."""
  174. MAX_MESSAGE_LENGTH = 10000
  175. _SPLIT_THRESHOLD = 9900
  176. splits_long_messages = True
  177. _chatto_client: ChattoClient
  178. def __init__(self, config: PlatformConfig, **kwargs):
  179. """Signature needs to be compatible with BasePlatformAdapter.__init__ """
  180. # "extra" has been pre-processed by Hermes-Framework to be a dict of extra config values
  181. chatto_config: ChattoConfig = ChattoConfig(config.extra)
  182. super().__init__(config=config, platform=Platform(chatto_config._PLATFORM))
  183. # --- Configuration from our configuration data class with some logic ---
  184. self.chatto_config = chatto_config
  185. # ------ State -------
  186. # SDK runtime handle (injected by Hermes); annotate for Pylance
  187. self.sdk: Any = getattr(self, "sdk", None)
  188. # --- Runtime state ---
  189. self._user_id: str = ""
  190. self._user_display: str = ""
  191. self._room_names: Dict[str, str] = {}
  192. self._room_kinds: Dict[str, str] = {}
  193. self._our_thread_roots: set = set() # thread root event IDs we created
  194. self._our_message_ids: set = set() # message IDs we sent (for thread root detection)
  195. self._seen: Dict[str, OrderedDict] = {} # room_id -> OrderedDict(event_id -> None)
  196. self._resume_cursor: Optional[str] = None
  197. self._watch_room_ids: List[str] = []
  198. self._ws_task: Optional[asyncio.Task] = None
  199. self._ws_ready: Optional[asyncio.Event] = None
  200. self._ws_active = False
  201. self._ws_ref = None # reference to open websocket for dynamic resubscribe
  202. # Persistent typing indicator loops per room
  203. self._typing_tasks: Dict[str, asyncio.Task] = {}
  204. # Liveness probe (REST health check)
  205. self._liveness_interval_seconds = 60.0
  206. self._liveness_failure_threshold = 3
  207. self._liveness_task: Optional[asyncio.Task] = None
  208. # Member directory cache: user_id -> user info dict
  209. self._user_cache: Dict[str, dict] = {}
  210. # ------------------------------------------------------------------ #
  211. # Auth
  212. # ------------------------------------------------------------------ #
  213. async def _get_chatto_client(self: ChattoAdapter) -> Optional[ChattoClient]:
  214. """Get or create a ChattoClient instance using thread-safe async singleton."""
  215. # Fast path: return cached client if available
  216. if self._chatto_client is not None:
  217. return self._chatto_client
  218. if not self.chatto_config._base_url or not self.chatto_config._login or not self.chatto_config._password:
  219. logger.error("Chatto: missing configuration (URL, login, or password)")
  220. return None
  221. try:
  222. client = await self._chatto_client.get(
  223. lambda: ChattoClient.login(
  224. self.chatto_config._login,
  225. self.chatto_config._password,
  226. base_url=self.chatto_config._base_url,
  227. )
  228. )
  229. self._chatto_client = client # Cache for direct access
  230. logger.info("Chatto: logged in as %s via chattolib", self._login)
  231. return client
  232. except ChattoAuthError as e:
  233. logger.error("Chatto: authentication failed: %s", e)
  234. return None
  235. except Exception as e:
  236. logger.error("Chatto: failed to create client: %s", e)
  237. return None
  238. async def _require_client(self) -> Any:
  239. """Return a ChattoClient or raise RuntimeError if unavailable.
  240. Use this helper when the caller expects a client to exist and
  241. wants a single canonical failure path. Methods that prefer a
  242. soft-fail can catch RuntimeError and return gracefully.
  243. """
  244. client = await self._get_chatto_client()
  245. if client is None:
  246. raise RuntimeError("Chatto client unavailable")
  247. return cast(Any, client)
  248. async def _client(self) -> Optional[ChattoClient]:
  249. """Helper to get the chattolib client. Thread-safe via AsyncSingletonSlot."""
  250. return await self._get_chatto_client()
  251. async def _ensure_token(self) -> bool:
  252. """Login via chattolib."""
  253. if self.chatto_config._token:
  254. return True
  255. client = await self._get_chatto_client()
  256. if client is not None:
  257. self._token = client.token
  258. return True
  259. return False
  260. async def _relogin(self) -> bool:
  261. """Force re-login (token expired)."""
  262. self._token = None
  263. return await self._ensure_token()
  264. # ------------------------------------------------------------------ #
  265. # Connection
  266. # ------------------------------------------------------------------ #
  267. async def _open_client(
  268. *,
  269. base_url: str,
  270. login: str,
  271. password: str,
  272. token: Optional[str] = None,
  273. ) -> Any:
  274. """Return a connected ``ChattoClient`` using token or login/password."""
  275. if token:
  276. return ChattoClient(token=token, base_url=base_url)
  277. return await ChattoClient.login(login, password, base_url=base_url)
  278. async def connect(self, *, is_reconnect: bool = False) -> bool:
  279. """Log into Chatto and start the realtime event stream. Codemancer"""
  280. if not self.chatto_config._base_url:
  281. logger.error("Chatto: CHATTO_BASE_URL not configured")
  282. return False
  283. if not self.chatto_config._login and not self.chatto_config._password:
  284. logger.error(
  285. "Chatto: CHATTO_LOGIN/PASSWORD is not set"
  286. )
  287. return False
  288. self._closing = False
  289. try:
  290. self._client = await self._open_client(
  291. base_url=self.chatto_config._base_url,
  292. login=self.chatto_config._login,
  293. password=self.chatto_config._password,
  294. token=self.chatto_config._token,
  295. )
  296. except Exception as exc:
  297. logger.error("Chatto: failed to construct client: %s", exc)
  298. return False
  299. # Best-effort presence broadcast (non-fatal on failure).
  300. try:
  301. await self._client.update_presence(_PRESENCE_ONLINE)
  302. except Exception:
  303. logger.debug("Chatto: update_presence not supported / failed (non-fatal)")
  304. self._stream_task = asyncio.create_task(
  305. self._run_event_stream(), name="chatto-event-stream",
  306. )
  307. self._mark_connected()
  308. return True
  309. async def connect_too_long(self, *, is_reconnect: bool = False) -> bool:
  310. """Login, discover rooms, start WebSocket realtime connection."""
  311. if not await self._ensure_token():
  312. return False
  313. try:
  314. client = await self._require_client()
  315. except RuntimeError:
  316. self._set_fatal_error("connect_failed", "Chatto client not available", retryable=True)
  317. return False
  318. # Get our own user info
  319. try:
  320. me = await client.me()
  321. self._user_id = str(me.id)
  322. self._user_login = str(me.login)
  323. self._user_display = str(me.display_name or "")
  324. logger.info("Chatto: got user info: %s", self._user_login)
  325. except Exception as e:
  326. logger.error("Chatto: failed to get user info: %s", e)
  327. return False
  328. # Discover rooms
  329. try:
  330. rooms_list = await client.list_rooms()
  331. rooms = []
  332. for room_with_state in rooms_list:
  333. # Defensive checks: vendored types may be partially populated
  334. if not room_with_state:
  335. continue
  336. room_obj = getattr(room_with_state, "room", None)
  337. if not room_obj:
  338. continue
  339. entry = {
  340. "room": {
  341. "id": str(getattr(room_obj, "id", "")),
  342. "name": str(getattr(room_obj, "name", "")),
  343. "kind": str(getattr(getattr(room_obj, "kind", None), "value", "")) if getattr(room_obj, "kind", None) else "",
  344. },
  345. "viewerState": {
  346. "isMember": getattr(getattr(room_with_state, "viewer_state", None), "is_member", False),
  347. }
  348. }
  349. rooms.append(entry)
  350. logger.info("Chatto: got %d rooms", len(rooms))
  351. except Exception as e:
  352. logger.error("Chatto: failed to list rooms: %s", e)
  353. return False
  354. all_room_ids = []
  355. for entry in rooms:
  356. room = entry.get("room", {})
  357. rid = str(room.get("id", ""))
  358. if not rid:
  359. continue
  360. name = str(room.get("name", rid))
  361. kind = str(room.get("kind", ""))
  362. self._room_names[rid] = name
  363. self._room_kinds[rid] = kind
  364. viewer = entry.get("viewerState", {})
  365. is_member = viewer.get("isMember", False)
  366. # If user-specified channels, only watch those; otherwise watch all joined rooms
  367. if self._channel_ids:
  368. if rid in self.chatto_config._channel_ids and not is_member:
  369. await self._join_room(rid)
  370. all_room_ids.append(rid)
  371. elif is_member:
  372. all_room_ids.append(rid)
  373. if self.chatto_config._channel_ids:
  374. watch = list(self.chatto_config._channel_ids)
  375. else:
  376. watch = all_room_ids
  377. if not watch:
  378. logger.error("Chatto: no rooms to watch (join a room or set CHATTO_CHANNELS)")
  379. return False
  380. # Ensure we're a member of each watched room
  381. for rid in watch:
  382. if self._room_kinds.get(rid) != "ROOM_KIND_DM":
  383. await self._join_room(rid)
  384. # Pick home channel
  385. if not self._home_channel:
  386. self._home_channel = watch[0]
  387. self._watch_room_ids = watch
  388. # Initialize seen for each room — seed from REST to avoid replaying history
  389. for rid in watch:
  390. self._seen[rid] = OrderedDict()
  391. await self._seed_room(rid)
  392. # Start WebSocket realtime connection
  393. if not await self._start_chattolib_realtime():
  394. return False
  395. self._mark_connected()
  396. self._start_liveness_probe()
  397. logger.info(
  398. "Chatto: connected to %s as %s, watching %d room(s) via WebSocket",
  399. self._base_url,
  400. self._user_display or self._user_login,
  401. len(watch),
  402. )
  403. # Broadcast online presence so the bot appears online in the member list
  404. try:
  405. await self.set_presence("online")
  406. except Exception:
  407. logger.debug("Chatto: set_presence(online) failed on connect", exc_info=True)
  408. return True
  409. async def disconnect(self) -> None:
  410. """Stop WebSocket, liveness probe, typing tasks, and clear state."""
  411. # Broadcast away presence before tearing down
  412. try:
  413. await self.set_presence("away")
  414. except Exception:
  415. logger.debug("Chatto: set_presence(away) failed on disconnect", exc_info=True)
  416. self._mark_disconnected()
  417. self._ws_active = False
  418. # Cancel liveness probe
  419. await self._cancel_liveness_task()
  420. # Cancel all typing tasks
  421. for chat_id in list(self._typing_tasks.keys()):
  422. await self.stop_typing(chat_id)
  423. if self._ws_task and not self._ws_task.done():
  424. self._ws_task.cancel()
  425. try:
  426. await self._ws_task
  427. except (asyncio.CancelledError, Exception):
  428. pass
  429. self._ws_task = None
  430. self._token = None
  431. self._chatto_client = None
  432. self._client_slot.reset()
  433. # ------------------------------------------------------------------ #
  434. # Liveness probe
  435. # ------------------------------------------------------------------ #
  436. def _start_liveness_probe(self) -> None:
  437. """Start the periodic REST health probe."""
  438. if (
  439. self._liveness_interval_seconds <= 0
  440. or self._liveness_failure_threshold <= 0
  441. ):
  442. return
  443. if self._liveness_task and not self._liveness_task.done():
  444. return
  445. self._liveness_task = asyncio.create_task(self._liveness_loop())
  446. async def _cancel_liveness_task(self) -> None:
  447. """Cancel the liveness probe task."""
  448. task = self._liveness_task
  449. self._liveness_task = None
  450. if task and not task.done():
  451. task.cancel()
  452. try:
  453. await task
  454. except (asyncio.CancelledError, Exception):
  455. pass
  456. async def _liveness_loop(self) -> None:
  457. """Periodically check if the REST API is alive via ViewerService/GetViewer.
  458. Also refreshes presence status on each successful probe so the bot
  459. stays showing as online — Chatto's presence expires if not refreshed.
  460. On ``threshold`` consecutive failures, set a fatal error with
  461. ``retryable=True`` so the gateway runner rebuilds the adapter.
  462. """
  463. interval = self._liveness_interval_seconds
  464. threshold = self._liveness_failure_threshold
  465. failures = 0
  466. while self._running:
  467. try:
  468. await asyncio.sleep(interval)
  469. except asyncio.CancelledError:
  470. return
  471. if not self._running:
  472. return
  473. try:
  474. try:
  475. client = await self._require_client()
  476. except RuntimeError:
  477. reason = "no_client"
  478. failures += 1
  479. logger.warning(
  480. "Chatto: liveness probe aborted - no client available (%s, %d/%d)",
  481. reason,
  482. failures,
  483. threshold,
  484. )
  485. continue
  486. await client.get_viewer()
  487. failures = 0
  488. # Refresh presence to keep showing as online
  489. try:
  490. await self.set_presence("online")
  491. except Exception:
  492. logger.debug("Chatto: presence refresh failed", exc_info=True)
  493. continue
  494. except asyncio.CancelledError:
  495. return
  496. except Exception as e:
  497. reason = str(e)
  498. failures += 1
  499. logger.warning(
  500. "Chatto: liveness probe failed (%s, %d/%d)",
  501. reason,
  502. failures,
  503. threshold,
  504. )
  505. if failures < threshold:
  506. continue
  507. # Threshold exceeded — force reconnect
  508. logger.error(
  509. "Chatto: liveness probe failed %d times consecutively; forcing reconnect",
  510. failures,
  511. )
  512. self._set_fatal_error(
  513. "chatto_liveness_failed",
  514. f"Chatto REST API liveness check failed: {reason}",
  515. retryable=True,
  516. )
  517. # Cancel the WebSocket to trigger reconnect
  518. if self._ws_task and not self._ws_task.done():
  519. self._ws_task.cancel()
  520. return
  521. async def _join_room(self, room_id: str) -> None:
  522. """Join a room if not already a member."""
  523. try:
  524. try:
  525. client = await self._require_client()
  526. except RuntimeError:
  527. logger.debug("Chatto: _join_room aborted - no client available for %s", room_id)
  528. return
  529. await client.join_room(room_id=room_id)
  530. logger.debug("Chatto: joined room %s (%s)", room_id, self._room_names.get(room_id, room_id))
  531. except ChattoError as e:
  532. if "permission_denied" in str(e).lower() or "403" in str(e):
  533. logger.debug("Chatto: already a member of %s or cannot join", room_id)
  534. else:
  535. logger.debug("Chatto: join room %s failed: %s", room_id, e)
  536. async def _seed_room(self, room_id: str) -> None:
  537. """Seed high-water mark from the newest events so a restart doesn't replay history."""
  538. try:
  539. # These should be available from check_chatto_requirements()
  540. try:
  541. client = await self._require_client()
  542. except RuntimeError:
  543. logger.debug("Chatto: _seed_room aborted - no client available for %s", room_id)
  544. return
  545. resp = await client.services.rooms.get_room_events(
  546. cast(Any, room_service_pb2).GetRoomEventsRequest(room_id=room_id),
  547. headers=client._headers(),
  548. )
  549. data = pb_to_dict(resp)
  550. events = data.get("page", {}).get("events", [])
  551. for ev in events:
  552. ev_id = str(ev.get("id", ""))
  553. if ev_id:
  554. self._mark_seen(room_id, ev_id)
  555. logger.debug("Chatto: seeded room %s with %d events", room_id, len(events))
  556. except Exception as e:
  557. logger.debug("Chatto: get room events failed for %s: %s", room_id, e)
  558. def _mark_seen(self, room_id: str, event_id: str) -> None:
  559. seen = self._seen.setdefault(room_id, OrderedDict())
  560. seen[event_id] = None
  561. while len(seen) > _SEEN_CAP:
  562. seen.popitem(last=False)
  563. def _is_seen(self, room_id: str, event_id: str) -> bool:
  564. return event_id in self._seen.get(room_id, {})
  565. # ------------------------------------------------------------------ #
  566. # WebSocket Realtime Transport
  567. # ------------------------------------------------------------------ #
  568. def _websocket_url(self) -> str:
  569. """Build the WebSocket URL from the base HTTP URL."""
  570. parsed = urlsplit(self.chatto_config._base_url.strip())
  571. scheme = {"http": "ws", "https": "wss"}.get(parsed.scheme, parsed.scheme)
  572. if scheme not in ("ws", "wss") or not parsed.netloc:
  573. raise ValueError(f"Chatto URL must use http(s) or ws(s), got {parsed.scheme}")
  574. path = parsed.path.rstrip("/") + _WS_PATH
  575. return urlunsplit((scheme, parsed.netloc, path, parsed.query, ""))
  576. async def _start_chattolib_realtime(self) -> bool:
  577. """Start realtime connection using chattolib's stream_events."""
  578. self._ws_ready = asyncio.Event()
  579. self._ws_task = asyncio.create_task(self._chattolib_event_loop())
  580. try:
  581. await asyncio.wait_for(self._ws_ready.wait(), timeout=_WS_AUTH_TIMEOUT + 10)
  582. except (asyncio.TimeoutError, TimeoutError):
  583. logger.warning("Chatto: chattolib realtime did not connect in time")
  584. self._ws_active = False
  585. if self._ws_task and not self._ws_task.done():
  586. self._ws_task.cancel()
  587. try:
  588. await self._ws_task
  589. except asyncio.CancelledError:
  590. pass
  591. self._ws_task = None
  592. return False
  593. return True
  594. async def _chattolib_event_loop(self) -> None:
  595. """Event loop using chattolib's stream_events.
  596. This replaces the manual WebSocket loop with chattolib's high-level
  597. stream_events() which provides pre-decoded RealtimeEvent objects.
  598. """
  599. try:
  600. client = await self._require_client()
  601. except RuntimeError:
  602. logger.warning("Chatto: chattolib event loop aborted - no client available")
  603. return
  604. # Ensure ws_ready is available for synchronization with starter
  605. assert self._ws_ready is not None
  606. backoff = _WS_RECONNECT_INITIAL_BACKOFF
  607. try:
  608. while True:
  609. try:
  610. logger.info("Chatto: starting chattolib event stream with %d rooms", len(self._watch_room_ids))
  611. # Start streaming events
  612. async for event in stream_events(client):
  613. # Signal that we're connected and ready
  614. if not self._ws_ready.is_set():
  615. self._ws_active = True
  616. self._ws_ready.set()
  617. backoff = _WS_RECONNECT_INITIAL_BACKOFF
  618. # Handle different event kinds
  619. if event.kind == "projection_event":
  620. # Convert chattolib RealtimeEvent to our format
  621. # event.payload is the RealtimeProjectionEvent protobuf
  622. try:
  623. # Extract the raw bytes for compatibility with existing handler
  624. # For now, we'll use the existing _handle_projection_event
  625. # which expects bytes. We need to convert.
  626. #
  627. # Actually, let's create a new handler that works with
  628. # chattolib's event objects directly.
  629. await self._handle_chattolib_projection_event(event)
  630. except Exception as e:
  631. logger.warning("Chatto: failed to handle projection event: %s", e)
  632. elif event.kind == "caught_up":
  633. # Update resume cursor
  634. if hasattr(event.payload, 'cursor'):
  635. self._resume_cursor = event.payload.cursor
  636. logger.debug("Chatto: caught_up received, cursor=%s", self._resume_cursor or "(none)")
  637. elif event.kind in ("message_posted", "mention_notification",
  638. "new_direct_message_notification", "user_joined_room",
  639. "room_created", "user_left_room", "message_edited",
  640. "message_retracted", "session_terminated"):
  641. # Transient events - convert to envelope format for existing handler
  642. await self._handle_chattolib_transient_event(event)
  643. elif event.kind in ("heartbeat", "pong", "subscribed"):
  644. # Ignore these
  645. logger.debug("Chatto: %s event received", event.kind)
  646. elif event.kind == "error":
  647. logger.warning("Chatto: server error event: %s", event.payload)
  648. elif event.kind == "close":
  649. logger.info("Chatto: server sent close event")
  650. raise ConnectionError("Server closed connection")
  651. else:
  652. logger.debug("Chatto: unknown event kind: %s", event.kind)
  653. except ChattoRealtimeCloseError as e:
  654. logger.warning("Chatto: realtime closed by server: %s (reconnect=%s)", e.message, e.reconnect)
  655. if e.reconnect:
  656. self._ws_active = False
  657. await asyncio.sleep(backoff)
  658. backoff = min(backoff * 2, _WS_RECONNECT_MAX_BACKOFF)
  659. continue
  660. raise
  661. except ChattoRealtimeError as e:
  662. logger.warning("Chatto: realtime error: %s (fatal=%s)", e.message, e.fatal)
  663. if e.fatal:
  664. raise
  665. except (ConnectionError, asyncio.CancelledError):
  666. raise
  667. except Exception as e:
  668. self._ws_active = False
  669. logger.warning("Chatto: event stream error: %s, retrying in %.1fs", e, backoff)
  670. await asyncio.sleep(backoff)
  671. backoff = min(backoff * 2, _WS_RECONNECT_MAX_BACKOFF)
  672. finally:
  673. self._ws_active = False
  674. async def _handle_chattolib_projection_event(self, event: RealtimeEvent) -> None:
  675. """Handle a chattolib RealtimeEvent with kind='projection_event'.
  676. This is a wrapper that converts chattolib's event to the format
  677. expected by _handle_projection_event.
  678. """
  679. try:
  680. # event.payload is a RealtimeProjectionEvent protobuf message
  681. # We need to convert it to the dict format that _handle_projection_event expects
  682. # pb_to_dict should be available from check_chatto_requirements()
  683. pe_dict = pb_to_dict(event.payload)
  684. # Extract operations from the projection event
  685. operations = []
  686. if "operations" in pe_dict:
  687. for op_pb in event.payload.operations:
  688. op_dict = pb_to_dict(op_pb)
  689. operations.append(op_dict)
  690. # Build the event dict in the format expected by _handle_projection_event
  691. event_data = {
  692. "id": pe_dict.get("id", ""),
  693. "created_at": pe_dict.get("createdAt", ""),
  694. "actor_id": pe_dict.get("actorId", ""),
  695. "resume_cursor": pe_dict.get("resumeCursor", ""),
  696. "operations": operations,
  697. }
  698. await self._handle_projection_event_from_dict(event_data)
  699. except Exception as e:
  700. logger.warning("Chatto: failed to convert chattolib projection event: %s", e)
  701. async def _handle_projection_event_from_dict(self, event: dict) -> None:
  702. """Handle a projection event from a dict (used by chattolib wrapper)."""
  703. # Update resume cursor if provided
  704. cursor = event.get("resume_cursor") or event.get("resumeCursor")
  705. if cursor:
  706. self._resume_cursor = cursor
  707. operations = event.get("operations", [])
  708. for op in operations:
  709. op_type = op.get("type", "")
  710. if op_type == "room_timeline_event_upsert":
  711. await self._handle_timeline_event_upsert_from_dict(op)
  712. async def _handle_timeline_event_upsert_from_dict(self, op: dict) -> None:
  713. """Handle a room_timeline_event_upsert operation from dict."""
  714. room_id = op.get("room_id", "") or op.get("roomId", "")
  715. event_data = op.get("event", {})
  716. if not event_data:
  717. return
  718. ev_id = str(event_data.get("id", ""))
  719. if not ev_id:
  720. return
  721. # De-dupe
  722. if self._is_seen(room_id, ev_id):
  723. return
  724. self._mark_seen(room_id, ev_id)
  725. # Only handle messagePosted events
  726. posted = event_data.get("messagePosted", {})
  727. if not posted:
  728. return
  729. msg = posted.get("message", {})
  730. if not msg:
  731. return
  732. await self._dispatch_message(msg, room_id)
  733. async def _handle_chattolib_transient_event(self, event: RealtimeEvent) -> None:
  734. """Handle a chattolib RealtimeEvent with transient event kinds.
  735. This converts chattolib's event to the envelope format expected by
  736. _handle_transient_event.
  737. """
  738. try:
  739. # pb_to_dict should be available from check_chatto_requirements()
  740. # Build envelope dict based on event kind
  741. envelope = {
  742. "type": event.kind,
  743. "actorId": event.actor_id or "",
  744. "id": event.id or "",
  745. "data": pb_to_dict(event.payload) if event.payload else {},
  746. }
  747. # Convert data field names to match expected format
  748. data = envelope["data"]
  749. if event.kind == "message_posted":
  750. data["roomId"] = data.get("roomId", "")
  751. data["messageEventId"] = data.get("eventId", data.get("id", ""))
  752. data["threadRootEventId"] = data.get("threadRootEventId", "")
  753. elif event.kind == "mention_notification":
  754. data["roomId"] = data.get("roomId", "")
  755. data["eventId"] = data.get("eventId", data.get("id", ""))
  756. elif event.kind == "new_direct_message_notification":
  757. data["roomId"] = data.get("roomId", "")
  758. data["eventId"] = data.get("eventId", data.get("id", ""))
  759. elif event.kind in ("user_joined_room", "room_created", "user_left_room"):
  760. data["roomId"] = data.get("roomId", data.get("room_id", ""))
  761. data["actorId"] = data.get("actorId", data.get("actor_id", ""))
  762. elif event.kind == "message_edited":
  763. data["roomId"] = data.get("roomId", "")
  764. data["messageEventId"] = data.get("eventId", data.get("id", ""))
  765. elif event.kind == "message_retracted":
  766. data["roomId"] = data.get("roomId", "")
  767. data["messageEventId"] = data.get("messageEventId", data.get("eventId", ""))
  768. data["reason"] = data.get("reason", "")
  769. elif event.kind == "session_terminated":
  770. data["reason"] = data.get("reason", "")
  771. await self._handle_transient_event_from_dict(envelope)
  772. except Exception as e:
  773. logger.warning("Chatto: failed to handle chattolib transient event: %s", e)
  774. async def _handle_transient_event_from_dict(self, envelope: dict) -> None:
  775. """Handle a transient event from a dict (used by chattolib wrapper)."""
  776. event_type = envelope.get("type", "unknown")
  777. event_data = envelope.get("data", {})
  778. actor_id = envelope.get("actorId", "")
  779. if event_type == "message_posted":
  780. room_id = event_data.get("roomId", "")
  781. event_id = event_data.get("messageEventId", "")
  782. thread_root = event_data.get("threadRootEventId", "")
  783. logger.info("Chatto WS: message_posted in room %s, event %s (thread=%s)", room_id, event_id, thread_root or "none")
  784. if room_id and event_id and not self._is_seen(room_id, event_id):
  785. await self._fetch_and_dispatch_event(room_id, event_id, thread_root)
  786. elif event_type == "mention_notification":
  787. room_id = event_data.get("roomId", "")
  788. event_id = event_data.get("eventId", "")
  789. logger.info("Chatto WS: mention notification in room %s for event %s", room_id, event_id)
  790. if room_id and event_id and not self._is_seen(room_id, event_id):
  791. await self._fetch_and_dispatch_event(room_id, event_id)
  792. elif event_type == "new_direct_message_notification":
  793. room_id = event_data.get("roomId", "")
  794. event_id = event_data.get("eventId", "")
  795. logger.info("Chatto WS: new DM notification in room %s for event %s", room_id, event_id)
  796. if room_id and event_id and not self._is_seen(room_id, event_id):
  797. await self._fetch_and_dispatch_event(room_id, event_id)
  798. elif event_type == "user_joined_room":
  799. room_id = event_data.get("roomId", "")
  800. actor_id = envelope.get("actorId", "")
  801. logger.info("Chatto WS: user_joined_room room=%s actor=%s", room_id, actor_id)
  802. if room_id and room_id not in self._watch_room_ids:
  803. await self._refresh_rooms()
  804. elif event_type == "room_created":
  805. room_id = event_data.get("roomId", "")
  806. logger.info("Chatto WS: room_created room=%s", room_id)
  807. if room_id and room_id not in self._watch_room_ids:
  808. await self._refresh_rooms()
  809. elif event_type == "user_left_room":
  810. room_id = event_data.get("roomId", "")
  811. actor_id = envelope.get("actorId", "")
  812. logger.info("Chatto WS: user_left_room room=%s actor=%s", room_id, actor_id)
  813. if room_id and actor_id == self._user_id and room_id in self._watch_room_ids:
  814. self._watch_room_ids.remove(room_id)
  815. logger.info("Chatto WS: stopped watching room %s (we left)", room_id)
  816. elif event_type == "message_edited":
  817. room_id = event_data.get("roomId", "")
  818. event_id = event_data.get("messageEventId", "")
  819. logger.info("Chatto WS: message_edited in room %s, event %s", room_id, event_id)
  820. elif event_type == "message_retracted":
  821. room_id = event_data.get("roomId", "")
  822. event_id = event_data.get("messageEventId", "")
  823. reason = event_data.get("reason", "")
  824. logger.info("Chatto WS: message_retracted in room %s, event %s (reason=%s)", room_id, event_id, reason or "none")
  825. if room_id and event_id:
  826. self._mark_seen(room_id, event_id)
  827. elif event_type == "session_terminated":
  828. reason = event_data.get("reason", "")
  829. logger.warning("Chatto WS: session terminated by server (reason=%s) — forcing reconnect", reason or "none")
  830. # We can't close _ws_ref here since we're using chattolib
  831. # The reconnect will happen automatically in _chattolib_event_loop
  832. else:
  833. logger.debug("Chatto WS: unknown transient event type: %s", event_type)
  834. # Update resume cursor if provided
  835. cursor = envelope.get("resume_cursor") or envelope.get("resumeCursor")
  836. if cursor:
  837. self._resume_cursor = cursor
  838. operations = envelope.get("operations", [])
  839. for op in operations:
  840. if op.get("type") == "room_timeline_event_upsert":
  841. await self._handle_timeline_event_upsert(op)
  842. # Other operation types (room_upsert, room_member_upsert, etc.) are
  843. # not relevant to message delivery — ignore them.
  844. async def _handle_timeline_event_upsert(self, op: dict) -> None:
  845. """Handle a room_timeline_event_upsert operation."""
  846. room_id = op.get("room_id", "")
  847. event = op.get("event", {})
  848. if not event:
  849. return
  850. ev_id = str(event.get("id", ""))
  851. if not ev_id:
  852. return
  853. # De-dupe: skip events we've already seen
  854. if self._is_seen(room_id, ev_id):
  855. return
  856. self._mark_seen(room_id, ev_id)
  857. # Only handle messagePosted events
  858. posted = event.get("messagePosted")
  859. if not posted:
  860. return
  861. msg = posted.get("message", {})
  862. if not msg:
  863. return
  864. await self._dispatch_message(msg, room_id)
  865. async def _handle_transient_event(self, data: bytes) -> None:
  866. """Handle a transient RealtimeEventEnvelope (message_posted, mentions, DMs).
  867. These are signal-only events — they contain room_id and event_id but NOT
  868. the message body. We fetch the actual message via REST as a fallback.
  869. """
  870. if not data:
  871. return
  872. try:
  873. envelope = _decode_event_envelope(data)
  874. except (ValueError, IndexError) as e:
  875. logger.warning("Chatto WS: failed to decode transient event: %s", e)
  876. return
  877. event_type = envelope.get("type", "unknown")
  878. event_data = envelope.get("data", {})
  879. if event_type == "message_posted":
  880. room_id = event_data.get("roomId", "")
  881. event_id = event_data.get("messageEventId", "")
  882. thread_root = event_data.get("threadRootEventId", "")
  883. logger.info("Chatto WS: message_posted in room %s, event %s (thread=%s)", room_id, event_id, thread_root or "none")
  884. if room_id and event_id and not self._is_seen(room_id, event_id):
  885. await self._fetch_and_dispatch_event(room_id, event_id, thread_root)
  886. elif event_type == "mention_notification":
  887. room_id = event_data.get("roomId", "")
  888. event_id = event_data.get("eventId", "")
  889. logger.info("Chatto WS: mention notification in room %s for event %s", room_id, event_id)
  890. if room_id and event_id and not self._is_seen(room_id, event_id):
  891. await self._fetch_and_dispatch_event(room_id, event_id)
  892. elif event_type == "new_direct_message_notification":
  893. room_id = event_data.get("roomId", "")
  894. event_id = event_data.get("eventId", "")
  895. logger.info("Chatto WS: new DM notification in room %s for event %s", room_id, event_id)
  896. if room_id and event_id and not self._is_seen(room_id, event_id):
  897. await self._fetch_and_dispatch_event(room_id, event_id)
  898. elif event_type == "user_joined_room":
  899. room_id = event_data.get("roomId", "")
  900. actor_id = envelope.get("actorId", "")
  901. logger.info("Chatto WS: user_joined_room room=%s actor=%s", room_id, actor_id)
  902. # If WE joined a room (or someone else joined and we should watch it),
  903. # refresh room list and resubscribe
  904. if room_id and room_id not in self._watch_room_ids:
  905. await self._refresh_rooms()
  906. elif event_type == "room_created":
  907. room_id = event_data.get("roomId", "")
  908. logger.info("Chatto WS: room_created room=%s", room_id)
  909. # A new room was created — check if we should join/watch it
  910. if room_id and room_id not in self._watch_room_ids:
  911. await self._refresh_rooms()
  912. elif event_type == "user_left_room":
  913. room_id = event_data.get("roomId", "")
  914. actor_id = envelope.get("actorId", "")
  915. logger.info("Chatto WS: user_left_room room=%s actor=%s", room_id, actor_id)
  916. # If WE left a room, stop watching it
  917. if room_id and actor_id == self._user_id and room_id in self._watch_room_ids:
  918. self._watch_room_ids.remove(room_id)
  919. logger.info("Chatto WS: stopped watching room %s (we left)", room_id)
  920. elif event_type == "message_edited":
  921. room_id = event_data.get("roomId", "")
  922. event_id = event_data.get("messageEventId", "")
  923. logger.info("Chatto WS: message_edited in room %s, event %s", room_id, event_id)
  924. # Log edit — could re-fetch for context if needed in the future
  925. elif event_type == "message_retracted":
  926. room_id = event_data.get("roomId", "")
  927. event_id = event_data.get("messageEventId", "")
  928. reason = event_data.get("reason", "")
  929. logger.info("Chatto WS: message_retracted in room %s, event %s (reason=%s)", room_id, event_id, reason or "none")
  930. # Mark the message as seen so we don't try to dispatch it later
  931. if room_id and event_id:
  932. self._mark_seen(room_id, event_id)
  933. elif event_type == "session_terminated":
  934. reason = event_data.get("reason", "")
  935. logger.warning("Chatto WS: session terminated by server (reason=%s) — forcing reconnect", reason or "none")
  936. # Close the websocket to trigger reconnect with backoff
  937. if self._ws_ref:
  938. try:
  939. ws = cast(Any, self._ws_ref)
  940. await ws.close()
  941. except Exception:
  942. pass
  943. else:
  944. logger.debug("Chatto WS: unknown transient event type: %s", event_type)
  945. async def _fetch_and_dispatch_event(self, room_id: str, event_id: str, thread_root_event_id: str = "") -> None:
  946. """Fetch a single event by ID via REST and dispatch it.
  947. Used as a fallback when the projection_event for a transient
  948. notification (mention/DM) hasn't arrived yet.
  949. When thread_root_event_id is set, fetches from the thread timeline
  950. instead of the room timeline.
  951. """
  952. self._mark_seen(room_id, event_id)
  953. try:
  954. # Ensure we have a client
  955. try:
  956. client = await self._require_client()
  957. except RuntimeError:
  958. logger.warning("Chatto WS: REST fallback fetch aborted - no client available for event %s", event_id)
  959. return
  960. # Fetch events either from the thread timeline or the room timeline
  961. if thread_root_event_id:
  962. resp = await client.services.threads.get_thread_events(
  963. cast(Any, thread_service_pb2).GetThreadEventsRequest(
  964. room_id=room_id,
  965. thread_root_event_id=thread_root_event_id,
  966. ),
  967. headers=client._headers(),
  968. )
  969. else:
  970. resp = await client.services.rooms.get_room_events(
  971. cast(Any, room_service_pb2).GetRoomEventsRequest(room_id=room_id),
  972. headers=client._headers(),
  973. )
  974. data = pb_to_dict(resp)
  975. events = data.get("page", {}).get("events", [])
  976. for ev in events:
  977. ev_id = str(ev.get("id", ""))
  978. if ev_id == event_id:
  979. posted = ev.get("messagePosted")
  980. if posted:
  981. msg = posted.get("message", {})
  982. if msg:
  983. # Ensure thread info is set on the message so
  984. # _dispatch_message can extract the thread root ID.
  985. if thread_root_event_id and not msg.get("thread"):
  986. msg["thread"] = {"threadRootEventId": thread_root_event_id}
  987. logger.info("Chatto WS: dispatching event %s via REST fallback (thread=%s)", event_id, thread_root_event_id or "none")
  988. await self._dispatch_message(msg, room_id)
  989. return
  990. logger.warning("Chatto WS: event %s not found in room %s events (thread=%s)", event_id, room_id, thread_root_event_id or "none")
  991. except Exception as e:
  992. logger.warning("Chatto WS: REST fallback fetch failed for event %s: %s", event_id, e)
  993. # If the REST fallback didn't find the event, refresh known rooms
  994. # by delegating to a dedicated helper. This avoids duplicate logic
  995. # across transient/projection event handlers.
  996. await self._refresh_rooms()
  997. async def _dispatch_message(self, msg: dict, room_id: str) -> None:
  998. """Build a MessageEvent and hand it to the base class handler.
  999. This method is identical to the polling version — it receives a
  1000. message dict (decoded from protobuf) and dispatches it through the
  1001. standard Hermes message pipeline.
  1002. """
  1003. if not self._message_handler:
  1004. return
  1005. actor_id = str(msg.get("actorId", ""))
  1006. # Skip our own messages
  1007. if actor_id == self._user_id:
  1008. return
  1009. # Best-effort: cache the sender's display name for richer message context
  1010. if actor_id and actor_id not in self._user_cache:
  1011. try:
  1012. await self.get_user(actor_id)
  1013. except Exception:
  1014. logger.debug("Chatto: get_user(%s) failed during dispatch", actor_id, exc_info=True)
  1015. body = str(msg.get("body", ""))
  1016. if not body:
  1017. return
  1018. msg_id = str(msg.get("id", ""))
  1019. chat_type = "dm" if self._room_kinds.get(room_id) == "ROOM_KIND_DM" else "group"
  1020. # Mention detection
  1021. is_dm = chat_type == "dm"
  1022. mentioned = False
  1023. if self._user_login:
  1024. mentioned = f"@{self._user_login}" in body
  1025. if self._user_display:
  1026. mentioned = mentioned or f"@{self._user_display}" in body
  1027. if self.chatto_config._require_mention and not is_dm and not mentioned:
  1028. # Allow free-response rooms (like Discord's free_response_channels)
  1029. if room_id not in self.chatto_config._free_response_channels:
  1030. return
  1031. # For DMs, always respond. For rooms with require_mention, only respond when mentioned.
  1032. # Strip the mention from the text for the agent
  1033. text = body
  1034. if mentioned and not is_dm:
  1035. # Remove mention prefix if present
  1036. if self._user_login and text.startswith(f"@{self._user_login}"):
  1037. text = text[len(f"@{self._user_login}"):].lstrip()
  1038. elif self._user_display and text.startswith(f"@{self._user_display}"):
  1039. text = text[len(f"@{self._user_display}"):].lstrip()
  1040. # Resolve user display name from actorLogin or actorDisplayName
  1041. user_name = str(msg.get("actorLogin", "")) or str(msg.get("actorDisplayName", actor_id))
  1042. thread_id = None
  1043. thread_info = msg.get("thread", {})
  1044. if thread_info and str(thread_info.get("threadRootEventId", "")) != msg_id:
  1045. thread_id = str(thread_info.get("threadRootEventId", ""))
  1046. # Hermes SDK: Propagate thread context if thread_id is set
  1047. try:
  1048. # Hermes injects the SDK into the plugin context as self.sdk
  1049. propagate_context_to_thread = self.sdk.thread_context.propagate_context_to_thread
  1050. propagate_context_to_thread(thread_id)
  1051. except AttributeError:
  1052. logger.debug("Hermes SDK thread_context not available in plugin context")
  1053. except Exception as e:
  1054. logger.warning("Failed to propagate thread context: %s", e, exc_info=True)
  1055. source = self.build_source(
  1056. chat_id=room_id,
  1057. chat_name=self._room_names.get(room_id, room_id),
  1058. chat_type=chat_type,
  1059. user_id=actor_id,
  1060. user_name=user_name,
  1061. thread_id=thread_id,
  1062. )
  1063. created_at_str = str(msg.get("createdAt", ""))
  1064. try:
  1065. timestamp = datetime.fromisoformat(created_at_str.replace("Z", "+00:00")) if created_at_str else datetime.now()
  1066. except (ValueError, TypeError):
  1067. timestamp = datetime.now()
  1068. event = MessageEvent(
  1069. text=text,
  1070. message_type=MessageType.TEXT,
  1071. source=source,
  1072. message_id=msg_id,
  1073. timestamp=timestamp,
  1074. raw_message=msg,
  1075. )
  1076. await self.handle_message(event)
  1077. async def _refresh_rooms(self) -> None:
  1078. """Refresh room list via REST, join and seed any newly discovered rooms."""
  1079. try:
  1080. try:
  1081. client = await self._require_client()
  1082. except RuntimeError:
  1083. logger.warning("Chatto WS: _refresh_rooms aborted - no client available")
  1084. return
  1085. rooms_list = await client.list_rooms()
  1086. new_room_ids: List[str] = []
  1087. for room_with_state in rooms_list:
  1088. if not room_with_state:
  1089. continue
  1090. room_obj = getattr(room_with_state, "room", None)
  1091. if not room_obj:
  1092. continue
  1093. rid = str(getattr(room_obj, "id", ""))
  1094. name = str(getattr(room_obj, "name", ""))
  1095. kind = str(getattr(getattr(room_obj, "kind", None), "value", "")) if getattr(room_obj, "kind", None) else ""
  1096. is_member = getattr(getattr(room_with_state, "viewer_state", None), "is_member", False)
  1097. self._room_names[rid] = name
  1098. self._room_kinds[rid] = kind
  1099. if is_member and rid not in self._watch_room_ids:
  1100. new_room_ids.append(rid)
  1101. if not new_room_ids:
  1102. return
  1103. logger.info("Chatto WS: discovered %d new room(s): %s", len(new_room_ids), new_room_ids)
  1104. for rid in new_room_ids:
  1105. if self._room_kinds.get(rid) != "ROOM_KIND_DM":
  1106. await self._join_room(rid)
  1107. self._seen[rid] = OrderedDict()
  1108. await self._seed_room(rid)
  1109. self._watch_room_ids.append(rid)
  1110. logger.info("Chatto WS: updated watch list with %d room(s)", len(self._watch_room_ids))
  1111. except Exception:
  1112. logger.warning("Chatto WS: _refresh_rooms failed", exc_info=True)
  1113. # ------------------------------------------------------------------ #
  1114. # Read state & notification dismissal (best-effort, Chatto-unique)
  1115. # ------------------------------------------------------------------ #
  1116. # Best-effort: mark all watched rooms as read (room_id may be undefined here)
  1117. for _rid in list(self._watch_room_ids):
  1118. try:
  1119. await self.mark_room_as_read(_rid)
  1120. except Exception:
  1121. logger.debug("Chatto: mark_room_as_read failed for %s", _rid, exc_info=True)
  1122. try:
  1123. await self.dismiss_all_notifications()
  1124. except Exception:
  1125. logger.debug("Chatto: dismiss_all_notifications failed", exc_info=True)
  1126. # ------------------------------------------------------------------ #
  1127. # Sending (REST — unchanged from polling version)
  1128. # ------------------------------------------------------------------ #
  1129. async def send(
  1130. self,
  1131. chat_id: str,
  1132. content: str,
  1133. reply_to: Optional[str] = None,
  1134. metadata: Optional[Dict[str, Any]] = None,
  1135. ) -> SendResult:
  1136. """Send a message to a Chatto room.
  1137. Long messages are split into chunks via ``truncate_message`` and
  1138. each chunk is sent as a separate CreateMessage call. The first
  1139. chunk's message ID is returned as ``message_id``.
  1140. When ``auto_thread`` is enabled and the incoming message was a
  1141. regular room message (not already in a thread), the first chunk is
  1142. sent as a room message and its ID becomes the thread root. Subsequent
  1143. chunks are sent in that thread. This mirrors Discord's auto_thread
  1144. behavior.
  1145. """
  1146. if not content:
  1147. return SendResult(success=False, error="Empty message")
  1148. formatted = self.format_message(content) if hasattr(self, "format_message") else content
  1149. chunks = self.truncate_message(formatted, self.MAX_MESSAGE_LENGTH)
  1150. # Thread support — resolve thread_id once
  1151. # DM rooms don't support threads, so skip threading for DMs
  1152. thread_id = (metadata or {}).get("thread_id")
  1153. # Only use reply_to as thread_id if auto_thread is enabled.
  1154. # When auto_thread=false, responses go directly in the room
  1155. # without threading under the incoming message.
  1156. if reply_to and self.chatto_config._auto_thread:
  1157. # reply_to might be the incoming message ID. If we already have
  1158. # thread_id from metadata, keep it (it's the thread root).
  1159. # Only use reply_to as thread_id if we don't already have one.
  1160. if not thread_id:
  1161. thread_id = reply_to
  1162. # Check if this is a DM room — DMs don't support threads
  1163. room_kind = self._room_kinds.get(str(chat_id), "")
  1164. is_dm = room_kind == "ROOM_KIND_DM" or room_kind == "dm"
  1165. if is_dm:
  1166. thread_id = None
  1167. # Auto-thread: by default, Chatto creates a thread for replies to room
  1168. # messages (not DMs, not already in a thread). This keeps conversations
  1169. # organized in the room. Can be disabled via extra.auto_thread=false.
  1170. auto_thread_enabled = os.getenv("CHATTO_AUTO_THREAD", "").strip().lower()
  1171. if auto_thread_enabled:
  1172. auto_thread_enabled = auto_thread_enabled in ("true", "1", "yes")
  1173. else:
  1174. auto_thread_enabled = True # default: enabled
  1175. use_auto_thread = auto_thread_enabled and not thread_id and not is_dm
  1176. message_ids: List[str] = []
  1177. last_resp: Optional[dict] = None
  1178. last_error: Optional[str] = None
  1179. retryable = False
  1180. try:
  1181. client = await self._require_client()
  1182. except RuntimeError:
  1183. return SendResult(success=False, error="Chatto client not available", retryable=True)
  1184. for i, chunk in enumerate(chunks):
  1185. try:
  1186. msg_obj = await client.post_message(
  1187. room_id=str(chat_id),
  1188. body=chunk,
  1189. thread_root_event_id=str(thread_id) if thread_id else "",
  1190. )
  1191. msg_id = str(msg_obj.id)
  1192. last_resp = {"message": {"id": msg_id}}
  1193. except ChattoError as e:
  1194. last_error = str(e)
  1195. retryable = True
  1196. break
  1197. except Exception as e:
  1198. last_error = str(e)
  1199. retryable = True
  1200. break
  1201. if msg_id:
  1202. self._mark_seen(str(chat_id), msg_id)
  1203. message_ids.append(msg_id)
  1204. self._our_message_ids.add(msg_id)
  1205. # If we sent a message WITHOUT a thread_id, this message could
  1206. # become a thread root if someone replies to it
  1207. if not thread_id:
  1208. self._our_thread_roots.add(msg_id)
  1209. # Auto-thread: first chunk becomes the thread root,
  1210. # subsequent chunks go in the thread
  1211. if use_auto_thread and i == 0 and not thread_id:
  1212. thread_id = msg_id
  1213. if last_error and not message_ids:
  1214. return SendResult(success=False, error=last_error, retryable=retryable)
  1215. first_id = message_ids[0] if message_ids else ""
  1216. # ------------------------------------------------------------------ #
  1217. # Thread following (best-effort, Chatto-unique)
  1218. # ------------------------------------------------------------------ #
  1219. if thread_id and message_ids:
  1220. try:
  1221. await self._follow_thread(str(chat_id), str(thread_id))
  1222. except Exception:
  1223. logger.debug("Chatto: _follow_thread failed for room=%s thread=%s",
  1224. chat_id, thread_id, exc_info=True)
  1225. return SendResult(success=True, message_id=first_id, raw_response=last_resp)
  1226. async def send_typing(self, chat_id: str, metadata=None) -> None:
  1227. """Start a persistent typing indicator for a room.
  1228. Sends a typing ping every 10 seconds (Chatto's indicator likely
  1229. lasts ~8-10s). The background loop runs until ``stop_typing()``
  1230. is called or the task is cancelled.
  1231. """
  1232. if chat_id in self._typing_tasks:
  1233. return # already running
  1234. async def _typing_loop() -> None:
  1235. try:
  1236. while True:
  1237. try:
  1238. try:
  1239. client = await self._require_client()
  1240. except RuntimeError:
  1241. return
  1242. await client.update_typing_indicator(room_id=str(chat_id))
  1243. except asyncio.CancelledError:
  1244. return
  1245. except Exception:
  1246. pass
  1247. await asyncio.sleep(10)
  1248. except asyncio.CancelledError:
  1249. pass
  1250. finally:
  1251. self._typing_tasks.pop(chat_id, None)
  1252. self._typing_tasks[chat_id] = asyncio.create_task(_typing_loop())
  1253. async def stop_typing(self, chat_id: str) -> None:
  1254. """Stop the persistent typing indicator for a room."""
  1255. task = self._typing_tasks.pop(chat_id, None)
  1256. if task:
  1257. task.cancel()
  1258. try:
  1259. await task
  1260. except (asyncio.CancelledError, Exception):
  1261. pass
  1262. async def get_chat_info(self, chat_id: str) -> Dict[str, Any]:
  1263. """Get information about a chat/room."""
  1264. name = self._room_names.get(chat_id, chat_id)
  1265. kind = self._room_kinds.get(chat_id, "")
  1266. chat_type = "dm" if kind == "ROOM_KIND_DM" else "group"
  1267. return {
  1268. "name": name,
  1269. "type": chat_type,
  1270. }
  1271. # ------------------------------------------------------------------ #
  1272. # Reactions
  1273. # ------------------------------------------------------------------ #
  1274. @staticmethod
  1275. def _emoji_to_shortcode(emoji: str) -> str:
  1276. """Convert a unicode emoji to a Chatto shortcode name.
  1277. If the emoji is already a shortcode (no unicode mapping found),
  1278. return it as-is.
  1279. """
  1280. shortcode = _EMOJI_TO_SHORTCODE.get(emoji)
  1281. if shortcode:
  1282. return shortcode
  1283. # Already a shortcode like "thumbsup" — return as-is
  1284. return emoji
  1285. async def send_reaction(self, chat_id: str, message_id: str, emoji: str) -> bool:
  1286. """Add a reaction to a message via MessageService/AddReaction."""
  1287. shortcode = self._emoji_to_shortcode(emoji)
  1288. try:
  1289. try:
  1290. client = await self._require_client()
  1291. except RuntimeError:
  1292. return False
  1293. result = await client.add_reaction(
  1294. room_id=str(chat_id),
  1295. message_event_id=str(message_id),
  1296. emoji=shortcode,
  1297. )
  1298. return result
  1299. except ChattoError as e:
  1300. logger.debug("Chatto: AddReaction failed: %s", e)
  1301. return False
  1302. except Exception as e:
  1303. logger.debug("Chatto: AddReaction error: %s", e)
  1304. return False
  1305. async def remove_reaction(self, chat_id: str, message_id: str, emoji: str) -> bool:
  1306. """Remove a reaction from a message via MessageService/RemoveReaction."""
  1307. shortcode = self._emoji_to_shortcode(emoji)
  1308. try:
  1309. try:
  1310. client = await self._require_client()
  1311. except RuntimeError:
  1312. return False
  1313. result = await client.remove_reaction(
  1314. room_id=str(chat_id),
  1315. message_event_id=str(message_id),
  1316. emoji=shortcode,
  1317. )
  1318. return result
  1319. except ChattoError as e:
  1320. logger.debug("Chatto: RemoveReaction failed: %s", e)
  1321. return False
  1322. except Exception as e:
  1323. logger.debug("Chatto: RemoveReaction error: %s", e)
  1324. return False
  1325. # ------------------------------------------------------------------ #
  1326. # Read state management (Chatto-unique)
  1327. # ------------------------------------------------------------------ #
  1328. async def mark_room_as_read(self, room_id: str) -> bool:
  1329. """Mark a room as read via RoomService/MarkRoomAsRead."""
  1330. try:
  1331. try:
  1332. client = await self._require_client()
  1333. except RuntimeError:
  1334. return False
  1335. await client.mark_room_as_read(room_id=str(room_id))
  1336. return True
  1337. except ChattoError as e:
  1338. logger.debug("Chatto: MarkRoomAsRead failed: %s", e)
  1339. return False
  1340. except Exception as e:
  1341. logger.debug("Chatto: MarkRoomAsRead error: %s", e)
  1342. return False
  1343. async def mark_thread_as_read(self, room_id: str, thread_root_event_id: str) -> bool:
  1344. """Mark a thread as read via ThreadService/MarkThreadAsRead."""
  1345. try:
  1346. try:
  1347. client = await self._require_client()
  1348. except RuntimeError:
  1349. return False
  1350. await client.mark_thread_as_read(
  1351. room_id=str(room_id), thread_root_event_id=str(thread_root_event_id)
  1352. )
  1353. return True
  1354. except ChattoError as e:
  1355. logger.debug("Chatto: MarkThreadAsRead failed: %s", e)
  1356. return False
  1357. except Exception as e:
  1358. logger.debug("Chatto: MarkThreadAsRead error: %s", e)
  1359. return False
  1360. # ------------------------------------------------------------------ #
  1361. # DM initiation (Chatto-unique)
  1362. # ------------------------------------------------------------------ #
  1363. async def start_dm(self, user_id: str) -> Optional[str]:
  1364. """Start a direct message with a user via RoomService/StartDM.
  1365. Returns the room ID on success, or None on failure.
  1366. """
  1367. if not user_id:
  1368. return None
  1369. try:
  1370. try:
  1371. client = await self._require_client()
  1372. except RuntimeError:
  1373. return None
  1374. room = await client.start_dm(participant_ids=[str(user_id)])
  1375. rid = str(getattr(cast(Any, room), "id", "")) if room else ""
  1376. if rid:
  1377. self._room_names[rid] = self._room_names.get(rid, "")
  1378. self._room_kinds[rid] = "ROOM_KIND_DM"
  1379. return rid
  1380. logger.debug("Chatto: StartDM returned no room id")
  1381. return None
  1382. except ChattoError as e:
  1383. logger.debug("Chatto: StartDM failed: %s", e)
  1384. return None
  1385. except Exception as e:
  1386. logger.debug("Chatto: StartDM error: %s", e)
  1387. return None
  1388. # ------------------------------------------------------------------ #
  1389. # Thread following (Chatto-unique)
  1390. # ------------------------------------------------------------------ #
  1391. async def _follow_thread(self, room_id: str, thread_root_event_id: str) -> None:
  1392. """Best-effort: follow a thread via ThreadService/FollowThread."""
  1393. try:
  1394. try:
  1395. client = await self._require_client()
  1396. except RuntimeError:
  1397. return
  1398. await client.follow_thread(
  1399. room_id=str(room_id), thread_root_event_id=str(thread_root_event_id)
  1400. )
  1401. except ChattoError as e:
  1402. logger.debug("Chatto: FollowThread failed: %s", e)
  1403. except Exception as e:
  1404. logger.debug("Chatto: FollowThread error: %s", e)
  1405. # ------------------------------------------------------------------ #
  1406. # Room creation (Chatto-unique)
  1407. # ------------------------------------------------------------------ #
  1408. async def create_room(
  1409. self,
  1410. name: str,
  1411. description: str = "",
  1412. group_id: str = "",
  1413. universal: bool = True,
  1414. ) -> Optional[str]:
  1415. """Create an ad-hoc room via RoomService/CreateRoom.
  1416. Returns the room ID on success, or None on failure.
  1417. """
  1418. try:
  1419. try:
  1420. client = await self._require_client()
  1421. except RuntimeError:
  1422. return None
  1423. room = await client.create_room(
  1424. name=name,
  1425. group_id=group_id or "",
  1426. description=description,
  1427. universal=universal,
  1428. )
  1429. rid = str(room.id) if room else ""
  1430. if rid:
  1431. self._room_names[rid] = name
  1432. self._room_kinds[rid] = "ROOM_KIND_GROUP"
  1433. return rid
  1434. logger.debug("Chatto: CreateRoom returned no room id")
  1435. return None
  1436. except ChattoError as e:
  1437. logger.debug("Chatto: CreateRoom failed: %s", e)
  1438. return None
  1439. except Exception as e:
  1440. logger.debug("Chatto: CreateRoom error: %s", e)
  1441. return None
  1442. # ------------------------------------------------------------------ #
  1443. # Notification dismissal (Chatto-unique)
  1444. # ------------------------------------------------------------------ #
  1445. async def dismiss_all_notifications(self) -> bool:
  1446. """Dismiss all notifications via NotificationService/DismissAllNotifications."""
  1447. try:
  1448. try:
  1449. client = await self._require_client()
  1450. except RuntimeError:
  1451. return False
  1452. await client.dismiss_all_notifications()
  1453. return True
  1454. except ChattoError as e:
  1455. logger.debug("Chatto: DismissAllNotifications failed: %s", e)
  1456. return False
  1457. except Exception as e:
  1458. logger.debug("Chatto: DismissAllNotifications error: %s", e)
  1459. return False
  1460. async def dismiss_notification(self, notification_id: str) -> bool:
  1461. """Dismiss a single notification via NotificationService/DismissNotification."""
  1462. try:
  1463. try:
  1464. client = await self._require_client()
  1465. except RuntimeError:
  1466. return False
  1467. await client.dismiss_notification(notification_id=str(notification_id))
  1468. return True
  1469. except ChattoError as e:
  1470. logger.debug("Chatto: DismissNotification failed: %s", e)
  1471. return False
  1472. except Exception as e:
  1473. logger.debug("Chatto: DismissNotification error: %s", e)
  1474. return False
  1475. # ------------------------------------------------------------------ #
  1476. # Message editing and deletion
  1477. # ------------------------------------------------------------------ #
  1478. async def edit_message(
  1479. self,
  1480. chat_id: str,
  1481. message_id: str,
  1482. new_content: str,
  1483. metadata: Optional[Dict[str, Any]] = None,
  1484. ) -> bool:
  1485. """Edit a previously sent message via MessageService/UpdateMessage."""
  1486. try:
  1487. try:
  1488. client = await self._require_client()
  1489. except RuntimeError:
  1490. return False
  1491. await client.update_message(
  1492. room_id=str(chat_id),
  1493. event_id=str(message_id),
  1494. body=new_content,
  1495. )
  1496. return True
  1497. except ChattoError as e:
  1498. logger.debug("Chatto: UpdateMessage failed: %s", e)
  1499. return False
  1500. except Exception as e:
  1501. logger.debug("Chatto: UpdateMessage error: %s", e)
  1502. return False
  1503. async def delete_message(
  1504. self,
  1505. chat_id: str,
  1506. message_id: str,
  1507. metadata: Optional[Dict[str, Any]] = None,
  1508. ) -> bool:
  1509. """Delete a previously sent message via MessageService/DeleteMessage."""
  1510. try:
  1511. try:
  1512. client = await self._require_client()
  1513. except RuntimeError:
  1514. return False
  1515. result = await client.delete_message(
  1516. room_id=str(chat_id),
  1517. event_id=str(message_id),
  1518. )
  1519. return result
  1520. except ChattoError as e:
  1521. logger.debug("Chatto: DeleteMessage failed: %s", e)
  1522. return False
  1523. except Exception as e:
  1524. logger.debug("Chatto: DeleteMessage error: %s", e)
  1525. return False
  1526. # ------------------------------------------------------------------ #
  1527. # Processing lifecycle hooks (reactions-based, like Discord)
  1528. # ------------------------------------------------------------------ #
  1529. def _reactions_enabled(self) -> bool:
  1530. """Check if processing reactions are enabled."""
  1531. return os.getenv("CHATTO_REACTIONS", "true").lower() not in {"false", "0", "no"}
  1532. def _event_room_and_message_id(self, event: MessageEvent) -> Tuple[str, str]:
  1533. """Extract room_id and message_id from a MessageEvent."""
  1534. chat_id = ""
  1535. message_id = str(event.message_id or "")
  1536. source = event.source
  1537. if source:
  1538. chat_id = str(getattr(source, "chat_id", "") or "")
  1539. # Fallback: try raw_message dict
  1540. if not chat_id or not message_id:
  1541. raw = event.raw_message
  1542. if isinstance(raw, dict):
  1543. if not chat_id:
  1544. chat_id = str(raw.get("roomId", "") or "")
  1545. if not message_id:
  1546. message_id = str(raw.get("id", "") or "")
  1547. return chat_id, message_id
  1548. async def on_processing_start(self, event: MessageEvent) -> None:
  1549. """Add an 👀 (eyes) reaction to the incoming message."""
  1550. if not self._reactions_enabled():
  1551. return
  1552. chat_id, message_id = self._event_room_and_message_id(event)
  1553. if not chat_id or not message_id:
  1554. return
  1555. await self.send_reaction(chat_id, message_id, "👀")
  1556. async def on_processing_complete(
  1557. self, event: MessageEvent, outcome: ProcessingOutcome
  1558. ) -> None:
  1559. """Swap the 👀 reaction for ✅ (success) or ❌ (failure)."""
  1560. if not self._reactions_enabled():
  1561. return
  1562. chat_id, message_id = self._event_room_and_message_id(event)
  1563. if not chat_id or not message_id:
  1564. return
  1565. # Remove the processing eyes reaction
  1566. await self.remove_reaction(chat_id, message_id, "👀")
  1567. # Add the outcome reaction
  1568. if outcome == ProcessingOutcome.SUCCESS:
  1569. await self.send_reaction(chat_id, message_id, "✅")
  1570. elif outcome == ProcessingOutcome.FAILURE:
  1571. await self.send_reaction(chat_id, message_id, "❌")
  1572. # ------------------------------------------------------------------ #
  1573. # Asset upload (chunked)
  1574. # ------------------------------------------------------------------ #
  1575. async def _upload_asset(self, room_id: str, file_path: str) -> Optional[str]:
  1576. """Upload a file via the chunked AssetUploadService.
  1577. Returns the asset ID on success, or None on failure.
  1578. """
  1579. try:
  1580. with open(file_path, "rb") as f:
  1581. file_data = f.read()
  1582. except Exception as e:
  1583. logger.error("Chatto: failed to read file %s — %s", file_path, e)
  1584. return None
  1585. if not file_data:
  1586. logger.error("Chatto: file %s is empty", file_path)
  1587. return None
  1588. file_size = len(file_data)
  1589. file_name = os.path.basename(file_path)
  1590. mime_type = mimetypes.guess_type(file_path)[0] or "application/octet-stream"
  1591. sha256_hash = hashlib.sha256(file_data).hexdigest()
  1592. try:
  1593. try:
  1594. client = await self._require_client()
  1595. except RuntimeError:
  1596. logger.error("Chatto: upload aborted - no client available")
  1597. return None
  1598. # Step 1: Create upload session
  1599. upload = await client.create_upload(
  1600. room_id=room_id,
  1601. filename=file_name,
  1602. size=file_size,
  1603. sha256=sha256_hash,
  1604. content_type=mime_type,
  1605. )
  1606. upload_id = str(getattr(cast(Any, upload), "id", ""))
  1607. if not upload_id:
  1608. logger.error("Chatto: CreateUpload returned no upload ID")
  1609. return None
  1610. # Step 2: Upload chunks
  1611. offset = 0
  1612. while offset < file_size:
  1613. chunk = file_data[offset:offset + _UPLOAD_CHUNK_SIZE]
  1614. chunk_sha256 = hashlib.sha256(chunk).hexdigest()
  1615. await client.upload_chunk(
  1616. upload_id=upload_id,
  1617. offset=offset,
  1618. content=chunk,
  1619. chunk_sha256=chunk_sha256,
  1620. )
  1621. offset += len(chunk)
  1622. # Step 3: Complete upload
  1623. upload, asset = await client.complete_upload(upload_id=upload_id)
  1624. if not asset:
  1625. logger.error("Chatto: CompleteUpload returned no asset")
  1626. return None
  1627. asset_id = str(getattr(cast(Any, asset), "id", ""))
  1628. logger.info("Chatto: uploaded %s as asset %s (%d bytes)", file_name, asset_id, file_size)
  1629. return asset_id
  1630. except ChattoError as e:
  1631. logger.error("Chatto: upload failed: %s", e)
  1632. return None
  1633. except Exception as e:
  1634. logger.error("Chatto: upload error: %s", e)
  1635. return None
  1636. async def send_image_file(
  1637. self,
  1638. chat_id: str,
  1639. file_path: str,
  1640. caption: Optional[str] = None,
  1641. reply_to: Optional[str] = None,
  1642. metadata: Optional[Dict[str, Any]] = None,
  1643. ) -> SendResult:
  1644. """Send a local image file via the chunked upload API."""
  1645. # Validate the path is safe
  1646. safe_path = self.validate_media_delivery_path(file_path)
  1647. if not safe_path:
  1648. logger.warning("Chatto: send_image_file — unsafe path %s", file_path)
  1649. text = "⚠️ Couldn't deliver the image attachment."
  1650. if caption:
  1651. text = f"{caption}\n{text}"
  1652. return await self.send(chat_id, text, reply_to=reply_to, metadata=metadata)
  1653. asset_id = await self._upload_asset(str(chat_id), safe_path)
  1654. if not asset_id:
  1655. # Fallback to a notice
  1656. text = "⚠️ Couldn't deliver the image attachment."
  1657. if caption:
  1658. text = f"{caption}\n{text}"
  1659. return await self.send(chat_id, text, reply_to=reply_to, metadata=metadata)
  1660. thread_id = (metadata or {}).get("thread_id")
  1661. if reply_to:
  1662. thread_id = reply_to
  1663. try:
  1664. try:
  1665. client = await self._require_client()
  1666. except RuntimeError:
  1667. return SendResult(success=False, error="Chatto client not available", retryable=True)
  1668. msg = await client.post_message(
  1669. room_id=str(chat_id),
  1670. body=caption or "",
  1671. attachment_asset_ids=[asset_id],
  1672. thread_root_event_id=str(thread_id) if thread_id else "",
  1673. )
  1674. msg_id = str(msg.id)
  1675. if msg_id:
  1676. self._mark_seen(str(chat_id), msg_id)
  1677. return SendResult(success=True, message_id=msg_id, raw_response=msg)
  1678. except ChattoError as e:
  1679. return SendResult(success=False, error=str(e), retryable=True)
  1680. except Exception as e:
  1681. return SendResult(success=False, error=str(e), retryable=False)
  1682. async def send_image(
  1683. self,
  1684. chat_id: str,
  1685. image_url: str,
  1686. caption: Optional[str] = None,
  1687. reply_to: Optional[str] = None,
  1688. metadata: Optional[Dict[str, Any]] = None,
  1689. ) -> SendResult:
  1690. """Send an image to a Chatto room.
  1691. Tries to download the image from the URL and upload it as a native
  1692. attachment. Falls back to sending the URL as a link (Chatto renders
  1693. link previews) if the download fails.
  1694. """
  1695. # Try downloading and uploading as attachment
  1696. try:
  1697. import tempfile
  1698. import urllib.request as _urllib_request
  1699. # Download to a temp file
  1700. parsed = urlsplit(image_url)
  1701. url_path = parsed.path
  1702. ext = os.path.splitext(url_path)[1] or ".png"
  1703. tmp_fd, tmp_path = tempfile.mkstemp(suffix=ext, prefix="chatto_img_")
  1704. try:
  1705. os.close(tmp_fd)
  1706. req = _urllib_request.Request(image_url, headers={"User-Agent": "Hermes/1.0"})
  1707. try:
  1708. import ssl
  1709. ctx = ssl.create_default_context()
  1710. except Exception:
  1711. ctx = None
  1712. with _urllib_request.urlopen(req, timeout=_HTTP_TIMEOUT, context=ctx) as resp:
  1713. with open(tmp_path, "wb") as f:
  1714. f.write(resp.read())
  1715. # Upload as attachment
  1716. result = await self.send_image_file(
  1717. chat_id, tmp_path, caption=caption,
  1718. reply_to=reply_to, metadata=metadata,
  1719. )
  1720. if result.success:
  1721. return result
  1722. finally:
  1723. try:
  1724. os.unlink(tmp_path)
  1725. except OSError:
  1726. pass
  1727. except Exception as e:
  1728. logger.debug("Chatto: send_image download/upload failed, falling back to link: %s", e)
  1729. # Fallback: send as link (Chatto renders link previews)
  1730. text = image_url
  1731. if caption:
  1732. text = f"{caption}\n{image_url}"
  1733. return await self.send(chat_id, text, reply_to=reply_to, metadata=metadata)
  1734. # ------------------------------------------------------------------ #
  1735. # Platform properties
  1736. # ------------------------------------------------------------------ #
  1737. @property
  1738. def platform_name(self) -> str:
  1739. return "chatto"
  1740. @property
  1741. def supports_markdown(self) -> bool:
  1742. return True
  1743. @property
  1744. def supports_reactions(self) -> bool:
  1745. return True
  1746. # ------------------------------------------------------------------ #
  1747. # Member directory — user lookup and mention resolution (Chatto-unique)
  1748. # ------------------------------------------------------------------ #
  1749. async def list_users(self) -> list:
  1750. """List all server members via UserService/ListUsers.
  1751. Returns a list of user dicts. Each dict typically contains
  1752. ``id``, ``login``, and ``displayName`` keys.
  1753. """
  1754. try:
  1755. try:
  1756. client = await self._require_client()
  1757. except RuntimeError:
  1758. return []
  1759. members, _ = await client.list_users()
  1760. users = []
  1761. # Cache all returned users and convert to dict format
  1762. for member in members:
  1763. if member and member.user:
  1764. user_dict = {
  1765. "id": str(member.user.id),
  1766. "login": str(member.user.login),
  1767. "displayName": str(member.user.display_name or ""),
  1768. }
  1769. uid = user_dict["id"]
  1770. if uid:
  1771. self._user_cache[uid] = user_dict
  1772. users.append(user_dict)
  1773. return users
  1774. except ChattoError as e:
  1775. logger.debug("Chatto: ListUsers failed: %s", e)
  1776. return []
  1777. except Exception as e:
  1778. logger.debug("Chatto: ListUsers error: %s", e)
  1779. return []
  1780. async def get_user(self, user_id: str) -> Optional[dict]:
  1781. """Get a single user by ID via UserService/GetUser.
  1782. Returns the user dict (containing ``id``, ``login``,
  1783. ``displayName``) or ``None`` on failure. Results are cached in
  1784. ``self._user_cache``.
  1785. """
  1786. if not user_id:
  1787. return None
  1788. # Return cached entry if available
  1789. if user_id in self._user_cache:
  1790. return self._user_cache[user_id]
  1791. try:
  1792. try:
  1793. client = await self._require_client()
  1794. except RuntimeError:
  1795. return None
  1796. user_obj = await client.get_user(user_id=str(user_id))
  1797. if user_obj and user_obj.user:
  1798. user_dict = {
  1799. "id": str(user_obj.user.id),
  1800. "login": str(user_obj.user.login),
  1801. "displayName": str(user_obj.user.display_name or ""),
  1802. }
  1803. uid = user_dict["id"]
  1804. if uid:
  1805. self._user_cache[uid] = user_dict
  1806. return user_dict
  1807. return None
  1808. except ChattoError as e:
  1809. logger.debug("Chatto: GetUser failed: %s", e)
  1810. return None
  1811. except Exception as e:
  1812. logger.debug("Chatto: GetUser error: %s", e)
  1813. return None
  1814. async def batch_get_users(self, user_ids: list) -> list:
  1815. """Batch-fetch multiple users via UserService/BatchGetUsers.
  1816. Returns a list of user dicts. Cached entries are reused and only
  1817. uncached IDs are fetched from the server.
  1818. """
  1819. if not user_ids:
  1820. return []
  1821. # Separate cached from uncached
  1822. cached: list = []
  1823. uncached_ids: list = []
  1824. for uid in user_ids:
  1825. uid_str = str(uid)
  1826. if uid_str in self._user_cache:
  1827. cached.append(self._user_cache[uid_str])
  1828. else:
  1829. uncached_ids.append(uid_str)
  1830. if not uncached_ids:
  1831. return cached
  1832. try:
  1833. try:
  1834. client = await self._require_client()
  1835. except RuntimeError:
  1836. return cached
  1837. members = await client.batch_get_users(user_ids=uncached_ids)
  1838. fetched = []
  1839. for member in members:
  1840. if member and member.user:
  1841. user_dict = {
  1842. "id": str(member.user.id),
  1843. "login": str(member.user.login),
  1844. "displayName": str(member.user.display_name or ""),
  1845. }
  1846. uid = user_dict["id"]
  1847. if uid:
  1848. self._user_cache[uid] = user_dict
  1849. fetched.append(user_dict)
  1850. return cached + fetched
  1851. except ChattoError as e:
  1852. logger.debug("Chatto: BatchGetUsers failed: %s", e)
  1853. return cached
  1854. except Exception as e:
  1855. logger.debug("Chatto: BatchGetUsers error: %s", e)
  1856. return cached
  1857. # ------------------------------------------------------------------ #
  1858. # Presence broadcasting (Chatto-unique)
  1859. # ------------------------------------------------------------------ #
  1860. async def set_presence(self, status: str) -> bool:
  1861. """Update the bot's presence status via MyAccountService/UpdatePresence.
  1862. Accepts string values ``"online"``, ``"away"``, ``"dnd"`` (or
  1863. ``"do_not_disturb"``) and maps them to Chatto's PresenceStatus enum.
  1864. Returns ``True`` on success.
  1865. """
  1866. # Map string status to chattolib PresenceStatus enum
  1867. status_map = {
  1868. "online": PresenceStatus.ONLINE,
  1869. "away": PresenceStatus.AWAY,
  1870. "dnd": PresenceStatus.DO_NOT_DISTURB,
  1871. "do_not_disturb": PresenceStatus.DO_NOT_DISTURB,
  1872. }
  1873. status_lower = status.lower().strip()
  1874. presence_status = status_map.get(status_lower)
  1875. if presence_status is None:
  1876. logger.warning("Chatto: unknown presence status %r", status)
  1877. return False
  1878. try:
  1879. try:
  1880. client = await self._require_client()
  1881. except RuntimeError:
  1882. return False
  1883. await client.update_presence(status=presence_status)
  1884. logger.debug("Chatto: presence set to %s", status_lower)
  1885. return True
  1886. except ChattoError as e:
  1887. logger.debug("Chatto: UpdatePresence failed: %s", e)
  1888. return False
  1889. except Exception as e:
  1890. logger.debug("Chatto: UpdatePresence error: %s", e)
  1891. return False
  1892. # ------------------------------------------------------------------ #
  1893. # Custom status messages (Chatto-unique)
  1894. # ------------------------------------------------------------------ #
  1895. async def set_custom_status(self, text: str) -> bool:
  1896. """Set a custom status message via MyAccountService/UpdateCustomStatus.
  1897. The status text is a plain string (max ~100 chars). Useful for
  1898. indicating long-running operations, e.g. ``"Processing..."``.
  1899. Returns ``True`` on success.
  1900. """
  1901. if not text:
  1902. return False
  1903. # Truncate to a reasonable length
  1904. status_text = text.strip()[:100]
  1905. if not status_text:
  1906. return False
  1907. try:
  1908. try:
  1909. client = await self._require_client()
  1910. except RuntimeError:
  1911. return False
  1912. await client.update_custom_status(emoji="", text=status_text)
  1913. logger.debug("Chatto: custom status set to %r", status_text)
  1914. return True
  1915. except ChattoError as e:
  1916. logger.debug("Chatto: UpdateCustomStatus failed: %s", e)
  1917. return False
  1918. except Exception as e:
  1919. logger.debug("Chatto: UpdateCustomStatus error: %s", e)
  1920. return False
  1921. async def clear_custom_status(self) -> bool:
  1922. """Clear the custom status message via MyAccountService/DeleteCustomStatus.
  1923. Returns ``True`` on success.
  1924. """
  1925. try:
  1926. try:
  1927. client = await self._require_client()
  1928. except RuntimeError:
  1929. return False
  1930. await client.delete_custom_status()
  1931. logger.debug("Chatto: custom status cleared")
  1932. return True
  1933. except ChattoError as e:
  1934. logger.debug("Chatto: DeleteCustomStatus failed: %s", e)
  1935. return False
  1936. except Exception as e:
  1937. logger.debug("Chatto: DeleteCustomStatus error: %s", e)
  1938. return False
  1939. @property
  1940. def supports_threads(self) -> bool:
  1941. return True
  1942. # --------------------------------------------------------------------------- #
  1943. # Plugin registration
  1944. # --------------------------------------------------------------------------- #
  1945. def check_requirements() -> bool:
  1946. """Check if Chatto is configured and dependencies are available."""
  1947. # Then check configuration
  1948. return bool(
  1949. os.getenv("CHATTO_URL", "").strip()
  1950. and os.getenv("CHATTO_LOGIN", "").strip()
  1951. and os.getenv("CHATTO_PASSWORD", "").strip()
  1952. )
  1953. def validate_config(config: PlatformConfig) -> bool:
  1954. """Validate that the platform config has enough info to connect."""
  1955. extra = getattr(config, "extra", {}) or {}
  1956. url = os.getenv("CHATTO_URL") or str(extra.get("url", ""))
  1957. login = os.getenv("CHATTO_LOGIN", "").strip()
  1958. password = os.getenv("CHATTO_PASSWORD", "").strip()
  1959. return bool(url and login and password)
  1960. def is_connected(config) -> bool:
  1961. """Check whether Chatto is configured."""
  1962. return validate_config(config)
  1963. def _apply_yaml_config(yaml_cfg: dict, chatto_cfg: dict) -> Optional[dict]:
  1964. """Translate config.yaml chatto.extra keys into CHATTO_* env vars."""
  1965. if not isinstance(chatto_cfg, dict):
  1966. chatto_cfg = {}
  1967. extra = chatto_cfg.get("extra", {}) or {}
  1968. if not isinstance(extra, dict):
  1969. extra = {}
  1970. mapping = {
  1971. "url": "CHATTO_URL",
  1972. "home_channel": "CHATTO_HOME_CHANNEL",
  1973. "require_mention": "CHATTO_REQUIRE_MENTION",
  1974. "free_response_channels": "CHATTO_FREE_RESPONSE_CHANNELS",
  1975. "auto_thread": "CHATTO_AUTO_THREAD",
  1976. }
  1977. for yaml_key, env_key in mapping.items():
  1978. val = extra.get(yaml_key)
  1979. if val is not None and not os.getenv(env_key):
  1980. if isinstance(val, bool):
  1981. env_val: str = str(val).lower()
  1982. elif isinstance(val, list):
  1983. env_val = ",".join(str(v) for v in val)
  1984. else:
  1985. env_val = str(val)
  1986. os.environ[env_key] = env_val
  1987. channels = extra.get("channels")
  1988. if isinstance(channels, list) and not os.getenv("CHATTO_CHANNELS"):
  1989. channels_str: str = ",".join(str(c) for c in channels)
  1990. os.environ["CHATTO_CHANNELS"] = channels_str
  1991. allowed = extra.get("allowed_users")
  1992. if isinstance(allowed, list) and not os.getenv("CHATTO_ALLOWED_USERS"):
  1993. allowed_str: str = ",".join(str(u) for u in allowed)
  1994. os.environ["CHATTO_ALLOWED_USERS"] = allowed_str
  1995. if "allow_all_users" in extra and not os.getenv("CHATTO_ALLOW_ALL_USERS"):
  1996. allow_all_val: str = str(extra["allow_all_users"]).lower()
  1997. os.environ["CHATTO_ALLOW_ALL_USERS"] = allow_all_val
  1998. # Return nothing to merge — all config flows through env
  1999. return None
  2000. def _env_enablement() -> Optional[dict]:
  2001. """Seed PlatformConfig.extra from env vars for env-only setups.
  2002. Returns a dict compatible with the PlatformConfig merge hook (or None
  2003. when no env-provided values are present).
  2004. """
  2005. extra_dict: dict = {}
  2006. env_url = os.getenv("CHATTO_URL", "").strip()
  2007. if env_url:
  2008. extra_dict["url"] = env_url
  2009. env_channels = os.getenv("CHATTO_CHANNELS", "").strip()
  2010. if env_channels:
  2011. extra_dict["channels"] = [c.strip() for c in env_channels.split(",") if c.strip()]
  2012. env_require_mention_str = str(utils.is_truthy_value(os.getenv("CHATTO_REQUIRE_MENTION", "")))
  2013. extra_dict["require_mention"] = env_require_mention_str
  2014. env_home_channel = os.getenv("CHATTO_HOME_CHANNEL", "").strip()
  2015. if env_home_channel:
  2016. home_dict = {"home_channel": env_home_channel}
  2017. return {**extra_dict, **home_dict}
  2018. async def _standalone_send(
  2019. pconfig,
  2020. chat_id,
  2021. message,
  2022. *,
  2023. thread_id=None,
  2024. media_files=None,
  2025. force_document=False,
  2026. ) -> SendResult:
  2027. """Deliver a message to Chatto without a running gateway adapter. Do not modify signature.
  2028. Used by cron / scheduled routines that run out-of-process. Creates a
  2029. short-lived chattolib client, posts, and closes.
  2030. """
  2031. pconfig.getattr("extra", {})
  2032. chatto_config = ChattoConfig(extra=pconfig.extra)
  2033. # Create a temporary client for standalone sending
  2034. client: ChattoClient
  2035. token = chatto_config._token
  2036. login = chatto_config._login
  2037. password = chatto_config._password
  2038. if not chatto_config._base_url or (not (login or not password) or not token):
  2039. return SendResult(success=False, error="Chatto: base URL or credentials missing")
  2040. try:
  2041. if token:
  2042. client = ChattoClient(base_url=chatto_config._base_url, token=token)
  2043. else:
  2044. client = await ChattoClient.login(
  2045. base_url=chatto_config._base_url, login=chatto_config._login, password=chatto_config._password,
  2046. )
  2047. except Exception as exc:
  2048. return SendResult(success=False, error=f"Chatto login failed: {exc}")
  2049. try:
  2050. kwargs: Dict[str, Any] = {}
  2051. if chatto_config._auto_thread and thread_id:
  2052. kwargs["in_reply_to"] = thread_id
  2053. if media_files and media_files.get("attachment_asset_ids"):
  2054. kwargs["attachment_asset_ids"] = list(media_files["attachment_asset_ids"])
  2055. try:
  2056. posted = await client.post_message(chat_id, message, **kwargs)
  2057. except Exception as exc:
  2058. return SendResult(success=False, error=str(exc))
  2059. return SendResult(success=True, message_id=getattr(posted, "id", "") or None)
  2060. finally:
  2061. try:
  2062. await client.close()
  2063. except Exception:
  2064. logger.error("Chatto standalone: error closing short-lived client")
  2065. await client.login(login=login, password=password)
  2066. def interactive_setup() -> None:
  2067. """Interactive setup wizard for Chatto."""
  2068. try:
  2069. from hermes_cli.gateway import prompt_env, set_env_var # type: ignore
  2070. except Exception:
  2071. # Fallback simple prompt implementations for environments where
  2072. # hermes_cli.gateway helpers are not available (tests / limited shells).
  2073. def prompt_env(prompt_text: str, password: bool = False):
  2074. if password:
  2075. import getpass
  2076. return getpass.getpass(prompt_text)
  2077. return input(prompt_text + " ")
  2078. def set_env_var(k: str, v: str) -> None:
  2079. os.environ[k] = v
  2080. url = prompt_env("Chatto server URL (e.g. https://chat.example.com):")
  2081. if url:
  2082. set_env_var("CHATTO_URL", url)
  2083. login = prompt_env("Chatto login (username):")
  2084. if login:
  2085. set_env_var("CHATTO_LOGIN", login)
  2086. password = prompt_env("Chatto password:", password=True)
  2087. if password:
  2088. set_env_var("CHATTO_PASSWORD", password)
  2089. channels = prompt_env("Room IDs to watch (comma-separated, or empty for all):")
  2090. if channels:
  2091. set_env_var("CHATTO_CHANNELS", channels)
  2092. home = prompt_env("Home room ID for notifications (or empty):")
  2093. if home:
  2094. set_env_var("CHATTO_HOME_CHANNEL", home)
  2095. allow_all = prompt_env("Allow all users? (true/false):")
  2096. if allow_all:
  2097. set_env_var("CHATTO_ALLOW_ALL_USERS", allow_all)
  2098. print("\n✓ Chatto configured. Restart the gateway to activate.")
  2099. def register(ctx) -> None:
  2100. """Plugin entry point — called by the Hermes plugin system."""
  2101. ctx.register_platform(
  2102. name="chatto",
  2103. label="Chatto Chat Server",
  2104. adapter_factory=lambda cfg: ChattoAdapter(cfg),
  2105. check_fn=check_requirements,
  2106. validate_config=validate_config,
  2107. is_connected=is_connected,
  2108. required_env=["CHATTO_URL", "CHATTO_LOGIN", "CHATTO_PASSWORD"],
  2109. install_hint="Requires a Chatto server. See https://docs.chatto.run",
  2110. env_enablement_fn=_env_enablement,
  2111. setup_fn=interactive_setup,
  2112. apply_yaml_config_fn=_apply_yaml_config,
  2113. cron_deliver_env_var="CHATTO_HOME_CHANNEL",
  2114. standalone_sender_fn=_standalone_send,
  2115. allowed_users_env="CHATTO_ALLOWED_USERS",
  2116. allow_all_env="CHATTO_ALLOW_ALL_USERS",
  2117. auto_thread_env="CHATTO_AUTO_THREAD",
  2118. max_message_length=_MAX_MESSAGE_LENGTH,
  2119. emoji="💬",
  2120. allow_update_command=True,
  2121. pii_safe=False,
  2122. platform_hint=(
  2123. "You are chatting in Chatto (a self-hosted team chat server). "
  2124. "Markdown IS supported. Users address you by @-mentioning your name "
  2125. "in rooms; direct messages reach you without a mention. "
  2126. "Keep responses conversational."
  2127. ),
  2128. )