client.py 79 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355
  1. """Main async client for the Chatto Connect API.
  2. Chatto migrated from GraphQL to a protobuf-first Connect API in v0.4.x
  3. (see ADR-042). The client speaks Connect binary protobuf via the official
  4. ``connectrpc`` Python package and the generated service stubs under
  5. ``chattolib._pb`` for all request/response operations. Realtime events
  6. live in ``chattolib.realtime``.
  7. """
  8. # mypy: disable-error-code="no-any-return"
  9. # Rationale: attribute access on generated protobuf messages is Any-typed
  10. # from mypy's perspective (the generated modules skip type checking via
  11. # follow_imports=skip). The runtime types are exactly what the return-type
  12. # annotations claim.
  13. from __future__ import annotations
  14. import hashlib
  15. from collections.abc import Awaitable
  16. from datetime import datetime
  17. from pathlib import Path
  18. from typing import Any, TypeVar
  19. import httpx
  20. # ConnectError isn't re-exported publicly by connectrpc.__init__; import from
  21. # its submodule so the top-level import path stays clean for callers.
  22. from connectrpc.errors import ConnectError # noqa: E402
  23. from chattolib._pb.chatto.admin.v1 import (
  24. event_log_pb2,
  25. room_layout_pb2,
  26. )
  27. from chattolib._pb.chatto.admin.v1 import (
  28. members_pb2 as admin_members_pb2,
  29. )
  30. from chattolib._pb.chatto.admin.v1 import (
  31. permissions_pb2 as admin_permissions_pb2,
  32. )
  33. from chattolib._pb.chatto.admin.v1 import (
  34. roles_pb2 as admin_roles_pb2,
  35. )
  36. from chattolib._pb.chatto.admin.v1 import (
  37. server_pb2 as admin_server_pb2,
  38. )
  39. from chattolib._pb.chatto.api.v1 import (
  40. account_pb2,
  41. asset_uploads_pb2,
  42. attachments_pb2,
  43. common_pb2,
  44. external_identities_pb2,
  45. link_previews_pb2,
  46. member_directory_pb2,
  47. messages_pb2,
  48. notification_preferences_pb2,
  49. notifications_pb2,
  50. pagination_pb2,
  51. presence_pb2,
  52. push_notifications_pb2,
  53. reactions_pb2,
  54. read_state_pb2,
  55. roles_pb2,
  56. room_directory_pb2,
  57. room_timeline_pb2,
  58. rooms_pb2,
  59. server_state_pb2,
  60. threads_pb2,
  61. user_status_pb2,
  62. viewer_pb2,
  63. voice_calls_pb2,
  64. )
  65. from chattolib._pb.chatto.auth.v1 import external_identity_auth_pb2
  66. from chattolib._pb.chatto.discovery.v1 import server_pb2 as discovery_server_pb2
  67. from chattolib._transport import (
  68. ServiceClients,
  69. build_service_clients,
  70. pb_to_dict,
  71. translate_connect_error,
  72. )
  73. from chattolib.exceptions import ChattoAuthError, ChattoError
  74. from chattolib.types import (
  75. ActiveCall,
  76. AdminEventLogEntry,
  77. AdminEventLogPage,
  78. AdminMember,
  79. AdminMemberDetail,
  80. AdminRole,
  81. AdminRoleDetail,
  82. AdminRoomLayoutGroup,
  83. AdminRoomLayoutItemKind,
  84. AdminSystemInfoSnapshot,
  85. Asset,
  86. AssetUpload,
  87. CustomUserStatus,
  88. DirectoryMember,
  89. ExternalIdentityAccount,
  90. ExternalIdentityProvider,
  91. FollowedThread,
  92. FollowedThreadsPage,
  93. ImageTransformOptions,
  94. LinkedExternalIdentity,
  95. LinkPreview,
  96. Message,
  97. Notification,
  98. NotificationLevel,
  99. NotificationPreference,
  100. NotificationsPage,
  101. Page,
  102. PendingExternalIdentity,
  103. PresenceStatus,
  104. Role,
  105. RolePermissionDecisions,
  106. RolePermissionMatrix,
  107. Room,
  108. RoomBan,
  109. RoomDirectoryScope,
  110. RoomGroup,
  111. RoomWithViewerState,
  112. ServerConfig,
  113. ServerLogin,
  114. ServerProfile,
  115. ServerRuntimeConfig,
  116. SidebarLink,
  117. TimeFormat,
  118. TimelinePage,
  119. User,
  120. UserPermissionDecisions,
  121. UserPermissionMatrix,
  122. UserSettings,
  123. ViewerSnapshot,
  124. ViewerUser,
  125. format_datetime,
  126. parse_datetime,
  127. )
  128. RES = TypeVar("RES")
  129. def _page_pb(
  130. limit: int | None, offset: int | None
  131. ) -> pagination_pb2.PageRequest | None:
  132. if limit is None and offset is None:
  133. return None
  134. return pagination_pb2.PageRequest(limit=limit or 0, offset=offset or 0)
  135. def _thumbnail_pb(
  136. opts: ImageTransformOptions | None,
  137. ) -> common_pb2.ImageTransformOptions | None:
  138. if opts is None:
  139. return None
  140. return common_pb2.ImageTransformOptions(
  141. width=opts.width, height=opts.height, fit=opts.fit.value
  142. )
  143. def _timestamp_pb(value: datetime | None) -> Any:
  144. from google.protobuf import timestamp_pb2
  145. if value is None:
  146. return None
  147. ts = timestamp_pb2.Timestamp()
  148. ts.FromJsonString(format_datetime(value))
  149. return ts
  150. class ChattoClient:
  151. """Async client for the Chatto Connect API.
  152. Usage::
  153. async with await ChattoClient.login("user", "pass") as client:
  154. viewer = await client.get_viewer()
  155. rooms = await client.list_rooms()
  156. # Or with an existing token:
  157. async with ChattoClient(token="cht_...") as client:
  158. ...
  159. """
  160. DEFAULT_BASE_URL = "https://chat.chatto.run"
  161. def __init__(
  162. self,
  163. token: str | None = None,
  164. *,
  165. base_url: str = DEFAULT_BASE_URL,
  166. session_cookie: str | None = None,
  167. service_clients: ServiceClients | None = None,
  168. ) -> None:
  169. self._base_url = base_url.rstrip("/")
  170. self._token = token
  171. self._session_cookie = session_cookie
  172. self._svc = service_clients or build_service_clients(self._base_url)
  173. self._owns_clients = service_clients is None
  174. @classmethod
  175. async def login(
  176. cls,
  177. login: str,
  178. password: str,
  179. *,
  180. base_url: str = DEFAULT_BASE_URL,
  181. ) -> ChattoClient:
  182. """Authenticate with username and password, returning a connected client.
  183. Uses Chatto's ``/auth/login`` HTTP endpoint (which is still exposed
  184. alongside the Connect API) and captures both the returned bearer token
  185. and any ``chatto_session`` cookie.
  186. """
  187. base = base_url.rstrip("/")
  188. async with httpx.AsyncClient() as http:
  189. resp = await http.post(
  190. f"{base}/auth/login",
  191. json={"login": login, "password": password},
  192. )
  193. if resp.status_code == 401:
  194. raise ChattoAuthError("Invalid credentials")
  195. resp.raise_for_status()
  196. body = resp.json()
  197. token = body.get("token")
  198. session_cookie = None
  199. if "set-cookie" in resp.headers:
  200. for cookie_header in resp.headers.get_list("set-cookie"):
  201. if cookie_header.startswith("chatto_session="):
  202. session_cookie = cookie_header.split(";")[0].split("=", 1)[1]
  203. break
  204. return cls(token=token, base_url=base_url, session_cookie=session_cookie)
  205. async def __aenter__(self) -> ChattoClient:
  206. return self
  207. async def __aexit__(self, *exc: Any) -> None:
  208. await self.close()
  209. async def close(self) -> None:
  210. if self._owns_clients:
  211. await self._svc.close()
  212. # --- Transport ------------------------------------------------------
  213. @property
  214. def base_url(self) -> str:
  215. return self._base_url
  216. @property
  217. def token(self) -> str | None:
  218. return self._token
  219. @property
  220. def session_cookie(self) -> str | None:
  221. return self._session_cookie
  222. @property
  223. def services(self) -> ServiceClients:
  224. """Direct access to the underlying ConnectRPC service clients.
  225. Useful when a caller wants to reach an RPC that this class doesn't
  226. expose yet, or wants full protobuf messages instead of the
  227. dataclass views returned by the high-level helpers.
  228. """
  229. return self._svc
  230. def _headers(self) -> dict[str, str]:
  231. headers: dict[str, str] = {}
  232. if self._token:
  233. headers["Authorization"] = f"Bearer {self._token}"
  234. if self._session_cookie:
  235. headers["Cookie"] = f"chatto_session={self._session_cookie}"
  236. return headers
  237. async def _rpc(self, coro: Awaitable[RES]) -> RES:
  238. """Await a ConnectRPC coroutine, translating errors."""
  239. try:
  240. return await coro
  241. except ConnectError as exc:
  242. raise translate_connect_error(exc) from exc
  243. # --- Server discovery ----------------------------------------------
  244. async def get_server(self) -> tuple[ServerProfile, ServerLogin]:
  245. """Public server profile and login options. Does not require auth."""
  246. resp = await self._rpc(
  247. self._svc.server_discovery.get_server(
  248. discovery_server_pb2.GetServerRequest(),
  249. headers=self._headers(),
  250. )
  251. )
  252. return (
  253. ServerProfile.parse(pb_to_dict(resp.profile)),
  254. ServerLogin.parse(pb_to_dict(resp.login)),
  255. )
  256. async def get_motd(self) -> str | None:
  257. resp = await self._rpc(
  258. self._svc.server.get_motd(
  259. server_state_pb2.GetMotdRequest(), headers=self._headers()
  260. )
  261. )
  262. if not resp.HasField("motd"):
  263. return None
  264. return resp.motd
  265. async def get_runtime_config(self) -> ServerRuntimeConfig:
  266. resp = await self._rpc(
  267. self._svc.server.get_runtime_config(
  268. server_state_pb2.GetRuntimeConfigRequest(),
  269. headers=self._headers(),
  270. )
  271. )
  272. return ServerRuntimeConfig.parse(pb_to_dict(resp.runtime))
  273. # --- Viewer ---------------------------------------------------------
  274. async def get_viewer(self) -> ViewerSnapshot:
  275. """Full authenticated viewer snapshot.
  276. Returns a :class:`~chattolib.types.ViewerSnapshot` containing the
  277. current user, capabilities, notification preferences, permission
  278. decisions, and viewer state. Callers that only need the current user's
  279. public profile should use :meth:`me` for a lightweight result.
  280. """
  281. resp = await self._rpc(
  282. self._svc.viewer.get_viewer(
  283. viewer_pb2.GetViewerRequest(), headers=self._headers()
  284. )
  285. )
  286. return ViewerSnapshot.parse(pb_to_dict(resp))
  287. async def viewer_user(self) -> ViewerUser | None:
  288. snapshot = await self.get_viewer()
  289. return snapshot.user
  290. async def me(self) -> User:
  291. """Return the authenticated user's public profile."""
  292. viewer = await self.viewer_user()
  293. if viewer is None or viewer.profile is None:
  294. raise ChattoAuthError("No authenticated viewer")
  295. return viewer.profile
  296. # --- MyAccount -----------------------------------------------------
  297. async def update_profile(
  298. self,
  299. *,
  300. display_name: str | None = None,
  301. login: str | None = None,
  302. ) -> User:
  303. req = account_pb2.UpdateProfileRequest()
  304. if display_name is not None:
  305. req.display_name = display_name
  306. if login is not None:
  307. req.login = login
  308. resp = await self._rpc(
  309. self._svc.account.update_profile(req, headers=self._headers())
  310. )
  311. user = User.parse(pb_to_dict(resp.user))
  312. assert user is not None
  313. return user
  314. async def upload_avatar(
  315. self,
  316. file_path: str | Path,
  317. *,
  318. content_type: str = "image/png",
  319. ) -> User:
  320. p = Path(file_path)
  321. req = account_pb2.UploadAvatarRequest(
  322. image=common_pb2.ImageUpload(
  323. image=p.read_bytes(),
  324. filename=p.name,
  325. content_type=content_type,
  326. )
  327. )
  328. resp = await self._rpc(
  329. self._svc.account.upload_avatar(req, headers=self._headers())
  330. )
  331. user = User.parse(pb_to_dict(resp.user))
  332. assert user is not None
  333. return user
  334. async def delete_avatar(self) -> User:
  335. resp = await self._rpc(
  336. self._svc.account.delete_avatar(
  337. account_pb2.DeleteAvatarRequest(), headers=self._headers()
  338. )
  339. )
  340. user = User.parse(pb_to_dict(resp.user))
  341. assert user is not None
  342. return user
  343. async def update_password(self, new_password: str, current_password: str = "") -> User:
  344. resp = await self._rpc(
  345. self._svc.account.update_password(
  346. account_pb2.UpdatePasswordRequest(
  347. password=new_password, current_password=current_password
  348. ),
  349. headers=self._headers(),
  350. )
  351. )
  352. user = User.parse(pb_to_dict(resp.user))
  353. assert user is not None
  354. return user
  355. async def update_settings(
  356. self,
  357. *,
  358. timezone: str | None = None,
  359. time_format: TimeFormat | None = None,
  360. ) -> UserSettings:
  361. req = account_pb2.UpdateSettingsRequest()
  362. if timezone is not None:
  363. req.timezone = timezone
  364. if time_format is not None:
  365. req.time_format = time_format.value
  366. resp = await self._rpc(
  367. self._svc.account.update_settings(req, headers=self._headers())
  368. )
  369. return UserSettings.parse(pb_to_dict(resp.settings))
  370. async def update_presence(
  371. self,
  372. status: PresenceStatus,
  373. *,
  374. user_selected: bool = True,
  375. ) -> PresenceStatus:
  376. if status in (PresenceStatus.UNSPECIFIED, PresenceStatus.OFFLINE):
  377. raise ValueError(
  378. "UNSPECIFIED and OFFLINE cannot be set as presence status; "
  379. "stop refreshing to go offline"
  380. )
  381. req = presence_pb2.UpdatePresenceRequest(
  382. status=status.value, user_selected=user_selected
  383. )
  384. resp = await self._rpc(
  385. self._svc.account.update_presence(req, headers=self._headers())
  386. )
  387. name = presence_pb2.PresenceStatus.Name(resp.status)
  388. return PresenceStatus(name)
  389. async def update_custom_status(
  390. self,
  391. emoji: str,
  392. text: str,
  393. *,
  394. expires_at: datetime | None = None,
  395. ) -> CustomUserStatus | None:
  396. req = user_status_pb2.UpdateCustomStatusRequest(emoji=emoji, text=text)
  397. if expires_at is not None:
  398. req.expires_at.CopyFrom(_timestamp_pb(expires_at))
  399. resp = await self._rpc(
  400. self._svc.account.update_custom_status(req, headers=self._headers())
  401. )
  402. return CustomUserStatus.parse(pb_to_dict(resp).get("status"))
  403. async def delete_custom_status(self) -> CustomUserStatus | None:
  404. resp = await self._rpc(
  405. self._svc.account.delete_custom_status(
  406. user_status_pb2.DeleteCustomStatusRequest(), headers=self._headers()
  407. )
  408. )
  409. return CustomUserStatus.parse(pb_to_dict(resp).get("status"))
  410. async def request_account_deletion(self) -> str:
  411. resp = await self._rpc(
  412. self._svc.account.request_account_deletion(
  413. account_pb2.RequestAccountDeletionRequest(),
  414. headers=self._headers(),
  415. )
  416. )
  417. return resp.confirmation_token
  418. async def delete_my_account(self, confirmation_token: str) -> bool:
  419. resp = await self._rpc(
  420. self._svc.account.delete_my_account(
  421. account_pb2.DeleteMyAccountRequest(
  422. confirmation_token=confirmation_token
  423. ),
  424. headers=self._headers(),
  425. )
  426. )
  427. return resp.deleted
  428. # --- Users ----------------------------------------------------------
  429. async def list_users(
  430. self,
  431. *,
  432. search: str = "",
  433. limit: int | None = None,
  434. offset: int | None = None,
  435. ) -> tuple[list[DirectoryMember], Page]:
  436. req = member_directory_pb2.ListUsersRequest(search=search)
  437. page = _page_pb(limit, offset)
  438. if page is not None:
  439. req.page.CopyFrom(page)
  440. resp = await self._rpc(
  441. self._svc.users.list_users(req, headers=self._headers())
  442. )
  443. data = pb_to_dict(resp)
  444. users = [
  445. u
  446. for u in (DirectoryMember.parse(row) for row in data.get("users") or [])
  447. if u is not None
  448. ]
  449. return users, Page.parse(data.get("page"))
  450. async def get_user(
  451. self, *, user_id: str | None = None, login: str | None = None
  452. ) -> DirectoryMember | None:
  453. if bool(user_id) == bool(login):
  454. raise ValueError("get_user requires exactly one of user_id or login")
  455. req = member_directory_pb2.GetUserRequest()
  456. if user_id:
  457. req.user_id = user_id
  458. else:
  459. assert login is not None
  460. req.login = login
  461. resp = await self._rpc(
  462. self._svc.users.get_user(req, headers=self._headers())
  463. )
  464. return DirectoryMember.parse(pb_to_dict(resp.user))
  465. async def batch_get_users(self, user_ids: list[str]) -> list[DirectoryMember]:
  466. resp = await self._rpc(
  467. self._svc.users.batch_get_users(
  468. member_directory_pb2.BatchGetUsersRequest(user_ids=user_ids),
  469. headers=self._headers(),
  470. )
  471. )
  472. data = pb_to_dict(resp)
  473. return [
  474. u
  475. for u in (DirectoryMember.parse(row) for row in data.get("users") or [])
  476. if u is not None
  477. ]
  478. # --- Roles (public) -----------------------------------------------
  479. async def list_roles(self) -> list[Role]:
  480. resp = await self._rpc(
  481. self._svc.roles.list_roles(
  482. roles_pb2.ListRolesRequest(), headers=self._headers()
  483. )
  484. )
  485. data = pb_to_dict(resp)
  486. return [
  487. r for r in (Role.parse(row) for row in data.get("roles") or []) if r is not None
  488. ]
  489. async def get_role(self, name: str) -> Role | None:
  490. resp = await self._rpc(
  491. self._svc.roles.get_role(
  492. roles_pb2.GetRoleRequest(name=name), headers=self._headers()
  493. )
  494. )
  495. return Role.parse(pb_to_dict(resp.role))
  496. async def batch_get_roles(self, names: list[str]) -> list[Role]:
  497. resp = await self._rpc(
  498. self._svc.roles.batch_get_roles(
  499. roles_pb2.BatchGetRolesRequest(names=names), headers=self._headers()
  500. )
  501. )
  502. data = pb_to_dict(resp)
  503. return [
  504. r for r in (Role.parse(row) for row in data.get("roles") or []) if r is not None
  505. ]
  506. # --- Room directory ------------------------------------------------
  507. async def list_rooms(
  508. self, scope: RoomDirectoryScope = RoomDirectoryScope.ALL
  509. ) -> list[RoomWithViewerState]:
  510. resp = await self._rpc(
  511. self._svc.room_directory.list_rooms(
  512. room_directory_pb2.ListRoomsRequest(scope=scope.value),
  513. headers=self._headers(),
  514. )
  515. )
  516. data = pb_to_dict(resp)
  517. return [
  518. r
  519. for r in (RoomWithViewerState.parse(row) for row in data.get("rooms") or [])
  520. if r is not None
  521. ]
  522. async def list_room_groups(self) -> list[RoomGroup]:
  523. resp = await self._rpc(
  524. self._svc.room_directory.list_room_groups(
  525. room_directory_pb2.ListRoomGroupsRequest(),
  526. headers=self._headers(),
  527. )
  528. )
  529. data = pb_to_dict(resp)
  530. return [
  531. g
  532. for g in (RoomGroup.parse(row) for row in data.get("groups") or [])
  533. if g is not None
  534. ]
  535. async def get_room_group(self, group_id: str) -> RoomGroup | None:
  536. resp = await self._rpc(
  537. self._svc.room_directory.get_room_group(
  538. room_directory_pb2.GetRoomGroupRequest(group_id=group_id),
  539. headers=self._headers(),
  540. )
  541. )
  542. return RoomGroup.parse(pb_to_dict(resp.group))
  543. async def batch_get_room_groups(self, group_ids: list[str]) -> list[RoomGroup]:
  544. resp = await self._rpc(
  545. self._svc.room_directory.batch_get_room_groups(
  546. room_directory_pb2.BatchGetRoomGroupsRequest(group_ids=group_ids),
  547. headers=self._headers(),
  548. )
  549. )
  550. data = pb_to_dict(resp)
  551. return [
  552. g
  553. for g in (RoomGroup.parse(row) for row in data.get("groups") or [])
  554. if g is not None
  555. ]
  556. async def get_room(self, room_id: str) -> RoomWithViewerState | None:
  557. resp = await self._rpc(
  558. self._svc.room_directory.get_room(
  559. room_directory_pb2.GetRoomRequest(room_id=room_id),
  560. headers=self._headers(),
  561. )
  562. )
  563. return RoomWithViewerState.parse(pb_to_dict(resp.room))
  564. async def batch_get_rooms(self, room_ids: list[str]) -> list[RoomWithViewerState]:
  565. resp = await self._rpc(
  566. self._svc.room_directory.batch_get_rooms(
  567. room_directory_pb2.BatchGetRoomsRequest(room_ids=room_ids),
  568. headers=self._headers(),
  569. )
  570. )
  571. data = pb_to_dict(resp)
  572. return [
  573. r
  574. for r in (RoomWithViewerState.parse(row) for row in data.get("rooms") or [])
  575. if r is not None
  576. ]
  577. # --- Room lifecycle & membership -----------------------------------
  578. async def create_room(
  579. self,
  580. name: str,
  581. group_id: str,
  582. *,
  583. description: str = "",
  584. universal: bool = False,
  585. ) -> Room:
  586. resp = await self._rpc(
  587. self._svc.rooms.create_room(
  588. rooms_pb2.CreateRoomRequest(
  589. name=name,
  590. group_id=group_id,
  591. description=description,
  592. universal=universal,
  593. ),
  594. headers=self._headers(),
  595. )
  596. )
  597. room = Room.parse(pb_to_dict(resp.room))
  598. assert room is not None
  599. return room
  600. async def update_room(
  601. self,
  602. room_id: str,
  603. *,
  604. name: str | None = None,
  605. description: str | None = None,
  606. universal: bool | None = None,
  607. ) -> Room:
  608. req = rooms_pb2.UpdateRoomRequest(room_id=room_id)
  609. if name is not None:
  610. req.name = name
  611. if description is not None:
  612. req.description = description
  613. if universal is not None:
  614. req.universal = universal
  615. resp = await self._rpc(
  616. self._svc.rooms.update_room(req, headers=self._headers())
  617. )
  618. room = Room.parse(pb_to_dict(resp.room))
  619. assert room is not None
  620. return room
  621. async def archive_room(self, room_id: str) -> Room:
  622. resp = await self._rpc(
  623. self._svc.rooms.archive_room(
  624. rooms_pb2.ArchiveRoomRequest(room_id=room_id),
  625. headers=self._headers(),
  626. )
  627. )
  628. room = Room.parse(pb_to_dict(resp.room))
  629. assert room is not None
  630. return room
  631. async def unarchive_room(self, room_id: str) -> Room:
  632. resp = await self._rpc(
  633. self._svc.rooms.unarchive_room(
  634. rooms_pb2.UnarchiveRoomRequest(room_id=room_id),
  635. headers=self._headers(),
  636. )
  637. )
  638. room = Room.parse(pb_to_dict(resp.room))
  639. assert room is not None
  640. return room
  641. async def join_room(self, room_id: str) -> Room:
  642. resp = await self._rpc(
  643. self._svc.rooms.join_room(
  644. rooms_pb2.JoinRoomRequest(room_id=room_id),
  645. headers=self._headers(),
  646. )
  647. )
  648. room = Room.parse(pb_to_dict(resp.room))
  649. assert room is not None
  650. return room
  651. async def join_room_group(self, group_id: str) -> list[str]:
  652. resp = await self._rpc(
  653. self._svc.rooms.join_room_group(
  654. rooms_pb2.JoinRoomGroupRequest(group_id=group_id),
  655. headers=self._headers(),
  656. )
  657. )
  658. return list(resp.joined_room_ids)
  659. async def start_dm(self, participant_ids: list[str]) -> Room:
  660. resp = await self._rpc(
  661. self._svc.rooms.start_dm(
  662. rooms_pb2.StartDMRequest(participant_ids=participant_ids),
  663. headers=self._headers(),
  664. )
  665. )
  666. room = Room.parse(pb_to_dict(resp.room))
  667. assert room is not None
  668. return room
  669. async def leave_room(self, room_id: str) -> bool:
  670. resp = await self._rpc(
  671. self._svc.rooms.leave_room(
  672. rooms_pb2.LeaveRoomRequest(room_id=room_id),
  673. headers=self._headers(),
  674. )
  675. )
  676. return resp.left
  677. async def add_member(self, room_id: str, user_id: str) -> DirectoryMember | None:
  678. resp = await self._rpc(
  679. self._svc.rooms.add_member(
  680. rooms_pb2.AddMemberRequest(room_id=room_id, user_id=user_id),
  681. headers=self._headers(),
  682. )
  683. )
  684. return DirectoryMember.parse(pb_to_dict(resp.member))
  685. async def remove_member(self, room_id: str, user_id: str) -> bool:
  686. resp = await self._rpc(
  687. self._svc.rooms.remove_member(
  688. rooms_pb2.RemoveMemberRequest(room_id=room_id, user_id=user_id),
  689. headers=self._headers(),
  690. )
  691. )
  692. return resp.removed
  693. async def list_room_members(
  694. self,
  695. room_id: str,
  696. *,
  697. search: str = "",
  698. limit: int | None = None,
  699. offset: int | None = None,
  700. ) -> tuple[list[DirectoryMember], Page]:
  701. req = member_directory_pb2.ListRoomMembersRequest(room_id=room_id, search=search)
  702. page = _page_pb(limit, offset)
  703. if page is not None:
  704. req.page.CopyFrom(page)
  705. resp = await self._rpc(
  706. self._svc.rooms.list_members(req, headers=self._headers())
  707. )
  708. data = pb_to_dict(resp)
  709. members = [
  710. m
  711. for m in (DirectoryMember.parse(row) for row in data.get("members") or [])
  712. if m is not None
  713. ]
  714. return members, Page.parse(data.get("page"))
  715. async def get_room_member(self, room_id: str, user_id: str) -> DirectoryMember | None:
  716. resp = await self._rpc(
  717. self._svc.rooms.get_member(
  718. member_directory_pb2.GetRoomMemberRequest(
  719. room_id=room_id, user_id=user_id
  720. ),
  721. headers=self._headers(),
  722. )
  723. )
  724. return DirectoryMember.parse(pb_to_dict(resp.member))
  725. async def batch_get_room_members(
  726. self, room_id: str, user_ids: list[str]
  727. ) -> list[DirectoryMember]:
  728. resp = await self._rpc(
  729. self._svc.rooms.batch_get_members(
  730. member_directory_pb2.BatchGetRoomMembersRequest(
  731. room_id=room_id, user_ids=user_ids
  732. ),
  733. headers=self._headers(),
  734. )
  735. )
  736. data = pb_to_dict(resp)
  737. return [
  738. m
  739. for m in (DirectoryMember.parse(row) for row in data.get("members") or [])
  740. if m is not None
  741. ]
  742. async def ban_member(
  743. self,
  744. room_id: str,
  745. user_id: str,
  746. reason: str,
  747. *,
  748. expires_at: datetime | None = None,
  749. ) -> bool:
  750. req = rooms_pb2.BanMemberRequest(
  751. room_id=room_id, user_id=user_id, reason=reason
  752. )
  753. if expires_at is not None:
  754. req.expires_at.CopyFrom(_timestamp_pb(expires_at))
  755. resp = await self._rpc(
  756. self._svc.rooms.ban_member(req, headers=self._headers())
  757. )
  758. return resp.banned
  759. async def unban_member(self, room_id: str, user_id: str, reason: str) -> bool:
  760. resp = await self._rpc(
  761. self._svc.rooms.unban_member(
  762. rooms_pb2.UnbanMemberRequest(
  763. room_id=room_id, user_id=user_id, reason=reason
  764. ),
  765. headers=self._headers(),
  766. )
  767. )
  768. return resp.unbanned
  769. async def list_bans(
  770. self,
  771. *,
  772. room_id: str = "",
  773. limit: int | None = None,
  774. offset: int | None = None,
  775. ) -> tuple[list[RoomBan], Page]:
  776. req = rooms_pb2.ListBansRequest(room_id=room_id)
  777. page = _page_pb(limit, offset)
  778. if page is not None:
  779. req.page.CopyFrom(page)
  780. resp = await self._rpc(
  781. self._svc.rooms.list_bans(req, headers=self._headers())
  782. )
  783. data = pb_to_dict(resp)
  784. bans = [
  785. b
  786. for b in (RoomBan.parse(row) for row in data.get("bans") or [])
  787. if b is not None
  788. ]
  789. return bans, Page.parse(data.get("page"))
  790. async def update_typing_indicator(
  791. self, room_id: str, *, thread_root_event_id: str = ""
  792. ) -> bool:
  793. resp = await self._rpc(
  794. self._svc.rooms.update_typing_indicator(
  795. rooms_pb2.UpdateTypingIndicatorRequest(
  796. room_id=room_id, thread_root_event_id=thread_root_event_id
  797. ),
  798. headers=self._headers(),
  799. )
  800. )
  801. return resp.updated
  802. # --- Room timeline / read state -----------------------------------
  803. async def get_room_events(
  804. self,
  805. room_id: str,
  806. *,
  807. limit: int | None = None,
  808. before: str | None = None,
  809. after: str | None = None,
  810. ) -> TimelinePage:
  811. req = room_timeline_pb2.GetRoomEventsRequest(room_id=room_id)
  812. if limit is not None:
  813. req.limit = limit
  814. if before is not None:
  815. req.before = before
  816. elif after is not None:
  817. req.after = after
  818. resp = await self._rpc(
  819. self._svc.rooms.get_room_events(req, headers=self._headers())
  820. )
  821. return TimelinePage.parse(pb_to_dict(resp.page))
  822. async def get_room_events_around(
  823. self,
  824. room_id: str,
  825. event_id: str,
  826. *,
  827. limit: int | None = None,
  828. ) -> tuple[TimelinePage, int]:
  829. req = room_timeline_pb2.GetRoomEventsAroundRequest(
  830. room_id=room_id, event_id=event_id
  831. )
  832. if limit is not None:
  833. req.limit = limit
  834. resp = await self._rpc(
  835. self._svc.rooms.get_room_events_around(req, headers=self._headers())
  836. )
  837. return TimelinePage.parse(pb_to_dict(resp.page)), resp.target_index
  838. async def mark_room_as_read(
  839. self, room_id: str, up_to_event_id: str = ""
  840. ) -> tuple[datetime | None, datetime | None]:
  841. resp = await self._rpc(
  842. self._svc.rooms.mark_room_as_read(
  843. read_state_pb2.MarkRoomAsReadRequest(
  844. room_id=room_id, up_to_event_id=up_to_event_id
  845. ),
  846. headers=self._headers(),
  847. )
  848. )
  849. d = pb_to_dict(resp)
  850. return (
  851. parse_datetime(d.get("lastReadAt")),
  852. parse_datetime(d.get("previousLastReadAt")),
  853. )
  854. async def list_room_attachments(
  855. self,
  856. room_id: str,
  857. *,
  858. thumbnail: ImageTransformOptions | None = None,
  859. limit: int | None = None,
  860. offset: int | None = None,
  861. ) -> tuple[list[dict[str, Any]], Page]:
  862. req = rooms_pb2.ListRoomAttachmentsRequest(room_id=room_id)
  863. thumb = _thumbnail_pb(thumbnail)
  864. if thumb is not None:
  865. req.thumbnail.CopyFrom(thumb)
  866. page = _page_pb(limit, offset)
  867. if page is not None:
  868. req.page.CopyFrom(page)
  869. resp = await self._rpc(
  870. self._svc.rooms.list_room_attachments(req, headers=self._headers())
  871. )
  872. data = pb_to_dict(resp)
  873. return list(data.get("attachments") or []), Page.parse(data.get("page"))
  874. # --- Messages -------------------------------------------------------
  875. async def fetch_link_preview(self, url: str) -> tuple[LinkPreview | None, str]:
  876. resp = await self._rpc(
  877. self._svc.messages.fetch_link_preview(
  878. link_previews_pb2.FetchLinkPreviewRequest(url=url),
  879. headers=self._headers(),
  880. )
  881. )
  882. return LinkPreview.parse(pb_to_dict(resp.preview)), resp.preview_token
  883. async def post_message(
  884. self,
  885. room_id: str,
  886. body: str = "",
  887. *,
  888. attachment_asset_ids: list[str] | None = None,
  889. thread_root_event_id: str = "",
  890. in_reply_to: str = "",
  891. also_send_to_channel: bool = False,
  892. link_preview_token: str = "",
  893. ) -> Message:
  894. req = messages_pb2.CreateMessageRequest(
  895. room_id=room_id,
  896. body=body,
  897. thread_root_event_id=thread_root_event_id,
  898. in_reply_to=in_reply_to,
  899. also_send_to_channel=also_send_to_channel,
  900. link_preview_token=link_preview_token,
  901. )
  902. if attachment_asset_ids:
  903. req.attachment_asset_ids.extend(attachment_asset_ids)
  904. resp = await self._rpc(
  905. self._svc.messages.create_message(req, headers=self._headers())
  906. )
  907. message = Message.parse(pb_to_dict(resp.message))
  908. assert message is not None
  909. return message
  910. async def update_message(
  911. self,
  912. room_id: str,
  913. event_id: str,
  914. *,
  915. body: str | None = None,
  916. also_send_to_channel: bool | None = None,
  917. ) -> Message:
  918. req = messages_pb2.UpdateMessageRequest(room_id=room_id, event_id=event_id)
  919. if body is not None:
  920. req.body = body
  921. if also_send_to_channel is not None:
  922. req.also_send_to_channel = also_send_to_channel
  923. resp = await self._rpc(
  924. self._svc.messages.update_message(req, headers=self._headers())
  925. )
  926. message = Message.parse(pb_to_dict(resp.message))
  927. assert message is not None
  928. return message
  929. async def delete_message(self, room_id: str, event_id: str) -> bool:
  930. resp = await self._rpc(
  931. self._svc.messages.delete_message(
  932. messages_pb2.DeleteMessageRequest(room_id=room_id, event_id=event_id),
  933. headers=self._headers(),
  934. )
  935. )
  936. return resp.deleted
  937. async def delete_attachment(
  938. self, room_id: str, event_id: str, attachment_id: str
  939. ) -> bool:
  940. resp = await self._rpc(
  941. self._svc.messages.delete_attachment(
  942. messages_pb2.DeleteAttachmentRequest(
  943. room_id=room_id, event_id=event_id, attachment_id=attachment_id
  944. ),
  945. headers=self._headers(),
  946. )
  947. )
  948. return resp.deleted
  949. async def delete_link_preview(self, room_id: str, event_id: str, url: str) -> bool:
  950. resp = await self._rpc(
  951. self._svc.messages.delete_link_preview(
  952. messages_pb2.DeleteLinkPreviewRequest(
  953. room_id=room_id, event_id=event_id, url=url
  954. ),
  955. headers=self._headers(),
  956. )
  957. )
  958. return resp.deleted
  959. async def get_message(self, room_id: str, event_id: str) -> Message | None:
  960. resp = await self._rpc(
  961. self._svc.messages.get_message(
  962. messages_pb2.GetMessageRequest(room_id=room_id, event_id=event_id),
  963. headers=self._headers(),
  964. )
  965. )
  966. return Message.parse(pb_to_dict(resp.message))
  967. async def batch_get_messages(
  968. self, room_id: str, event_ids: list[str]
  969. ) -> list[Message]:
  970. resp = await self._rpc(
  971. self._svc.messages.batch_get_messages(
  972. messages_pb2.BatchGetMessagesRequest(
  973. room_id=room_id, event_ids=event_ids
  974. ),
  975. headers=self._headers(),
  976. )
  977. )
  978. data = pb_to_dict(resp)
  979. return [
  980. m
  981. for m in (Message.parse(row) for row in data.get("messages") or [])
  982. if m is not None
  983. ]
  984. async def add_reaction(
  985. self, room_id: str, message_event_id: str, emoji: str
  986. ) -> bool:
  987. resp = await self._rpc(
  988. self._svc.messages.add_reaction(
  989. reactions_pb2.AddReactionRequest(
  990. room_id=room_id, message_event_id=message_event_id, emoji=emoji
  991. ),
  992. headers=self._headers(),
  993. )
  994. )
  995. return resp.added
  996. async def remove_reaction(
  997. self, room_id: str, message_event_id: str, emoji: str
  998. ) -> bool:
  999. resp = await self._rpc(
  1000. self._svc.messages.remove_reaction(
  1001. reactions_pb2.RemoveReactionRequest(
  1002. room_id=room_id, message_event_id=message_event_id, emoji=emoji
  1003. ),
  1004. headers=self._headers(),
  1005. )
  1006. )
  1007. return resp.removed
  1008. # --- Threads --------------------------------------------------------
  1009. async def follow_thread(self, room_id: str, thread_root_event_id: str) -> bool:
  1010. resp = await self._rpc(
  1011. self._svc.threads.follow_thread(
  1012. threads_pb2.FollowThreadRequest(
  1013. room_id=room_id, thread_root_event_id=thread_root_event_id
  1014. ),
  1015. headers=self._headers(),
  1016. )
  1017. )
  1018. return resp.following
  1019. async def unfollow_thread(self, room_id: str, thread_root_event_id: str) -> bool:
  1020. resp = await self._rpc(
  1021. self._svc.threads.unfollow_thread(
  1022. threads_pb2.UnfollowThreadRequest(
  1023. room_id=room_id, thread_root_event_id=thread_root_event_id
  1024. ),
  1025. headers=self._headers(),
  1026. )
  1027. )
  1028. return resp.following
  1029. async def list_followed_threads(
  1030. self, *, limit: int | None = None, offset: int | None = None
  1031. ) -> FollowedThreadsPage:
  1032. req = threads_pb2.ListFollowedThreadsRequest()
  1033. page = _page_pb(limit, offset)
  1034. if page is not None:
  1035. req.page.CopyFrom(page)
  1036. resp = await self._rpc(
  1037. self._svc.threads.list_followed_threads(req, headers=self._headers())
  1038. )
  1039. data = pb_to_dict(resp)
  1040. threads = [FollowedThread.parse(t) for t in data.get("threads") or []]
  1041. users: dict[str, User] = {}
  1042. includes = data.get("includes") or {}
  1043. for uid, user_data in (includes.get("users") or {}).items():
  1044. parsed = User.parse(user_data)
  1045. if parsed is not None:
  1046. users[uid] = parsed
  1047. return FollowedThreadsPage(
  1048. threads=threads, page=Page.parse(data.get("page")), users_by_id=users
  1049. )
  1050. async def get_thread_events(
  1051. self,
  1052. room_id: str,
  1053. thread_root_event_id: str,
  1054. *,
  1055. limit: int | None = None,
  1056. before: str | None = None,
  1057. after: str | None = None,
  1058. ) -> TimelinePage:
  1059. req = room_timeline_pb2.GetThreadEventsRequest(
  1060. room_id=room_id, thread_root_event_id=thread_root_event_id
  1061. )
  1062. if limit is not None:
  1063. req.limit = limit
  1064. if before is not None:
  1065. req.before = before
  1066. elif after is not None:
  1067. req.after = after
  1068. resp = await self._rpc(
  1069. self._svc.threads.get_thread_events(req, headers=self._headers())
  1070. )
  1071. return TimelinePage.parse(pb_to_dict(resp.page))
  1072. async def get_thread_events_around(
  1073. self,
  1074. room_id: str,
  1075. thread_root_event_id: str,
  1076. event_id: str,
  1077. *,
  1078. limit: int | None = None,
  1079. ) -> tuple[TimelinePage, int]:
  1080. req = room_timeline_pb2.GetThreadEventsAroundRequest(
  1081. room_id=room_id,
  1082. thread_root_event_id=thread_root_event_id,
  1083. event_id=event_id,
  1084. )
  1085. if limit is not None:
  1086. req.limit = limit
  1087. resp = await self._rpc(
  1088. self._svc.threads.get_thread_events_around(req, headers=self._headers())
  1089. )
  1090. return TimelinePage.parse(pb_to_dict(resp.page)), resp.target_index
  1091. async def mark_thread_as_read(
  1092. self,
  1093. room_id: str,
  1094. thread_root_event_id: str,
  1095. up_to_event_id: str = "",
  1096. ) -> datetime | None:
  1097. resp = await self._rpc(
  1098. self._svc.threads.mark_thread_as_read(
  1099. read_state_pb2.MarkThreadAsReadRequest(
  1100. room_id=room_id,
  1101. thread_root_event_id=thread_root_event_id,
  1102. up_to_event_id=up_to_event_id,
  1103. ),
  1104. headers=self._headers(),
  1105. )
  1106. )
  1107. d = pb_to_dict(resp)
  1108. return parse_datetime(d.get("previousReadAt"))
  1109. # --- Notifications --------------------------------------------------
  1110. async def list_notifications(
  1111. self, *, limit: int | None = None, offset: int | None = None
  1112. ) -> NotificationsPage:
  1113. req = notifications_pb2.ListNotificationsRequest()
  1114. page = _page_pb(limit, offset)
  1115. if page is not None:
  1116. req.page.CopyFrom(page)
  1117. resp = await self._rpc(
  1118. self._svc.notifications.list_notifications(req, headers=self._headers())
  1119. )
  1120. data = pb_to_dict(resp)
  1121. notifications = [Notification.parse(n) for n in data.get("notifications") or []]
  1122. return NotificationsPage(
  1123. notifications=notifications, page=Page.parse(data.get("page"))
  1124. )
  1125. async def has_notifications(self) -> bool:
  1126. resp = await self._rpc(
  1127. self._svc.notifications.has_notifications(
  1128. notifications_pb2.HasNotificationsRequest(),
  1129. headers=self._headers(),
  1130. )
  1131. )
  1132. return resp.has_notifications
  1133. async def get_notification(self, notification_id: str) -> Notification | None:
  1134. resp = await self._rpc(
  1135. self._svc.notifications.get_notification(
  1136. notifications_pb2.GetNotificationRequest(notification_id=notification_id),
  1137. headers=self._headers(),
  1138. )
  1139. )
  1140. raw = pb_to_dict(resp).get("notification")
  1141. return Notification.parse(raw) if raw else None
  1142. async def batch_get_notifications(
  1143. self, notification_ids: list[str]
  1144. ) -> list[Notification]:
  1145. resp = await self._rpc(
  1146. self._svc.notifications.batch_get_notifications(
  1147. notifications_pb2.BatchGetNotificationsRequest(
  1148. notification_ids=notification_ids
  1149. ),
  1150. headers=self._headers(),
  1151. )
  1152. )
  1153. data = pb_to_dict(resp)
  1154. return [Notification.parse(n) for n in data.get("notifications") or []]
  1155. async def list_room_notifications(
  1156. self,
  1157. room_id: str,
  1158. *,
  1159. limit: int | None = None,
  1160. offset: int | None = None,
  1161. ) -> NotificationsPage:
  1162. req = notifications_pb2.ListRoomNotificationsRequest(room_id=room_id)
  1163. page = _page_pb(limit, offset)
  1164. if page is not None:
  1165. req.page.CopyFrom(page)
  1166. resp = await self._rpc(
  1167. self._svc.notifications.list_room_notifications(
  1168. req, headers=self._headers()
  1169. )
  1170. )
  1171. data = pb_to_dict(resp)
  1172. return NotificationsPage(
  1173. notifications=[Notification.parse(n) for n in data.get("notifications") or []],
  1174. page=Page.parse(data.get("page")),
  1175. )
  1176. async def list_room_notification_counts(self) -> dict[str, int]:
  1177. resp = await self._rpc(
  1178. self._svc.notifications.list_room_notification_counts(
  1179. notifications_pb2.ListRoomNotificationCountsRequest(),
  1180. headers=self._headers(),
  1181. )
  1182. )
  1183. return {row.room_id: row.total_count for row in resp.room_counts}
  1184. async def dismiss_notification(self, notification_id: str) -> bool:
  1185. resp = await self._rpc(
  1186. self._svc.notifications.dismiss_notification(
  1187. notifications_pb2.DismissNotificationRequest(
  1188. notification_id=notification_id
  1189. ),
  1190. headers=self._headers(),
  1191. )
  1192. )
  1193. return resp.dismissed
  1194. async def dismiss_all_notifications(self) -> int:
  1195. resp = await self._rpc(
  1196. self._svc.notifications.dismiss_all_notifications(
  1197. notifications_pb2.DismissAllNotificationsRequest(),
  1198. headers=self._headers(),
  1199. )
  1200. )
  1201. return resp.dismissed_count
  1202. # --- Notification preferences --------------------------------------
  1203. async def get_server_notification_preference(self) -> NotificationPreference:
  1204. resp = await self._rpc(
  1205. self._svc.notification_prefs.get_server_notification_preference(
  1206. notification_preferences_pb2.GetServerNotificationPreferenceRequest(),
  1207. headers=self._headers(),
  1208. )
  1209. )
  1210. return NotificationPreference.parse(pb_to_dict(resp.preference))
  1211. async def update_server_notification_preference(
  1212. self, level: NotificationLevel
  1213. ) -> NotificationPreference:
  1214. resp = await self._rpc(
  1215. self._svc.notification_prefs.update_server_notification_preference(
  1216. notification_preferences_pb2.UpdateServerNotificationPreferenceRequest(
  1217. level=level.value
  1218. ),
  1219. headers=self._headers(),
  1220. )
  1221. )
  1222. return NotificationPreference.parse(pb_to_dict(resp.preference))
  1223. async def get_room_notification_preference(
  1224. self, room_id: str
  1225. ) -> NotificationPreference:
  1226. resp = await self._rpc(
  1227. self._svc.notification_prefs.get_room_notification_preference(
  1228. notification_preferences_pb2.GetRoomNotificationPreferenceRequest(
  1229. room_id=room_id
  1230. ),
  1231. headers=self._headers(),
  1232. )
  1233. )
  1234. return NotificationPreference.parse(pb_to_dict(resp.preference))
  1235. async def update_room_notification_preference(
  1236. self, room_id: str, level: NotificationLevel
  1237. ) -> NotificationPreference:
  1238. resp = await self._rpc(
  1239. self._svc.notification_prefs.update_room_notification_preference(
  1240. notification_preferences_pb2.UpdateRoomNotificationPreferenceRequest(
  1241. room_id=room_id, level=level.value
  1242. ),
  1243. headers=self._headers(),
  1244. )
  1245. )
  1246. return NotificationPreference.parse(pb_to_dict(resp.preference))
  1247. # --- Push notifications --------------------------------------------
  1248. async def subscribe_push(
  1249. self,
  1250. endpoint: str,
  1251. p256dh: str,
  1252. auth: str,
  1253. *,
  1254. user_agent: str | None = None,
  1255. ) -> bool:
  1256. req = push_notifications_pb2.SubscribePushRequest(
  1257. endpoint=endpoint, p256dh=p256dh, auth=auth
  1258. )
  1259. if user_agent is not None:
  1260. req.user_agent = user_agent
  1261. resp = await self._rpc(
  1262. self._svc.push.subscribe(req, headers=self._headers())
  1263. )
  1264. return resp.subscribed
  1265. async def unsubscribe_push(self, endpoint: str) -> bool:
  1266. resp = await self._rpc(
  1267. self._svc.push.unsubscribe(
  1268. push_notifications_pb2.UnsubscribePushRequest(endpoint=endpoint),
  1269. headers=self._headers(),
  1270. )
  1271. )
  1272. return resp.unsubscribed
  1273. # --- Assets ---------------------------------------------------------
  1274. async def get_asset(
  1275. self,
  1276. room_id: str,
  1277. asset_id: str,
  1278. *,
  1279. thumbnail: ImageTransformOptions | None = None,
  1280. ) -> Asset | None:
  1281. req = attachments_pb2.GetAssetRequest(room_id=room_id, asset_id=asset_id)
  1282. thumb = _thumbnail_pb(thumbnail)
  1283. if thumb is not None:
  1284. req.thumbnail.CopyFrom(thumb)
  1285. resp = await self._rpc(
  1286. self._svc.assets.get_asset(req, headers=self._headers())
  1287. )
  1288. return Asset.parse(pb_to_dict(resp.asset))
  1289. async def batch_get_assets(
  1290. self,
  1291. room_id: str,
  1292. asset_ids: list[str],
  1293. *,
  1294. thumbnail: ImageTransformOptions | None = None,
  1295. ) -> list[Asset]:
  1296. req = attachments_pb2.BatchGetAssetsRequest(
  1297. room_id=room_id, asset_ids=asset_ids
  1298. )
  1299. thumb = _thumbnail_pb(thumbnail)
  1300. if thumb is not None:
  1301. req.thumbnail.CopyFrom(thumb)
  1302. resp = await self._rpc(
  1303. self._svc.assets.batch_get_assets(req, headers=self._headers())
  1304. )
  1305. data = pb_to_dict(resp)
  1306. return [
  1307. a
  1308. for a in (Asset.parse(row) for row in data.get("assets") or [])
  1309. if a is not None
  1310. ]
  1311. # --- Asset uploads ------------------------------------------------
  1312. async def create_upload(
  1313. self,
  1314. room_id: str,
  1315. filename: str,
  1316. size: int,
  1317. sha256: str,
  1318. *,
  1319. content_type: str = "",
  1320. ) -> AssetUpload:
  1321. resp = await self._rpc(
  1322. self._svc.asset_uploads.create_upload(
  1323. asset_uploads_pb2.CreateUploadRequest(
  1324. room_id=room_id,
  1325. filename=filename,
  1326. content_type=content_type,
  1327. size=size,
  1328. sha256=sha256,
  1329. ),
  1330. headers=self._headers(),
  1331. )
  1332. )
  1333. upload = AssetUpload.parse(pb_to_dict(resp.upload))
  1334. assert upload is not None
  1335. return upload
  1336. async def upload_chunk(
  1337. self, upload_id: str, offset: int, content: bytes, chunk_sha256: str
  1338. ) -> AssetUpload:
  1339. resp = await self._rpc(
  1340. self._svc.asset_uploads.upload_chunk(
  1341. asset_uploads_pb2.UploadChunkRequest(
  1342. upload_id=upload_id,
  1343. offset=offset,
  1344. content=content,
  1345. chunk_sha256=chunk_sha256,
  1346. ),
  1347. headers=self._headers(),
  1348. )
  1349. )
  1350. upload = AssetUpload.parse(pb_to_dict(resp.upload))
  1351. assert upload is not None
  1352. return upload
  1353. async def get_upload(self, upload_id: str) -> AssetUpload:
  1354. resp = await self._rpc(
  1355. self._svc.asset_uploads.get_upload(
  1356. asset_uploads_pb2.GetUploadRequest(upload_id=upload_id),
  1357. headers=self._headers(),
  1358. )
  1359. )
  1360. upload = AssetUpload.parse(pb_to_dict(resp.upload))
  1361. assert upload is not None
  1362. return upload
  1363. async def complete_upload(
  1364. self, upload_id: str
  1365. ) -> tuple[AssetUpload, Asset | None]:
  1366. resp = await self._rpc(
  1367. self._svc.asset_uploads.complete_upload(
  1368. asset_uploads_pb2.CompleteUploadRequest(upload_id=upload_id),
  1369. headers=self._headers(),
  1370. )
  1371. )
  1372. upload = AssetUpload.parse(pb_to_dict(resp.upload))
  1373. assert upload is not None
  1374. return upload, Asset.parse(pb_to_dict(resp.asset))
  1375. async def cancel_upload(self, upload_id: str) -> AssetUpload:
  1376. resp = await self._rpc(
  1377. self._svc.asset_uploads.cancel_upload(
  1378. asset_uploads_pb2.CancelUploadRequest(upload_id=upload_id),
  1379. headers=self._headers(),
  1380. )
  1381. )
  1382. upload = AssetUpload.parse(pb_to_dict(resp.upload))
  1383. assert upload is not None
  1384. return upload
  1385. async def upload_attachment(
  1386. self,
  1387. room_id: str,
  1388. file_path: str | Path,
  1389. *,
  1390. content_type: str = "",
  1391. filename: str | None = None,
  1392. ) -> Asset:
  1393. """Upload a file as a room attachment and return the resulting Asset."""
  1394. path = Path(file_path)
  1395. data = path.read_bytes()
  1396. size = len(data)
  1397. sha = hashlib.sha256(data).hexdigest()
  1398. upload = await self.create_upload(
  1399. room_id, filename or path.name, size, sha, content_type=content_type
  1400. )
  1401. chunk_size = upload.max_chunk_size or 512 * 1024
  1402. offset = upload.committed_offset
  1403. while offset < size:
  1404. end = min(offset + chunk_size, size)
  1405. chunk = data[offset:end]
  1406. chunk_sha = hashlib.sha256(chunk).hexdigest()
  1407. upload = await self.upload_chunk(upload.upload_id, offset, chunk, chunk_sha)
  1408. if upload.committed_offset <= offset:
  1409. raise ChattoError(
  1410. f"upload stalled at offset {offset} (server reported "
  1411. f"committed_offset={upload.committed_offset})"
  1412. )
  1413. offset = upload.committed_offset
  1414. upload, asset = await self.complete_upload(upload.upload_id)
  1415. if asset is None:
  1416. raise ChattoError("upload completed but server returned no asset")
  1417. return asset
  1418. # --- MyAccount external identities --------------------------------
  1419. async def list_external_identities(
  1420. self,
  1421. ) -> tuple[list[ExternalIdentityProvider], list[LinkedExternalIdentity]]:
  1422. resp = await self._rpc(
  1423. self._svc.account.list_external_identities(
  1424. external_identities_pb2.ListExternalIdentitiesRequest(),
  1425. headers=self._headers(),
  1426. )
  1427. )
  1428. data = pb_to_dict(resp)
  1429. providers = [
  1430. ExternalIdentityProvider.parse(p) for p in data.get("providers") or []
  1431. ]
  1432. linked = [
  1433. LinkedExternalIdentity.parse(li)
  1434. for li in data.get("linkedIdentities") or []
  1435. ]
  1436. return providers, linked
  1437. async def start_external_identity_link(
  1438. self,
  1439. provider_id: str,
  1440. *,
  1441. redirect_path: str = "",
  1442. current_password: str = "",
  1443. ) -> str:
  1444. resp = await self._rpc(
  1445. self._svc.account.start_external_identity_link(
  1446. external_identities_pb2.StartExternalIdentityLinkRequest(
  1447. provider_id=provider_id,
  1448. redirect_path=redirect_path,
  1449. current_password=current_password,
  1450. ),
  1451. headers=self._headers(),
  1452. )
  1453. )
  1454. return resp.start_url
  1455. async def disconnect_external_identity(
  1456. self, subject_hash: str, *, current_password: str = ""
  1457. ) -> bool:
  1458. resp = await self._rpc(
  1459. self._svc.account.disconnect_external_identity(
  1460. external_identities_pb2.DisconnectExternalIdentityRequest(
  1461. subject_hash=subject_hash, current_password=current_password
  1462. ),
  1463. headers=self._headers(),
  1464. )
  1465. )
  1466. return resp.disconnected
  1467. # --- ExternalIdentityAuthService (public OAuth handoff) -----------
  1468. async def get_pending_external_identity(
  1469. self, token: str
  1470. ) -> PendingExternalIdentity | None:
  1471. resp = await self._rpc(
  1472. self._svc.external_auth.get_pending_external_identity(
  1473. external_identity_auth_pb2.GetPendingExternalIdentityRequest(
  1474. token=token
  1475. ),
  1476. headers=self._headers(),
  1477. )
  1478. )
  1479. return PendingExternalIdentity.parse(pb_to_dict(resp).get("pending"))
  1480. async def create_external_identity_account(
  1481. self, token: str, login: str
  1482. ) -> ExternalIdentityAccount | None:
  1483. resp = await self._rpc(
  1484. self._svc.external_auth.create_external_identity_account(
  1485. external_identity_auth_pb2.CreateExternalIdentityAccountRequest(
  1486. token=token, login=login
  1487. ),
  1488. headers=self._headers(),
  1489. )
  1490. )
  1491. return ExternalIdentityAccount.parse(pb_to_dict(resp))
  1492. async def confirm_external_identity_link(
  1493. self, token: str
  1494. ) -> LinkedExternalIdentity | None:
  1495. resp = await self._rpc(
  1496. self._svc.external_auth.confirm_external_identity_link(
  1497. external_identity_auth_pb2.ConfirmExternalIdentityLinkRequest(
  1498. token=token
  1499. ),
  1500. headers=self._headers(),
  1501. )
  1502. )
  1503. data = pb_to_dict(resp).get("linkedIdentity")
  1504. return LinkedExternalIdentity.parse(data) if data else None
  1505. async def cancel_external_identity_flow(self, token: str) -> bool:
  1506. resp = await self._rpc(
  1507. self._svc.external_auth.cancel_external_identity_flow(
  1508. external_identity_auth_pb2.CancelExternalIdentityFlowRequest(
  1509. token=token
  1510. ),
  1511. headers=self._headers(),
  1512. )
  1513. )
  1514. return resp.cancelled
  1515. # --- Voice calls ----------------------------------------------------
  1516. async def list_active_calls(self) -> list[ActiveCall]:
  1517. resp = await self._rpc(
  1518. self._svc.voice_calls.list_active_calls(
  1519. voice_calls_pb2.ListActiveCallsRequest(), headers=self._headers()
  1520. )
  1521. )
  1522. data = pb_to_dict(resp)
  1523. return [ActiveCall.parse(c) for c in data.get("calls") or []]
  1524. async def get_active_call(self, room_id: str) -> ActiveCall | None:
  1525. resp = await self._rpc(
  1526. self._svc.voice_calls.get_active_call(
  1527. voice_calls_pb2.GetActiveCallRequest(room_id=room_id),
  1528. headers=self._headers(),
  1529. )
  1530. )
  1531. raw = pb_to_dict(resp).get("call")
  1532. return ActiveCall.parse(raw) if raw else None
  1533. async def batch_get_active_calls(self, room_ids: list[str]) -> list[ActiveCall]:
  1534. resp = await self._rpc(
  1535. self._svc.voice_calls.batch_get_active_calls(
  1536. voice_calls_pb2.BatchGetActiveCallsRequest(room_ids=room_ids),
  1537. headers=self._headers(),
  1538. )
  1539. )
  1540. data = pb_to_dict(resp)
  1541. return [ActiveCall.parse(c) for c in data.get("calls") or []]
  1542. async def join_call(self, room_id: str) -> bool:
  1543. resp = await self._rpc(
  1544. self._svc.voice_calls.join_call(
  1545. voice_calls_pb2.JoinCallRequest(room_id=room_id),
  1546. headers=self._headers(),
  1547. )
  1548. )
  1549. return resp.joined
  1550. async def leave_call(self, room_id: str) -> bool:
  1551. resp = await self._rpc(
  1552. self._svc.voice_calls.leave_call(
  1553. voice_calls_pb2.LeaveCallRequest(room_id=room_id),
  1554. headers=self._headers(),
  1555. )
  1556. )
  1557. return resp.left
  1558. async def get_call_token(self, room_id: str) -> str:
  1559. resp = await self._rpc(
  1560. self._svc.voice_calls.get_call_token(
  1561. voice_calls_pb2.GetCallTokenRequest(room_id=room_id),
  1562. headers=self._headers(),
  1563. )
  1564. )
  1565. return resp.token
  1566. # --- Admin: server --------------------------------------------------
  1567. async def admin_get_server_config(self) -> tuple[ServerConfig, ServerProfile]:
  1568. resp = await self._rpc(
  1569. self._svc.admin_server.get_server_config(
  1570. admin_server_pb2.GetServerConfigRequest(),
  1571. headers=self._headers(),
  1572. )
  1573. )
  1574. return (
  1575. ServerConfig.parse(pb_to_dict(resp.config)),
  1576. ServerProfile.parse(pb_to_dict(resp.public_profile)),
  1577. )
  1578. async def admin_update_server_config(
  1579. self,
  1580. *,
  1581. server_name: str | None = None,
  1582. description: str | None = None,
  1583. motd: str | None = None,
  1584. welcome_message: str | None = None,
  1585. ) -> tuple[ServerConfig, ServerProfile]:
  1586. req = admin_server_pb2.UpdateServerConfigRequest()
  1587. if server_name is not None:
  1588. req.server_name = server_name
  1589. if description is not None:
  1590. req.description = description
  1591. if motd is not None:
  1592. req.motd = motd
  1593. if welcome_message is not None:
  1594. req.welcome_message = welcome_message
  1595. resp = await self._rpc(
  1596. self._svc.admin_server.update_server_config(req, headers=self._headers())
  1597. )
  1598. return (
  1599. ServerConfig.parse(pb_to_dict(resp.config)),
  1600. ServerProfile.parse(pb_to_dict(resp.public_profile)),
  1601. )
  1602. async def admin_upload_server_logo(
  1603. self,
  1604. file_path: str | Path,
  1605. *,
  1606. content_type: str = "image/png",
  1607. ) -> ServerProfile:
  1608. p = Path(file_path)
  1609. req = admin_server_pb2.UploadServerLogoRequest(
  1610. image=common_pb2.ImageUpload(
  1611. image=p.read_bytes(), filename=p.name, content_type=content_type
  1612. )
  1613. )
  1614. resp = await self._rpc(
  1615. self._svc.admin_server.upload_server_logo(req, headers=self._headers())
  1616. )
  1617. return ServerProfile.parse(pb_to_dict(resp.public_profile))
  1618. async def admin_delete_server_logo(self) -> ServerProfile:
  1619. resp = await self._rpc(
  1620. self._svc.admin_server.delete_server_logo(
  1621. admin_server_pb2.DeleteServerLogoRequest(),
  1622. headers=self._headers(),
  1623. )
  1624. )
  1625. return ServerProfile.parse(pb_to_dict(resp.public_profile))
  1626. async def admin_upload_server_banner(
  1627. self,
  1628. file_path: str | Path,
  1629. *,
  1630. content_type: str = "image/png",
  1631. ) -> ServerProfile:
  1632. p = Path(file_path)
  1633. req = admin_server_pb2.UploadServerBannerRequest(
  1634. image=common_pb2.ImageUpload(
  1635. image=p.read_bytes(), filename=p.name, content_type=content_type
  1636. )
  1637. )
  1638. resp = await self._rpc(
  1639. self._svc.admin_server.upload_server_banner(req, headers=self._headers())
  1640. )
  1641. return ServerProfile.parse(pb_to_dict(resp.public_profile))
  1642. async def admin_delete_server_banner(self) -> ServerProfile:
  1643. resp = await self._rpc(
  1644. self._svc.admin_server.delete_server_banner(
  1645. admin_server_pb2.DeleteServerBannerRequest(),
  1646. headers=self._headers(),
  1647. )
  1648. )
  1649. return ServerProfile.parse(pb_to_dict(resp.public_profile))
  1650. async def admin_get_server_security_config(self) -> list[str]:
  1651. resp = await self._rpc(
  1652. self._svc.admin_server.get_server_security_config(
  1653. admin_server_pb2.GetServerSecurityConfigRequest(),
  1654. headers=self._headers(),
  1655. )
  1656. )
  1657. return list(resp.blocked_usernames)
  1658. async def admin_update_blocked_usernames(self, usernames: list[str]) -> list[str]:
  1659. resp = await self._rpc(
  1660. self._svc.admin_server.update_blocked_usernames(
  1661. admin_server_pb2.UpdateBlockedUsernamesRequest(
  1662. blocked_usernames=usernames
  1663. ),
  1664. headers=self._headers(),
  1665. )
  1666. )
  1667. return list(resp.blocked_usernames)
  1668. # --- Admin: room layout & sidebar links ---------------------------
  1669. async def admin_list_room_groups(self) -> list[AdminRoomLayoutGroup]:
  1670. resp = await self._rpc(
  1671. self._svc.admin_room_layout.list_room_groups(
  1672. room_layout_pb2.ListRoomGroupsRequest(),
  1673. headers=self._headers(),
  1674. )
  1675. )
  1676. data = pb_to_dict(resp)
  1677. return [
  1678. g
  1679. for g in (
  1680. AdminRoomLayoutGroup.parse(row) for row in data.get("groups") or []
  1681. )
  1682. if g is not None
  1683. ]
  1684. async def admin_create_room_group(
  1685. self, name: str, description: str = ""
  1686. ) -> AdminRoomLayoutGroup:
  1687. resp = await self._rpc(
  1688. self._svc.admin_room_layout.create_room_group(
  1689. room_layout_pb2.CreateRoomGroupRequest(
  1690. name=name, description=description
  1691. ),
  1692. headers=self._headers(),
  1693. )
  1694. )
  1695. group = AdminRoomLayoutGroup.parse(pb_to_dict(resp.group))
  1696. assert group is not None
  1697. return group
  1698. async def admin_update_room_group(
  1699. self,
  1700. group_id: str,
  1701. *,
  1702. name: str | None = None,
  1703. description: str | None = None,
  1704. ) -> AdminRoomLayoutGroup:
  1705. req = room_layout_pb2.UpdateRoomGroupRequest(group_id=group_id)
  1706. if name is not None:
  1707. req.name = name
  1708. if description is not None:
  1709. req.description = description
  1710. resp = await self._rpc(
  1711. self._svc.admin_room_layout.update_room_group(
  1712. req, headers=self._headers()
  1713. )
  1714. )
  1715. group = AdminRoomLayoutGroup.parse(pb_to_dict(resp.group))
  1716. assert group is not None
  1717. return group
  1718. async def admin_delete_room_group(self, group_id: str) -> bool:
  1719. resp = await self._rpc(
  1720. self._svc.admin_room_layout.delete_room_group(
  1721. room_layout_pb2.DeleteRoomGroupRequest(group_id=group_id),
  1722. headers=self._headers(),
  1723. )
  1724. )
  1725. return resp.deleted
  1726. async def admin_reorder_room_groups(
  1727. self, ordered_group_ids: list[str]
  1728. ) -> list[AdminRoomLayoutGroup]:
  1729. resp = await self._rpc(
  1730. self._svc.admin_room_layout.reorder_room_groups(
  1731. room_layout_pb2.ReorderRoomGroupsRequest(
  1732. ordered_group_ids=ordered_group_ids
  1733. ),
  1734. headers=self._headers(),
  1735. )
  1736. )
  1737. data = pb_to_dict(resp)
  1738. return [
  1739. g
  1740. for g in (
  1741. AdminRoomLayoutGroup.parse(row) for row in data.get("groups") or []
  1742. )
  1743. if g is not None
  1744. ]
  1745. async def admin_move_room_to_group(self, room_id: str, group_id: str) -> Room:
  1746. resp = await self._rpc(
  1747. self._svc.admin_room_layout.move_room_to_group(
  1748. room_layout_pb2.MoveRoomToGroupRequest(
  1749. room_id=room_id, group_id=group_id
  1750. ),
  1751. headers=self._headers(),
  1752. )
  1753. )
  1754. room = Room.parse(pb_to_dict(resp.room))
  1755. assert room is not None
  1756. return room
  1757. async def admin_reorder_sidebar_items_in_group(
  1758. self,
  1759. group_id: str,
  1760. items: list[tuple[AdminRoomLayoutItemKind, str]],
  1761. ) -> AdminRoomLayoutGroup:
  1762. req = room_layout_pb2.ReorderSidebarItemsInGroupRequest(group_id=group_id)
  1763. for kind, item_id in items:
  1764. item = req.items.add()
  1765. item.kind = kind.value
  1766. item.id = item_id
  1767. resp = await self._rpc(
  1768. self._svc.admin_room_layout.reorder_sidebar_items_in_group(
  1769. req, headers=self._headers()
  1770. )
  1771. )
  1772. group = AdminRoomLayoutGroup.parse(pb_to_dict(resp.group))
  1773. assert group is not None
  1774. return group
  1775. async def admin_create_sidebar_link(
  1776. self, group_id: str, label: str, url: str
  1777. ) -> SidebarLink | None:
  1778. resp = await self._rpc(
  1779. self._svc.admin_room_layout.create_sidebar_link(
  1780. room_layout_pb2.CreateSidebarLinkRequest(
  1781. group_id=group_id, label=label, url=url
  1782. ),
  1783. headers=self._headers(),
  1784. )
  1785. )
  1786. return SidebarLink.parse(pb_to_dict(resp).get("sidebarLink"))
  1787. async def admin_update_sidebar_link(
  1788. self,
  1789. link_id: str,
  1790. *,
  1791. label: str | None = None,
  1792. url: str | None = None,
  1793. ) -> SidebarLink | None:
  1794. req = room_layout_pb2.UpdateSidebarLinkRequest(link_id=link_id)
  1795. if label is not None:
  1796. req.label = label
  1797. if url is not None:
  1798. req.url = url
  1799. resp = await self._rpc(
  1800. self._svc.admin_room_layout.update_sidebar_link(
  1801. req, headers=self._headers()
  1802. )
  1803. )
  1804. return SidebarLink.parse(pb_to_dict(resp).get("sidebarLink"))
  1805. async def admin_delete_sidebar_link(self, link_id: str) -> bool:
  1806. resp = await self._rpc(
  1807. self._svc.admin_room_layout.delete_sidebar_link(
  1808. room_layout_pb2.DeleteSidebarLinkRequest(link_id=link_id),
  1809. headers=self._headers(),
  1810. )
  1811. )
  1812. return resp.deleted
  1813. async def admin_move_sidebar_link_to_group(
  1814. self, link_id: str, group_id: str
  1815. ) -> SidebarLink | None:
  1816. resp = await self._rpc(
  1817. self._svc.admin_room_layout.move_sidebar_link_to_group(
  1818. room_layout_pb2.MoveSidebarLinkToGroupRequest(
  1819. link_id=link_id, group_id=group_id
  1820. ),
  1821. headers=self._headers(),
  1822. )
  1823. )
  1824. return SidebarLink.parse(pb_to_dict(resp).get("sidebarLink"))
  1825. # --- Admin: users --------------------------------------------------
  1826. async def admin_list_members(
  1827. self,
  1828. *,
  1829. search: str = "",
  1830. limit: int | None = None,
  1831. offset: int | None = None,
  1832. ) -> tuple[list[AdminMember], list[Role], Page]:
  1833. req = admin_members_pb2.ListMembersRequest(search=search)
  1834. page = _page_pb(limit, offset)
  1835. if page is not None:
  1836. req.page.CopyFrom(page)
  1837. resp = await self._rpc(
  1838. self._svc.admin_users.list_members(req, headers=self._headers())
  1839. )
  1840. data = pb_to_dict(resp)
  1841. members = [
  1842. m
  1843. for m in (AdminMember.parse(row) for row in data.get("members") or [])
  1844. if m is not None
  1845. ]
  1846. roles = [
  1847. r for r in (Role.parse(row) for row in data.get("roles") or []) if r is not None
  1848. ]
  1849. return members, roles, Page.parse(data.get("page"))
  1850. async def admin_get_member(
  1851. self,
  1852. *,
  1853. user_id: str | None = None,
  1854. login: str | None = None,
  1855. ) -> AdminMemberDetail:
  1856. if bool(user_id) == bool(login):
  1857. raise ValueError("admin_get_member requires exactly one of user_id or login")
  1858. req = admin_members_pb2.GetMemberRequest()
  1859. if user_id:
  1860. req.user_id = user_id
  1861. else:
  1862. assert login is not None
  1863. req.login = login
  1864. resp = await self._rpc(
  1865. self._svc.admin_users.get_member(req, headers=self._headers())
  1866. )
  1867. return AdminMemberDetail.parse(pb_to_dict(resp))
  1868. async def admin_batch_get_members(self, user_ids: list[str]) -> list[AdminMember]:
  1869. resp = await self._rpc(
  1870. self._svc.admin_users.batch_get_members(
  1871. admin_members_pb2.BatchGetMembersRequest(user_ids=user_ids),
  1872. headers=self._headers(),
  1873. )
  1874. )
  1875. data = pb_to_dict(resp)
  1876. return [
  1877. m
  1878. for m in (AdminMember.parse(row) for row in data.get("members") or [])
  1879. if m is not None
  1880. ]
  1881. async def admin_assign_role(
  1882. self, user_id: str, role_name: str
  1883. ) -> AdminMember | None:
  1884. resp = await self._rpc(
  1885. self._svc.admin_users.assign_role(
  1886. admin_members_pb2.AssignRoleRequest(
  1887. user_id=user_id, role_name=role_name
  1888. ),
  1889. headers=self._headers(),
  1890. )
  1891. )
  1892. return AdminMember.parse(pb_to_dict(resp.member))
  1893. async def admin_revoke_role(
  1894. self, user_id: str, role_name: str
  1895. ) -> AdminMember | None:
  1896. resp = await self._rpc(
  1897. self._svc.admin_users.revoke_role(
  1898. admin_members_pb2.RevokeRoleRequest(
  1899. user_id=user_id, role_name=role_name
  1900. ),
  1901. headers=self._headers(),
  1902. )
  1903. )
  1904. return AdminMember.parse(pb_to_dict(resp.member))
  1905. async def admin_update_user(
  1906. self,
  1907. user_id: str,
  1908. *,
  1909. display_name: str | None = None,
  1910. login: str | None = None,
  1911. ) -> tuple[User | None, AdminMember | None]:
  1912. req = admin_members_pb2.UpdateUserRequest(user_id=user_id)
  1913. if display_name is not None:
  1914. req.display_name = display_name
  1915. if login is not None:
  1916. req.login = login
  1917. resp = await self._rpc(
  1918. self._svc.admin_users.update_user(req, headers=self._headers())
  1919. )
  1920. return (
  1921. User.parse(pb_to_dict(resp.user)),
  1922. AdminMember.parse(pb_to_dict(resp.member)),
  1923. )
  1924. async def admin_update_user_password(
  1925. self, user_id: str, password: str
  1926. ) -> AdminMember | None:
  1927. resp = await self._rpc(
  1928. self._svc.admin_users.update_user_password(
  1929. admin_members_pb2.UpdateUserPasswordRequest(
  1930. user_id=user_id, password=password
  1931. ),
  1932. headers=self._headers(),
  1933. )
  1934. )
  1935. return AdminMember.parse(pb_to_dict(resp.member))
  1936. async def admin_clear_username_cooldown(self, user_id: str) -> bool:
  1937. resp = await self._rpc(
  1938. self._svc.admin_users.clear_username_cooldown(
  1939. admin_members_pb2.ClearUsernameCooldownRequest(user_id=user_id),
  1940. headers=self._headers(),
  1941. )
  1942. )
  1943. return resp.cleared
  1944. async def admin_delete_user(
  1945. self, user_id: str, *, current_password: str = ""
  1946. ) -> bool:
  1947. resp = await self._rpc(
  1948. self._svc.admin_users.delete_user(
  1949. admin_members_pb2.DeleteUserRequest(
  1950. user_id=user_id, current_password=current_password
  1951. ),
  1952. headers=self._headers(),
  1953. )
  1954. )
  1955. return resp.deleted
  1956. # --- Admin: roles --------------------------------------------------
  1957. async def admin_list_roles(self) -> list[AdminRole]:
  1958. resp = await self._rpc(
  1959. self._svc.admin_roles.list_roles(
  1960. admin_roles_pb2.ListRolesRequest(), headers=self._headers()
  1961. )
  1962. )
  1963. data = pb_to_dict(resp)
  1964. return [
  1965. r
  1966. for r in (AdminRole.parse(row) for row in data.get("roles") or [])
  1967. if r is not None
  1968. ]
  1969. async def admin_get_role(self, name: str) -> AdminRoleDetail:
  1970. resp = await self._rpc(
  1971. self._svc.admin_roles.get_role(
  1972. admin_roles_pb2.GetRoleRequest(name=name), headers=self._headers()
  1973. )
  1974. )
  1975. return AdminRoleDetail.parse(pb_to_dict(resp))
  1976. async def admin_create_role(
  1977. self,
  1978. name: str,
  1979. *,
  1980. display_name: str = "",
  1981. description: str = "",
  1982. pingable: bool = False,
  1983. ) -> AdminRole | None:
  1984. resp = await self._rpc(
  1985. self._svc.admin_roles.create_role(
  1986. admin_roles_pb2.CreateRoleRequest(
  1987. name=name,
  1988. display_name=display_name,
  1989. description=description,
  1990. pingable=pingable,
  1991. ),
  1992. headers=self._headers(),
  1993. )
  1994. )
  1995. return AdminRole.parse(pb_to_dict(resp.role))
  1996. async def admin_update_role(
  1997. self,
  1998. name: str,
  1999. *,
  2000. display_name: str | None = None,
  2001. description: str | None = None,
  2002. pingable: bool | None = None,
  2003. ) -> AdminRole | None:
  2004. req = admin_roles_pb2.UpdateRoleRequest(name=name)
  2005. if display_name is not None:
  2006. req.display_name = display_name
  2007. if description is not None:
  2008. req.description = description
  2009. if pingable is not None:
  2010. req.pingable = pingable
  2011. resp = await self._rpc(
  2012. self._svc.admin_roles.update_role(req, headers=self._headers())
  2013. )
  2014. return AdminRole.parse(pb_to_dict(resp.role))
  2015. async def admin_delete_role(self, name: str) -> bool:
  2016. resp = await self._rpc(
  2017. self._svc.admin_roles.delete_role(
  2018. admin_roles_pb2.DeleteRoleRequest(name=name),
  2019. headers=self._headers(),
  2020. )
  2021. )
  2022. return resp.deleted
  2023. async def admin_reorder_roles(self, role_names: list[str]) -> list[AdminRole]:
  2024. resp = await self._rpc(
  2025. self._svc.admin_roles.reorder_roles(
  2026. admin_roles_pb2.ReorderRolesRequest(role_names=role_names),
  2027. headers=self._headers(),
  2028. )
  2029. )
  2030. data = pb_to_dict(resp)
  2031. return [
  2032. r
  2033. for r in (AdminRole.parse(row) for row in data.get("roles") or [])
  2034. if r is not None
  2035. ]
  2036. # --- Admin: event log / diagnostics / permissions ----------------
  2037. async def admin_list_events(
  2038. self,
  2039. *,
  2040. limit: int | None = None,
  2041. before: str | None = None,
  2042. event_type: str | None = None,
  2043. actor_id: str | None = None,
  2044. ) -> AdminEventLogPage:
  2045. """List durable EVT entries newest-first.
  2046. Args:
  2047. limit: Maximum entries to return (server clamps to its diagnostic
  2048. limit).
  2049. before: Exclusive sequence cursor; returned entries are older than
  2050. this sequence.
  2051. event_type: Filter by event payload type name (e.g.
  2052. ``"MessagePostedEvent"``).
  2053. actor_id: Filter by actor user ID.
  2054. """
  2055. req = event_log_pb2.ListEventsRequest()
  2056. if limit is not None:
  2057. req.limit = limit
  2058. if before is not None:
  2059. req.before = before
  2060. if event_type or actor_id:
  2061. fltr = event_log_pb2.AdminEventLogFilter()
  2062. if event_type:
  2063. fltr.event_type = event_type
  2064. if actor_id:
  2065. fltr.actor_id = actor_id
  2066. req.filter.CopyFrom(fltr)
  2067. resp = await self._rpc(
  2068. self._svc.admin_event_log.list_events(req, headers=self._headers())
  2069. )
  2070. return AdminEventLogPage.parse(pb_to_dict(resp))
  2071. async def admin_list_event_types(self) -> list[str]:
  2072. resp = await self._rpc(
  2073. self._svc.admin_event_log.list_event_types(
  2074. event_log_pb2.ListEventTypesRequest(),
  2075. headers=self._headers(),
  2076. )
  2077. )
  2078. return list(resp.event_types)
  2079. async def admin_get_event(self, sequence: str) -> AdminEventLogEntry | None:
  2080. """Read one durable EVT entry by its stream sequence.
  2081. Args:
  2082. sequence: EVT stream sequence string (e.g. ``"42"``).
  2083. Returns:
  2084. The matching ``AdminEventLogEntry``, or ``None`` when the proto
  2085. ``entry`` field is absent.
  2086. """
  2087. resp = await self._rpc(
  2088. self._svc.admin_event_log.get_event(
  2089. event_log_pb2.GetEventRequest(sequence=sequence),
  2090. headers=self._headers(),
  2091. )
  2092. )
  2093. return AdminEventLogEntry.parse(pb_to_dict(resp).get("entry"))
  2094. async def admin_get_system_info(self) -> AdminSystemInfoSnapshot:
  2095. from chattolib._pb.chatto.admin.v1 import diagnostics_pb2
  2096. resp = await self._rpc(
  2097. self._svc.admin_diagnostics.get_system_info(
  2098. diagnostics_pb2.GetSystemInfoRequest(),
  2099. headers=self._headers(),
  2100. )
  2101. )
  2102. return AdminSystemInfoSnapshot.parse(pb_to_dict(resp))
  2103. async def admin_get_role_permission_matrix(
  2104. self, role_name: str
  2105. ) -> RolePermissionMatrix | None:
  2106. resp = await self._rpc(
  2107. self._svc.admin_permissions.get_role_permission_matrix(
  2108. admin_permissions_pb2.GetRolePermissionMatrixRequest(
  2109. role_name=role_name
  2110. ),
  2111. headers=self._headers(),
  2112. )
  2113. )
  2114. return RolePermissionMatrix.parse(pb_to_dict(resp).get("matrix"))
  2115. async def admin_list_role_permission_decisions(
  2116. self, role_name: str
  2117. ) -> RolePermissionDecisions:
  2118. resp = await self._rpc(
  2119. self._svc.admin_permissions.list_role_permission_decisions(
  2120. admin_permissions_pb2.ListRolePermissionDecisionsRequest(
  2121. role_name=role_name
  2122. ),
  2123. headers=self._headers(),
  2124. )
  2125. )
  2126. return RolePermissionDecisions.parse(pb_to_dict(resp))
  2127. async def admin_get_user_permission_matrix(
  2128. self, user_id: str
  2129. ) -> UserPermissionMatrix | None:
  2130. resp = await self._rpc(
  2131. self._svc.admin_permissions.get_user_permission_matrix(
  2132. admin_permissions_pb2.GetUserPermissionMatrixRequest(user_id=user_id),
  2133. headers=self._headers(),
  2134. )
  2135. )
  2136. return UserPermissionMatrix.parse(pb_to_dict(resp).get("matrix"))
  2137. async def admin_list_user_permission_decisions(
  2138. self, user_id: str
  2139. ) -> UserPermissionDecisions:
  2140. resp = await self._rpc(
  2141. self._svc.admin_permissions.list_user_permission_decisions(
  2142. admin_permissions_pb2.ListUserPermissionDecisionsRequest(
  2143. user_id=user_id
  2144. ),
  2145. headers=self._headers(),
  2146. )
  2147. )
  2148. return UserPermissionDecisions.parse(pb_to_dict(resp))