client.py 76 KB

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