client.py 77 KB

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