client.py 77 KB

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