| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875187618771878187918801881188218831884188518861887188818891890189118921893189418951896189718981899190019011902190319041905190619071908190919101911191219131914191519161917191819191920192119221923192419251926192719281929193019311932193319341935193619371938193919401941194219431944194519461947194819491950195119521953195419551956195719581959196019611962196319641965196619671968196919701971197219731974197519761977197819791980198119821983198419851986198719881989199019911992199319941995199619971998199920002001200220032004200520062007200820092010201120122013201420152016201720182019202020212022202320242025202620272028202920302031203220332034203520362037203820392040204120422043204420452046204720482049205020512052205320542055205620572058205920602061206220632064206520662067206820692070207120722073207420752076207720782079208020812082208320842085208620872088208920902091209220932094209520962097209820992100210121022103210421052106210721082109211021112112211321142115211621172118211921202121212221232124212521262127212821292130213121322133213421352136213721382139214021412142214321442145214621472148214921502151215221532154215521562157215821592160216121622163216421652166216721682169217021712172217321742175217621772178217921802181218221832184218521862187218821892190219121922193219421952196219721982199220022012202220322042205220622072208220922102211221222132214221522162217221822192220222122222223222422252226222722282229223022312232223322342235223622372238223922402241224222432244224522462247224822492250225122522253225422552256225722582259226022612262226322642265226622672268226922702271227222732274227522762277227822792280228122822283228422852286228722882289229022912292229322942295229622972298229923002301230223032304230523062307230823092310231123122313231423152316231723182319232023212322232323242325232623272328232923302331233223332334233523362337233823392340234123422343234423452346234723482349235023512352235323542355 |
- """Main async client for the Chatto Connect API.
- Chatto migrated from GraphQL to a protobuf-first Connect API in v0.4.x
- (see ADR-042). The client speaks Connect binary protobuf via the official
- ``connectrpc`` Python package and the generated service stubs under
- ``chattolib._pb`` for all request/response operations. Realtime events
- live in ``chattolib.realtime``.
- """
- # mypy: disable-error-code="no-any-return"
- # Rationale: attribute access on generated protobuf messages is Any-typed
- # from mypy's perspective (the generated modules skip type checking via
- # follow_imports=skip). The runtime types are exactly what the return-type
- # annotations claim.
- from __future__ import annotations
- import hashlib
- from collections.abc import Awaitable
- from datetime import datetime
- from pathlib import Path
- from typing import Any, TypeVar
- import httpx
- # ConnectError isn't re-exported publicly by connectrpc.__init__; import from
- # its submodule so the top-level import path stays clean for callers.
- from connectrpc.errors import ConnectError # noqa: E402
- from chattolib._pb.chatto.admin.v1 import (
- event_log_pb2,
- room_layout_pb2,
- )
- from chattolib._pb.chatto.admin.v1 import (
- members_pb2 as admin_members_pb2,
- )
- from chattolib._pb.chatto.admin.v1 import (
- permissions_pb2 as admin_permissions_pb2,
- )
- from chattolib._pb.chatto.admin.v1 import (
- roles_pb2 as admin_roles_pb2,
- )
- from chattolib._pb.chatto.admin.v1 import (
- server_pb2 as admin_server_pb2,
- )
- from chattolib._pb.chatto.api.v1 import (
- account_pb2,
- asset_uploads_pb2,
- attachments_pb2,
- common_pb2,
- external_identities_pb2,
- link_previews_pb2,
- member_directory_pb2,
- messages_pb2,
- notification_preferences_pb2,
- notifications_pb2,
- pagination_pb2,
- presence_pb2,
- push_notifications_pb2,
- reactions_pb2,
- read_state_pb2,
- roles_pb2,
- room_directory_pb2,
- room_timeline_pb2,
- rooms_pb2,
- server_state_pb2,
- threads_pb2,
- user_status_pb2,
- viewer_pb2,
- voice_calls_pb2,
- )
- from chattolib._pb.chatto.auth.v1 import external_identity_auth_pb2
- from chattolib._pb.chatto.discovery.v1 import server_pb2 as discovery_server_pb2
- from chattolib._transport import (
- ServiceClients,
- build_service_clients,
- pb_to_dict,
- translate_connect_error,
- )
- from chattolib.exceptions import ChattoAuthError, ChattoError
- from chattolib.types import (
- ActiveCall,
- AdminEventLogEntry,
- AdminEventLogPage,
- AdminMember,
- AdminMemberDetail,
- AdminRole,
- AdminRoleDetail,
- AdminRoomLayoutGroup,
- AdminRoomLayoutItemKind,
- AdminSystemInfoSnapshot,
- Asset,
- AssetUpload,
- CustomUserStatus,
- DirectoryMember,
- ExternalIdentityAccount,
- ExternalIdentityProvider,
- FollowedThread,
- FollowedThreadsPage,
- ImageTransformOptions,
- LinkedExternalIdentity,
- LinkPreview,
- Message,
- Notification,
- NotificationLevel,
- NotificationPreference,
- NotificationsPage,
- Page,
- PendingExternalIdentity,
- PresenceStatus,
- Role,
- RolePermissionDecisions,
- RolePermissionMatrix,
- Room,
- RoomBan,
- RoomDirectoryScope,
- RoomGroup,
- RoomWithViewerState,
- ServerConfig,
- ServerLogin,
- ServerProfile,
- ServerRuntimeConfig,
- SidebarLink,
- TimeFormat,
- TimelinePage,
- User,
- UserPermissionDecisions,
- UserPermissionMatrix,
- UserSettings,
- ViewerSnapshot,
- ViewerUser,
- format_datetime,
- parse_datetime,
- )
- RES = TypeVar("RES")
- def _page_pb(
- limit: int | None, offset: int | None
- ) -> pagination_pb2.PageRequest | None:
- if limit is None and offset is None:
- return None
- return pagination_pb2.PageRequest(limit=limit or 0, offset=offset or 0)
- def _thumbnail_pb(
- opts: ImageTransformOptions | None,
- ) -> common_pb2.ImageTransformOptions | None:
- if opts is None:
- return None
- return common_pb2.ImageTransformOptions(
- width=opts.width, height=opts.height, fit=opts.fit.value
- )
- def _timestamp_pb(value: datetime | None) -> Any:
- from google.protobuf import timestamp_pb2
- if value is None:
- return None
- ts = timestamp_pb2.Timestamp()
- ts.FromJsonString(format_datetime(value))
- return ts
- class ChattoClient:
- """Async client for the Chatto Connect API.
- Usage::
- async with await ChattoClient.login("user", "pass") as client:
- viewer = await client.get_viewer()
- rooms = await client.list_rooms()
- # Or with an existing token:
- async with ChattoClient(token="cht_...") as client:
- ...
- """
- DEFAULT_BASE_URL = "https://chat.chatto.run"
- def __init__(
- self,
- token: str | None = None,
- *,
- base_url: str = DEFAULT_BASE_URL,
- session_cookie: str | None = None,
- service_clients: ServiceClients | None = None,
- ) -> None:
- self._base_url = base_url.rstrip("/")
- self._token = token
- self._session_cookie = session_cookie
- self._svc = service_clients or build_service_clients(self._base_url)
- self._owns_clients = service_clients is None
- @classmethod
- async def login(
- cls,
- login: str,
- password: str,
- *,
- base_url: str = DEFAULT_BASE_URL,
- ) -> ChattoClient:
- """Authenticate with username and password, returning a connected client.
- Uses Chatto's ``/auth/login`` HTTP endpoint (which is still exposed
- alongside the Connect API) and captures both the returned bearer token
- and any ``chatto_session`` cookie.
- """
- base = base_url.rstrip("/")
- async with httpx.AsyncClient() as http:
- resp = await http.post(
- f"{base}/auth/login",
- json={"login": login, "password": password},
- )
- if resp.status_code == 401:
- raise ChattoAuthError("Invalid credentials")
- resp.raise_for_status()
- body = resp.json()
- token = body.get("token")
- session_cookie = None
- if "set-cookie" in resp.headers:
- for cookie_header in resp.headers.get_list("set-cookie"):
- if cookie_header.startswith("chatto_session="):
- session_cookie = cookie_header.split(";")[0].split("=", 1)[1]
- break
- return cls(token=token, base_url=base_url, session_cookie=session_cookie)
- async def __aenter__(self) -> ChattoClient:
- return self
- async def __aexit__(self, *exc: Any) -> None:
- await self.close()
- async def close(self) -> None:
- if self._owns_clients:
- await self._svc.close()
- # --- Transport ------------------------------------------------------
- @property
- def base_url(self) -> str:
- return self._base_url
- @property
- def token(self) -> str | None:
- return self._token
- @property
- def session_cookie(self) -> str | None:
- return self._session_cookie
- @property
- def services(self) -> ServiceClients:
- """Direct access to the underlying ConnectRPC service clients.
- Useful when a caller wants to reach an RPC that this class doesn't
- expose yet, or wants full protobuf messages instead of the
- dataclass views returned by the high-level helpers.
- """
- return self._svc
- def _headers(self) -> dict[str, str]:
- headers: dict[str, str] = {}
- if self._token:
- headers["Authorization"] = f"Bearer {self._token}"
- if self._session_cookie:
- headers["Cookie"] = f"chatto_session={self._session_cookie}"
- return headers
- async def _rpc(self, coro: Awaitable[RES]) -> RES:
- """Await a ConnectRPC coroutine, translating errors."""
- try:
- return await coro
- except ConnectError as exc:
- raise translate_connect_error(exc) from exc
- # --- Server discovery ----------------------------------------------
- async def get_server(self) -> tuple[ServerProfile, ServerLogin]:
- """Public server profile and login options. Does not require auth."""
- resp = await self._rpc(
- self._svc.server_discovery.get_server(
- discovery_server_pb2.GetServerRequest(),
- headers=self._headers(),
- )
- )
- return (
- ServerProfile.parse(pb_to_dict(resp.profile)),
- ServerLogin.parse(pb_to_dict(resp.login)),
- )
- async def get_motd(self) -> str | None:
- resp = await self._rpc(
- self._svc.server.get_motd(
- server_state_pb2.GetMotdRequest(), headers=self._headers()
- )
- )
- if not resp.HasField("motd"):
- return None
- return resp.motd
- async def get_runtime_config(self) -> ServerRuntimeConfig:
- resp = await self._rpc(
- self._svc.server.get_runtime_config(
- server_state_pb2.GetRuntimeConfigRequest(),
- headers=self._headers(),
- )
- )
- return ServerRuntimeConfig.parse(pb_to_dict(resp.runtime))
- # --- Viewer ---------------------------------------------------------
- async def get_viewer(self) -> ViewerSnapshot:
- """Full authenticated viewer snapshot.
- Returns a :class:`~chattolib.types.ViewerSnapshot` containing the
- current user, capabilities, notification preferences, permission
- decisions, and viewer state. Callers that only need the current user's
- public profile should use :meth:`me` for a lightweight result.
- """
- resp = await self._rpc(
- self._svc.viewer.get_viewer(
- viewer_pb2.GetViewerRequest(), headers=self._headers()
- )
- )
- return ViewerSnapshot.parse(pb_to_dict(resp))
- async def viewer_user(self) -> ViewerUser | None:
- snapshot = await self.get_viewer()
- return snapshot.user
- async def me(self) -> User:
- """Return the authenticated user's public profile."""
- viewer = await self.viewer_user()
- if viewer is None or viewer.profile is None:
- raise ChattoAuthError("No authenticated viewer")
- return viewer.profile
- # --- MyAccount -----------------------------------------------------
- async def update_profile(
- self,
- *,
- display_name: str | None = None,
- login: str | None = None,
- ) -> User:
- req = account_pb2.UpdateProfileRequest()
- if display_name is not None:
- req.display_name = display_name
- if login is not None:
- req.login = login
- resp = await self._rpc(
- self._svc.account.update_profile(req, headers=self._headers())
- )
- user = User.parse(pb_to_dict(resp.user))
- assert user is not None
- return user
- async def upload_avatar(
- self,
- file_path: str | Path,
- *,
- content_type: str = "image/png",
- ) -> User:
- p = Path(file_path)
- req = account_pb2.UploadAvatarRequest(
- image=common_pb2.ImageUpload(
- image=p.read_bytes(),
- filename=p.name,
- content_type=content_type,
- )
- )
- resp = await self._rpc(
- self._svc.account.upload_avatar(req, headers=self._headers())
- )
- user = User.parse(pb_to_dict(resp.user))
- assert user is not None
- return user
- async def delete_avatar(self) -> User:
- resp = await self._rpc(
- self._svc.account.delete_avatar(
- account_pb2.DeleteAvatarRequest(), headers=self._headers()
- )
- )
- user = User.parse(pb_to_dict(resp.user))
- assert user is not None
- return user
- async def update_password(self, new_password: str, current_password: str = "") -> User:
- resp = await self._rpc(
- self._svc.account.update_password(
- account_pb2.UpdatePasswordRequest(
- password=new_password, current_password=current_password
- ),
- headers=self._headers(),
- )
- )
- user = User.parse(pb_to_dict(resp.user))
- assert user is not None
- return user
- async def update_settings(
- self,
- *,
- timezone: str | None = None,
- time_format: TimeFormat | None = None,
- ) -> UserSettings:
- req = account_pb2.UpdateSettingsRequest()
- if timezone is not None:
- req.timezone = timezone
- if time_format is not None:
- req.time_format = time_format.value
- resp = await self._rpc(
- self._svc.account.update_settings(req, headers=self._headers())
- )
- return UserSettings.parse(pb_to_dict(resp.settings))
- async def update_presence(
- self,
- status: PresenceStatus,
- *,
- user_selected: bool = True,
- ) -> PresenceStatus:
- if status in (PresenceStatus.UNSPECIFIED, PresenceStatus.OFFLINE):
- raise ValueError(
- "UNSPECIFIED and OFFLINE cannot be set as presence status; "
- "stop refreshing to go offline"
- )
- req = presence_pb2.UpdatePresenceRequest(
- status=status.value, user_selected=user_selected
- )
- resp = await self._rpc(
- self._svc.account.update_presence(req, headers=self._headers())
- )
- name = presence_pb2.PresenceStatus.Name(resp.status)
- return PresenceStatus(name)
- async def update_custom_status(
- self,
- emoji: str,
- text: str,
- *,
- expires_at: datetime | None = None,
- ) -> CustomUserStatus | None:
- req = user_status_pb2.UpdateCustomStatusRequest(emoji=emoji, text=text)
- if expires_at is not None:
- req.expires_at.CopyFrom(_timestamp_pb(expires_at))
- resp = await self._rpc(
- self._svc.account.update_custom_status(req, headers=self._headers())
- )
- return CustomUserStatus.parse(pb_to_dict(resp).get("status"))
- async def delete_custom_status(self) -> CustomUserStatus | None:
- resp = await self._rpc(
- self._svc.account.delete_custom_status(
- user_status_pb2.DeleteCustomStatusRequest(), headers=self._headers()
- )
- )
- return CustomUserStatus.parse(pb_to_dict(resp).get("status"))
- async def request_account_deletion(self) -> str:
- resp = await self._rpc(
- self._svc.account.request_account_deletion(
- account_pb2.RequestAccountDeletionRequest(),
- headers=self._headers(),
- )
- )
- return resp.confirmation_token
- async def delete_my_account(self, confirmation_token: str) -> bool:
- resp = await self._rpc(
- self._svc.account.delete_my_account(
- account_pb2.DeleteMyAccountRequest(
- confirmation_token=confirmation_token
- ),
- headers=self._headers(),
- )
- )
- return resp.deleted
- # --- Users ----------------------------------------------------------
- async def list_users(
- self,
- *,
- search: str = "",
- limit: int | None = None,
- offset: int | None = None,
- ) -> tuple[list[DirectoryMember], Page]:
- req = member_directory_pb2.ListUsersRequest(search=search)
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.users.list_users(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- users = [
- u
- for u in (DirectoryMember.parse(row) for row in data.get("users") or [])
- if u is not None
- ]
- return users, Page.parse(data.get("page"))
- async def get_user(
- self, *, user_id: str | None = None, login: str | None = None
- ) -> DirectoryMember | None:
- if bool(user_id) == bool(login):
- raise ValueError("get_user requires exactly one of user_id or login")
- req = member_directory_pb2.GetUserRequest()
- if user_id:
- req.user_id = user_id
- else:
- assert login is not None
- req.login = login
- resp = await self._rpc(
- self._svc.users.get_user(req, headers=self._headers())
- )
- return DirectoryMember.parse(pb_to_dict(resp.user))
- async def batch_get_users(self, user_ids: list[str]) -> list[DirectoryMember]:
- resp = await self._rpc(
- self._svc.users.batch_get_users(
- member_directory_pb2.BatchGetUsersRequest(user_ids=user_ids),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- u
- for u in (DirectoryMember.parse(row) for row in data.get("users") or [])
- if u is not None
- ]
- # --- Roles (public) -----------------------------------------------
- async def list_roles(self) -> list[Role]:
- resp = await self._rpc(
- self._svc.roles.list_roles(
- roles_pb2.ListRolesRequest(), headers=self._headers()
- )
- )
- data = pb_to_dict(resp)
- return [
- r for r in (Role.parse(row) for row in data.get("roles") or []) if r is not None
- ]
- async def get_role(self, name: str) -> Role | None:
- resp = await self._rpc(
- self._svc.roles.get_role(
- roles_pb2.GetRoleRequest(name=name), headers=self._headers()
- )
- )
- return Role.parse(pb_to_dict(resp.role))
- async def batch_get_roles(self, names: list[str]) -> list[Role]:
- resp = await self._rpc(
- self._svc.roles.batch_get_roles(
- roles_pb2.BatchGetRolesRequest(names=names), headers=self._headers()
- )
- )
- data = pb_to_dict(resp)
- return [
- r for r in (Role.parse(row) for row in data.get("roles") or []) if r is not None
- ]
- # --- Room directory ------------------------------------------------
- async def list_rooms(
- self, scope: RoomDirectoryScope = RoomDirectoryScope.ALL
- ) -> list[RoomWithViewerState]:
- resp = await self._rpc(
- self._svc.room_directory.list_rooms(
- room_directory_pb2.ListRoomsRequest(scope=scope.value),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- r
- for r in (RoomWithViewerState.parse(row) for row in data.get("rooms") or [])
- if r is not None
- ]
- async def list_room_groups(self) -> list[RoomGroup]:
- resp = await self._rpc(
- self._svc.room_directory.list_room_groups(
- room_directory_pb2.ListRoomGroupsRequest(),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- g
- for g in (RoomGroup.parse(row) for row in data.get("groups") or [])
- if g is not None
- ]
- async def get_room_group(self, group_id: str) -> RoomGroup | None:
- resp = await self._rpc(
- self._svc.room_directory.get_room_group(
- room_directory_pb2.GetRoomGroupRequest(group_id=group_id),
- headers=self._headers(),
- )
- )
- return RoomGroup.parse(pb_to_dict(resp.group))
- async def batch_get_room_groups(self, group_ids: list[str]) -> list[RoomGroup]:
- resp = await self._rpc(
- self._svc.room_directory.batch_get_room_groups(
- room_directory_pb2.BatchGetRoomGroupsRequest(group_ids=group_ids),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- g
- for g in (RoomGroup.parse(row) for row in data.get("groups") or [])
- if g is not None
- ]
- async def get_room(self, room_id: str) -> RoomWithViewerState | None:
- resp = await self._rpc(
- self._svc.room_directory.get_room(
- room_directory_pb2.GetRoomRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- return RoomWithViewerState.parse(pb_to_dict(resp.room))
- async def batch_get_rooms(self, room_ids: list[str]) -> list[RoomWithViewerState]:
- resp = await self._rpc(
- self._svc.room_directory.batch_get_rooms(
- room_directory_pb2.BatchGetRoomsRequest(room_ids=room_ids),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- r
- for r in (RoomWithViewerState.parse(row) for row in data.get("rooms") or [])
- if r is not None
- ]
- # --- Room lifecycle & membership -----------------------------------
- async def create_room(
- self,
- name: str,
- group_id: str,
- *,
- description: str = "",
- universal: bool = False,
- ) -> Room:
- resp = await self._rpc(
- self._svc.rooms.create_room(
- rooms_pb2.CreateRoomRequest(
- name=name,
- group_id=group_id,
- description=description,
- universal=universal,
- ),
- headers=self._headers(),
- )
- )
- room = Room.parse(pb_to_dict(resp.room))
- assert room is not None
- return room
- async def update_room(
- self,
- room_id: str,
- *,
- name: str | None = None,
- description: str | None = None,
- universal: bool | None = None,
- ) -> Room:
- req = rooms_pb2.UpdateRoomRequest(room_id=room_id)
- if name is not None:
- req.name = name
- if description is not None:
- req.description = description
- if universal is not None:
- req.universal = universal
- resp = await self._rpc(
- self._svc.rooms.update_room(req, headers=self._headers())
- )
- room = Room.parse(pb_to_dict(resp.room))
- assert room is not None
- return room
- async def archive_room(self, room_id: str) -> Room:
- resp = await self._rpc(
- self._svc.rooms.archive_room(
- rooms_pb2.ArchiveRoomRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- room = Room.parse(pb_to_dict(resp.room))
- assert room is not None
- return room
- async def unarchive_room(self, room_id: str) -> Room:
- resp = await self._rpc(
- self._svc.rooms.unarchive_room(
- rooms_pb2.UnarchiveRoomRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- room = Room.parse(pb_to_dict(resp.room))
- assert room is not None
- return room
- async def join_room(self, room_id: str) -> Room:
- resp = await self._rpc(
- self._svc.rooms.join_room(
- rooms_pb2.JoinRoomRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- room = Room.parse(pb_to_dict(resp.room))
- assert room is not None
- return room
- async def join_room_group(self, group_id: str) -> list[str]:
- resp = await self._rpc(
- self._svc.rooms.join_room_group(
- rooms_pb2.JoinRoomGroupRequest(group_id=group_id),
- headers=self._headers(),
- )
- )
- return list(resp.joined_room_ids)
- async def start_dm(self, participant_ids: list[str]) -> Room:
- resp = await self._rpc(
- self._svc.rooms.start_dm(
- rooms_pb2.StartDMRequest(participant_ids=participant_ids),
- headers=self._headers(),
- )
- )
- room = Room.parse(pb_to_dict(resp.room))
- assert room is not None
- return room
- async def leave_room(self, room_id: str) -> bool:
- resp = await self._rpc(
- self._svc.rooms.leave_room(
- rooms_pb2.LeaveRoomRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- return resp.left
- async def add_member(self, room_id: str, user_id: str) -> DirectoryMember | None:
- resp = await self._rpc(
- self._svc.rooms.add_member(
- rooms_pb2.AddMemberRequest(room_id=room_id, user_id=user_id),
- headers=self._headers(),
- )
- )
- return DirectoryMember.parse(pb_to_dict(resp.member))
- async def remove_member(self, room_id: str, user_id: str) -> bool:
- resp = await self._rpc(
- self._svc.rooms.remove_member(
- rooms_pb2.RemoveMemberRequest(room_id=room_id, user_id=user_id),
- headers=self._headers(),
- )
- )
- return resp.removed
- async def list_room_members(
- self,
- room_id: str,
- *,
- search: str = "",
- limit: int | None = None,
- offset: int | None = None,
- ) -> tuple[list[DirectoryMember], Page]:
- req = member_directory_pb2.ListRoomMembersRequest(room_id=room_id, search=search)
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.rooms.list_members(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- members = [
- m
- for m in (DirectoryMember.parse(row) for row in data.get("members") or [])
- if m is not None
- ]
- return members, Page.parse(data.get("page"))
- async def get_room_member(self, room_id: str, user_id: str) -> DirectoryMember | None:
- resp = await self._rpc(
- self._svc.rooms.get_member(
- member_directory_pb2.GetRoomMemberRequest(
- room_id=room_id, user_id=user_id
- ),
- headers=self._headers(),
- )
- )
- return DirectoryMember.parse(pb_to_dict(resp.member))
- async def batch_get_room_members(
- self, room_id: str, user_ids: list[str]
- ) -> list[DirectoryMember]:
- resp = await self._rpc(
- self._svc.rooms.batch_get_members(
- member_directory_pb2.BatchGetRoomMembersRequest(
- room_id=room_id, user_ids=user_ids
- ),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- m
- for m in (DirectoryMember.parse(row) for row in data.get("members") or [])
- if m is not None
- ]
- async def ban_member(
- self,
- room_id: str,
- user_id: str,
- reason: str,
- *,
- expires_at: datetime | None = None,
- ) -> bool:
- req = rooms_pb2.BanMemberRequest(
- room_id=room_id, user_id=user_id, reason=reason
- )
- if expires_at is not None:
- req.expires_at.CopyFrom(_timestamp_pb(expires_at))
- resp = await self._rpc(
- self._svc.rooms.ban_member(req, headers=self._headers())
- )
- return resp.banned
- async def unban_member(self, room_id: str, user_id: str, reason: str) -> bool:
- resp = await self._rpc(
- self._svc.rooms.unban_member(
- rooms_pb2.UnbanMemberRequest(
- room_id=room_id, user_id=user_id, reason=reason
- ),
- headers=self._headers(),
- )
- )
- return resp.unbanned
- async def list_bans(
- self,
- *,
- room_id: str = "",
- limit: int | None = None,
- offset: int | None = None,
- ) -> tuple[list[RoomBan], Page]:
- req = rooms_pb2.ListBansRequest(room_id=room_id)
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.rooms.list_bans(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- bans = [
- b
- for b in (RoomBan.parse(row) for row in data.get("bans") or [])
- if b is not None
- ]
- return bans, Page.parse(data.get("page"))
- async def update_typing_indicator(
- self, room_id: str, *, thread_root_event_id: str = ""
- ) -> bool:
- resp = await self._rpc(
- self._svc.rooms.update_typing_indicator(
- rooms_pb2.UpdateTypingIndicatorRequest(
- room_id=room_id, thread_root_event_id=thread_root_event_id
- ),
- headers=self._headers(),
- )
- )
- return resp.updated
- # --- Room timeline / read state -----------------------------------
- async def get_room_events(
- self,
- room_id: str,
- *,
- limit: int | None = None,
- before: str | None = None,
- after: str | None = None,
- ) -> TimelinePage:
- req = room_timeline_pb2.GetRoomEventsRequest(room_id=room_id)
- if limit is not None:
- req.limit = limit
- if before is not None:
- req.before = before
- elif after is not None:
- req.after = after
- resp = await self._rpc(
- self._svc.rooms.get_room_events(req, headers=self._headers())
- )
- return TimelinePage.parse(pb_to_dict(resp.page))
- async def get_room_events_around(
- self,
- room_id: str,
- event_id: str,
- *,
- limit: int | None = None,
- ) -> tuple[TimelinePage, int]:
- req = room_timeline_pb2.GetRoomEventsAroundRequest(
- room_id=room_id, event_id=event_id
- )
- if limit is not None:
- req.limit = limit
- resp = await self._rpc(
- self._svc.rooms.get_room_events_around(req, headers=self._headers())
- )
- return TimelinePage.parse(pb_to_dict(resp.page)), resp.target_index
- async def mark_room_as_read(
- self, room_id: str, up_to_event_id: str = ""
- ) -> tuple[datetime | None, datetime | None]:
- resp = await self._rpc(
- self._svc.rooms.mark_room_as_read(
- read_state_pb2.MarkRoomAsReadRequest(
- room_id=room_id, up_to_event_id=up_to_event_id
- ),
- headers=self._headers(),
- )
- )
- d = pb_to_dict(resp)
- return (
- parse_datetime(d.get("lastReadAt")),
- parse_datetime(d.get("previousLastReadAt")),
- )
- async def list_room_attachments(
- self,
- room_id: str,
- *,
- thumbnail: ImageTransformOptions | None = None,
- limit: int | None = None,
- offset: int | None = None,
- ) -> tuple[list[dict[str, Any]], Page]:
- req = rooms_pb2.ListRoomAttachmentsRequest(room_id=room_id)
- thumb = _thumbnail_pb(thumbnail)
- if thumb is not None:
- req.thumbnail.CopyFrom(thumb)
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.rooms.list_room_attachments(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- return list(data.get("attachments") or []), Page.parse(data.get("page"))
- # --- Messages -------------------------------------------------------
- async def fetch_link_preview(self, url: str) -> tuple[LinkPreview | None, str]:
- resp = await self._rpc(
- self._svc.messages.fetch_link_preview(
- link_previews_pb2.FetchLinkPreviewRequest(url=url),
- headers=self._headers(),
- )
- )
- return LinkPreview.parse(pb_to_dict(resp.preview)), resp.preview_token
- async def post_message(
- self,
- room_id: str,
- body: str = "",
- *,
- attachment_asset_ids: list[str] | None = None,
- thread_root_event_id: str = "",
- in_reply_to: str = "",
- also_send_to_channel: bool = False,
- link_preview_token: str = "",
- ) -> Message:
- req = messages_pb2.CreateMessageRequest(
- room_id=room_id,
- body=body,
- thread_root_event_id=thread_root_event_id,
- in_reply_to=in_reply_to,
- also_send_to_channel=also_send_to_channel,
- link_preview_token=link_preview_token,
- )
- if attachment_asset_ids:
- req.attachment_asset_ids.extend(attachment_asset_ids)
- resp = await self._rpc(
- self._svc.messages.create_message(req, headers=self._headers())
- )
- message = Message.parse(pb_to_dict(resp.message))
- assert message is not None
- return message
- async def update_message(
- self,
- room_id: str,
- event_id: str,
- *,
- body: str | None = None,
- also_send_to_channel: bool | None = None,
- ) -> Message:
- req = messages_pb2.UpdateMessageRequest(room_id=room_id, event_id=event_id)
- if body is not None:
- req.body = body
- if also_send_to_channel is not None:
- req.also_send_to_channel = also_send_to_channel
- resp = await self._rpc(
- self._svc.messages.update_message(req, headers=self._headers())
- )
- message = Message.parse(pb_to_dict(resp.message))
- assert message is not None
- return message
- async def delete_message(self, room_id: str, event_id: str) -> bool:
- resp = await self._rpc(
- self._svc.messages.delete_message(
- messages_pb2.DeleteMessageRequest(room_id=room_id, event_id=event_id),
- headers=self._headers(),
- )
- )
- return resp.deleted
- async def delete_attachment(
- self, room_id: str, event_id: str, attachment_id: str
- ) -> bool:
- resp = await self._rpc(
- self._svc.messages.delete_attachment(
- messages_pb2.DeleteAttachmentRequest(
- room_id=room_id, event_id=event_id, attachment_id=attachment_id
- ),
- headers=self._headers(),
- )
- )
- return resp.deleted
- async def delete_link_preview(self, room_id: str, event_id: str, url: str) -> bool:
- resp = await self._rpc(
- self._svc.messages.delete_link_preview(
- messages_pb2.DeleteLinkPreviewRequest(
- room_id=room_id, event_id=event_id, url=url
- ),
- headers=self._headers(),
- )
- )
- return resp.deleted
- async def get_message(self, room_id: str, event_id: str) -> Message | None:
- resp = await self._rpc(
- self._svc.messages.get_message(
- messages_pb2.GetMessageRequest(room_id=room_id, event_id=event_id),
- headers=self._headers(),
- )
- )
- return Message.parse(pb_to_dict(resp.message))
- async def batch_get_messages(
- self, room_id: str, event_ids: list[str]
- ) -> list[Message]:
- resp = await self._rpc(
- self._svc.messages.batch_get_messages(
- messages_pb2.BatchGetMessagesRequest(
- room_id=room_id, event_ids=event_ids
- ),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- m
- for m in (Message.parse(row) for row in data.get("messages") or [])
- if m is not None
- ]
- async def add_reaction(
- self, room_id: str, message_event_id: str, emoji: str
- ) -> bool:
- resp = await self._rpc(
- self._svc.messages.add_reaction(
- reactions_pb2.AddReactionRequest(
- room_id=room_id, message_event_id=message_event_id, emoji=emoji
- ),
- headers=self._headers(),
- )
- )
- return resp.added
- async def remove_reaction(
- self, room_id: str, message_event_id: str, emoji: str
- ) -> bool:
- resp = await self._rpc(
- self._svc.messages.remove_reaction(
- reactions_pb2.RemoveReactionRequest(
- room_id=room_id, message_event_id=message_event_id, emoji=emoji
- ),
- headers=self._headers(),
- )
- )
- return resp.removed
- # --- Threads --------------------------------------------------------
- async def follow_thread(self, room_id: str, thread_root_event_id: str) -> bool:
- resp = await self._rpc(
- self._svc.threads.follow_thread(
- threads_pb2.FollowThreadRequest(
- room_id=room_id, thread_root_event_id=thread_root_event_id
- ),
- headers=self._headers(),
- )
- )
- return resp.following
- async def unfollow_thread(self, room_id: str, thread_root_event_id: str) -> bool:
- resp = await self._rpc(
- self._svc.threads.unfollow_thread(
- threads_pb2.UnfollowThreadRequest(
- room_id=room_id, thread_root_event_id=thread_root_event_id
- ),
- headers=self._headers(),
- )
- )
- return resp.following
- async def list_followed_threads(
- self, *, limit: int | None = None, offset: int | None = None
- ) -> FollowedThreadsPage:
- req = threads_pb2.ListFollowedThreadsRequest()
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.threads.list_followed_threads(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- threads = [FollowedThread.parse(t) for t in data.get("threads") or []]
- users: dict[str, User] = {}
- includes = data.get("includes") or {}
- for uid, user_data in (includes.get("users") or {}).items():
- parsed = User.parse(user_data)
- if parsed is not None:
- users[uid] = parsed
- return FollowedThreadsPage(
- threads=threads, page=Page.parse(data.get("page")), users_by_id=users
- )
- async def get_thread_events(
- self,
- room_id: str,
- thread_root_event_id: str,
- *,
- limit: int | None = None,
- before: str | None = None,
- after: str | None = None,
- ) -> TimelinePage:
- req = room_timeline_pb2.GetThreadEventsRequest(
- room_id=room_id, thread_root_event_id=thread_root_event_id
- )
- if limit is not None:
- req.limit = limit
- if before is not None:
- req.before = before
- elif after is not None:
- req.after = after
- resp = await self._rpc(
- self._svc.threads.get_thread_events(req, headers=self._headers())
- )
- return TimelinePage.parse(pb_to_dict(resp.page))
- async def get_thread_events_around(
- self,
- room_id: str,
- thread_root_event_id: str,
- event_id: str,
- *,
- limit: int | None = None,
- ) -> tuple[TimelinePage, int]:
- req = room_timeline_pb2.GetThreadEventsAroundRequest(
- room_id=room_id,
- thread_root_event_id=thread_root_event_id,
- event_id=event_id,
- )
- if limit is not None:
- req.limit = limit
- resp = await self._rpc(
- self._svc.threads.get_thread_events_around(req, headers=self._headers())
- )
- return TimelinePage.parse(pb_to_dict(resp.page)), resp.target_index
- async def mark_thread_as_read(
- self,
- room_id: str,
- thread_root_event_id: str,
- up_to_event_id: str = "",
- ) -> datetime | None:
- resp = await self._rpc(
- self._svc.threads.mark_thread_as_read(
- read_state_pb2.MarkThreadAsReadRequest(
- room_id=room_id,
- thread_root_event_id=thread_root_event_id,
- up_to_event_id=up_to_event_id,
- ),
- headers=self._headers(),
- )
- )
- d = pb_to_dict(resp)
- return parse_datetime(d.get("previousReadAt"))
- # --- Notifications --------------------------------------------------
- async def list_notifications(
- self, *, limit: int | None = None, offset: int | None = None
- ) -> NotificationsPage:
- req = notifications_pb2.ListNotificationsRequest()
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.notifications.list_notifications(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- notifications = [Notification.parse(n) for n in data.get("notifications") or []]
- return NotificationsPage(
- notifications=notifications, page=Page.parse(data.get("page"))
- )
- async def has_notifications(self) -> bool:
- resp = await self._rpc(
- self._svc.notifications.has_notifications(
- notifications_pb2.HasNotificationsRequest(),
- headers=self._headers(),
- )
- )
- return resp.has_notifications
- async def get_notification(self, notification_id: str) -> Notification | None:
- resp = await self._rpc(
- self._svc.notifications.get_notification(
- notifications_pb2.GetNotificationRequest(notification_id=notification_id),
- headers=self._headers(),
- )
- )
- raw = pb_to_dict(resp).get("notification")
- return Notification.parse(raw) if raw else None
- async def batch_get_notifications(
- self, notification_ids: list[str]
- ) -> list[Notification]:
- resp = await self._rpc(
- self._svc.notifications.batch_get_notifications(
- notifications_pb2.BatchGetNotificationsRequest(
- notification_ids=notification_ids
- ),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [Notification.parse(n) for n in data.get("notifications") or []]
- async def list_room_notifications(
- self,
- room_id: str,
- *,
- limit: int | None = None,
- offset: int | None = None,
- ) -> NotificationsPage:
- req = notifications_pb2.ListRoomNotificationsRequest(room_id=room_id)
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.notifications.list_room_notifications(
- req, headers=self._headers()
- )
- )
- data = pb_to_dict(resp)
- return NotificationsPage(
- notifications=[Notification.parse(n) for n in data.get("notifications") or []],
- page=Page.parse(data.get("page")),
- )
- async def list_room_notification_counts(self) -> dict[str, int]:
- resp = await self._rpc(
- self._svc.notifications.list_room_notification_counts(
- notifications_pb2.ListRoomNotificationCountsRequest(),
- headers=self._headers(),
- )
- )
- return {row.room_id: row.total_count for row in resp.room_counts}
- async def dismiss_notification(self, notification_id: str) -> bool:
- resp = await self._rpc(
- self._svc.notifications.dismiss_notification(
- notifications_pb2.DismissNotificationRequest(
- notification_id=notification_id
- ),
- headers=self._headers(),
- )
- )
- return resp.dismissed
- async def dismiss_all_notifications(self) -> int:
- resp = await self._rpc(
- self._svc.notifications.dismiss_all_notifications(
- notifications_pb2.DismissAllNotificationsRequest(),
- headers=self._headers(),
- )
- )
- return resp.dismissed_count
- # --- Notification preferences --------------------------------------
- async def get_server_notification_preference(self) -> NotificationPreference:
- resp = await self._rpc(
- self._svc.notification_prefs.get_server_notification_preference(
- notification_preferences_pb2.GetServerNotificationPreferenceRequest(),
- headers=self._headers(),
- )
- )
- return NotificationPreference.parse(pb_to_dict(resp.preference))
- async def update_server_notification_preference(
- self, level: NotificationLevel
- ) -> NotificationPreference:
- resp = await self._rpc(
- self._svc.notification_prefs.update_server_notification_preference(
- notification_preferences_pb2.UpdateServerNotificationPreferenceRequest(
- level=level.value
- ),
- headers=self._headers(),
- )
- )
- return NotificationPreference.parse(pb_to_dict(resp.preference))
- async def get_room_notification_preference(
- self, room_id: str
- ) -> NotificationPreference:
- resp = await self._rpc(
- self._svc.notification_prefs.get_room_notification_preference(
- notification_preferences_pb2.GetRoomNotificationPreferenceRequest(
- room_id=room_id
- ),
- headers=self._headers(),
- )
- )
- return NotificationPreference.parse(pb_to_dict(resp.preference))
- async def update_room_notification_preference(
- self, room_id: str, level: NotificationLevel
- ) -> NotificationPreference:
- resp = await self._rpc(
- self._svc.notification_prefs.update_room_notification_preference(
- notification_preferences_pb2.UpdateRoomNotificationPreferenceRequest(
- room_id=room_id, level=level.value
- ),
- headers=self._headers(),
- )
- )
- return NotificationPreference.parse(pb_to_dict(resp.preference))
- # --- Push notifications --------------------------------------------
- async def subscribe_push(
- self,
- endpoint: str,
- p256dh: str,
- auth: str,
- *,
- user_agent: str | None = None,
- ) -> bool:
- req = push_notifications_pb2.SubscribePushRequest(
- endpoint=endpoint, p256dh=p256dh, auth=auth
- )
- if user_agent is not None:
- req.user_agent = user_agent
- resp = await self._rpc(
- self._svc.push.subscribe(req, headers=self._headers())
- )
- return resp.subscribed
- async def unsubscribe_push(self, endpoint: str) -> bool:
- resp = await self._rpc(
- self._svc.push.unsubscribe(
- push_notifications_pb2.UnsubscribePushRequest(endpoint=endpoint),
- headers=self._headers(),
- )
- )
- return resp.unsubscribed
- # --- Assets ---------------------------------------------------------
- async def get_asset(
- self,
- room_id: str,
- asset_id: str,
- *,
- thumbnail: ImageTransformOptions | None = None,
- ) -> Asset | None:
- req = attachments_pb2.GetAssetRequest(room_id=room_id, asset_id=asset_id)
- thumb = _thumbnail_pb(thumbnail)
- if thumb is not None:
- req.thumbnail.CopyFrom(thumb)
- resp = await self._rpc(
- self._svc.assets.get_asset(req, headers=self._headers())
- )
- return Asset.parse(pb_to_dict(resp.asset))
- async def batch_get_assets(
- self,
- room_id: str,
- asset_ids: list[str],
- *,
- thumbnail: ImageTransformOptions | None = None,
- ) -> list[Asset]:
- req = attachments_pb2.BatchGetAssetsRequest(
- room_id=room_id, asset_ids=asset_ids
- )
- thumb = _thumbnail_pb(thumbnail)
- if thumb is not None:
- req.thumbnail.CopyFrom(thumb)
- resp = await self._rpc(
- self._svc.assets.batch_get_assets(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- return [
- a
- for a in (Asset.parse(row) for row in data.get("assets") or [])
- if a is not None
- ]
- # --- Asset uploads ------------------------------------------------
- async def create_upload(
- self,
- room_id: str,
- filename: str,
- size: int,
- sha256: str,
- *,
- content_type: str = "",
- ) -> AssetUpload:
- resp = await self._rpc(
- self._svc.asset_uploads.create_upload(
- asset_uploads_pb2.CreateUploadRequest(
- room_id=room_id,
- filename=filename,
- content_type=content_type,
- size=size,
- sha256=sha256,
- ),
- headers=self._headers(),
- )
- )
- upload = AssetUpload.parse(pb_to_dict(resp.upload))
- assert upload is not None
- return upload
- async def upload_chunk(
- self, upload_id: str, offset: int, content: bytes, chunk_sha256: str
- ) -> AssetUpload:
- resp = await self._rpc(
- self._svc.asset_uploads.upload_chunk(
- asset_uploads_pb2.UploadChunkRequest(
- upload_id=upload_id,
- offset=offset,
- content=content,
- chunk_sha256=chunk_sha256,
- ),
- headers=self._headers(),
- )
- )
- upload = AssetUpload.parse(pb_to_dict(resp.upload))
- assert upload is not None
- return upload
- async def get_upload(self, upload_id: str) -> AssetUpload:
- resp = await self._rpc(
- self._svc.asset_uploads.get_upload(
- asset_uploads_pb2.GetUploadRequest(upload_id=upload_id),
- headers=self._headers(),
- )
- )
- upload = AssetUpload.parse(pb_to_dict(resp.upload))
- assert upload is not None
- return upload
- async def complete_upload(
- self, upload_id: str
- ) -> tuple[AssetUpload, Asset | None]:
- resp = await self._rpc(
- self._svc.asset_uploads.complete_upload(
- asset_uploads_pb2.CompleteUploadRequest(upload_id=upload_id),
- headers=self._headers(),
- )
- )
- upload = AssetUpload.parse(pb_to_dict(resp.upload))
- assert upload is not None
- return upload, Asset.parse(pb_to_dict(resp.asset))
- async def cancel_upload(self, upload_id: str) -> AssetUpload:
- resp = await self._rpc(
- self._svc.asset_uploads.cancel_upload(
- asset_uploads_pb2.CancelUploadRequest(upload_id=upload_id),
- headers=self._headers(),
- )
- )
- upload = AssetUpload.parse(pb_to_dict(resp.upload))
- assert upload is not None
- return upload
- async def upload_attachment(
- self,
- room_id: str,
- file_path: str | Path,
- *,
- content_type: str = "",
- filename: str | None = None,
- ) -> Asset:
- """Upload a file as a room attachment and return the resulting Asset."""
- path = Path(file_path)
- data = path.read_bytes()
- size = len(data)
- sha = hashlib.sha256(data).hexdigest()
- upload = await self.create_upload(
- room_id, filename or path.name, size, sha, content_type=content_type
- )
- chunk_size = upload.max_chunk_size or 512 * 1024
- offset = upload.committed_offset
- while offset < size:
- end = min(offset + chunk_size, size)
- chunk = data[offset:end]
- chunk_sha = hashlib.sha256(chunk).hexdigest()
- upload = await self.upload_chunk(upload.upload_id, offset, chunk, chunk_sha)
- if upload.committed_offset <= offset:
- raise ChattoError(
- f"upload stalled at offset {offset} (server reported "
- f"committed_offset={upload.committed_offset})"
- )
- offset = upload.committed_offset
- upload, asset = await self.complete_upload(upload.upload_id)
- if asset is None:
- raise ChattoError("upload completed but server returned no asset")
- return asset
- # --- MyAccount external identities --------------------------------
- async def list_external_identities(
- self,
- ) -> tuple[list[ExternalIdentityProvider], list[LinkedExternalIdentity]]:
- resp = await self._rpc(
- self._svc.account.list_external_identities(
- external_identities_pb2.ListExternalIdentitiesRequest(),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- providers = [
- ExternalIdentityProvider.parse(p) for p in data.get("providers") or []
- ]
- linked = [
- LinkedExternalIdentity.parse(li)
- for li in data.get("linkedIdentities") or []
- ]
- return providers, linked
- async def start_external_identity_link(
- self,
- provider_id: str,
- *,
- redirect_path: str = "",
- current_password: str = "",
- ) -> str:
- resp = await self._rpc(
- self._svc.account.start_external_identity_link(
- external_identities_pb2.StartExternalIdentityLinkRequest(
- provider_id=provider_id,
- redirect_path=redirect_path,
- current_password=current_password,
- ),
- headers=self._headers(),
- )
- )
- return resp.start_url
- async def disconnect_external_identity(
- self, subject_hash: str, *, current_password: str = ""
- ) -> bool:
- resp = await self._rpc(
- self._svc.account.disconnect_external_identity(
- external_identities_pb2.DisconnectExternalIdentityRequest(
- subject_hash=subject_hash, current_password=current_password
- ),
- headers=self._headers(),
- )
- )
- return resp.disconnected
- # --- ExternalIdentityAuthService (public OAuth handoff) -----------
- async def get_pending_external_identity(
- self, token: str
- ) -> PendingExternalIdentity | None:
- resp = await self._rpc(
- self._svc.external_auth.get_pending_external_identity(
- external_identity_auth_pb2.GetPendingExternalIdentityRequest(
- token=token
- ),
- headers=self._headers(),
- )
- )
- return PendingExternalIdentity.parse(pb_to_dict(resp).get("pending"))
- async def create_external_identity_account(
- self, token: str, login: str
- ) -> ExternalIdentityAccount | None:
- resp = await self._rpc(
- self._svc.external_auth.create_external_identity_account(
- external_identity_auth_pb2.CreateExternalIdentityAccountRequest(
- token=token, login=login
- ),
- headers=self._headers(),
- )
- )
- return ExternalIdentityAccount.parse(pb_to_dict(resp))
- async def confirm_external_identity_link(
- self, token: str
- ) -> LinkedExternalIdentity | None:
- resp = await self._rpc(
- self._svc.external_auth.confirm_external_identity_link(
- external_identity_auth_pb2.ConfirmExternalIdentityLinkRequest(
- token=token
- ),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp).get("linkedIdentity")
- return LinkedExternalIdentity.parse(data) if data else None
- async def cancel_external_identity_flow(self, token: str) -> bool:
- resp = await self._rpc(
- self._svc.external_auth.cancel_external_identity_flow(
- external_identity_auth_pb2.CancelExternalIdentityFlowRequest(
- token=token
- ),
- headers=self._headers(),
- )
- )
- return resp.cancelled
- # --- Voice calls ----------------------------------------------------
- async def list_active_calls(self) -> list[ActiveCall]:
- resp = await self._rpc(
- self._svc.voice_calls.list_active_calls(
- voice_calls_pb2.ListActiveCallsRequest(), headers=self._headers()
- )
- )
- data = pb_to_dict(resp)
- return [ActiveCall.parse(c) for c in data.get("calls") or []]
- async def get_active_call(self, room_id: str) -> ActiveCall | None:
- resp = await self._rpc(
- self._svc.voice_calls.get_active_call(
- voice_calls_pb2.GetActiveCallRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- raw = pb_to_dict(resp).get("call")
- return ActiveCall.parse(raw) if raw else None
- async def batch_get_active_calls(self, room_ids: list[str]) -> list[ActiveCall]:
- resp = await self._rpc(
- self._svc.voice_calls.batch_get_active_calls(
- voice_calls_pb2.BatchGetActiveCallsRequest(room_ids=room_ids),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [ActiveCall.parse(c) for c in data.get("calls") or []]
- async def join_call(self, room_id: str) -> bool:
- resp = await self._rpc(
- self._svc.voice_calls.join_call(
- voice_calls_pb2.JoinCallRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- return resp.joined
- async def leave_call(self, room_id: str) -> bool:
- resp = await self._rpc(
- self._svc.voice_calls.leave_call(
- voice_calls_pb2.LeaveCallRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- return resp.left
- async def get_call_token(self, room_id: str) -> str:
- resp = await self._rpc(
- self._svc.voice_calls.get_call_token(
- voice_calls_pb2.GetCallTokenRequest(room_id=room_id),
- headers=self._headers(),
- )
- )
- return resp.token
- # --- Admin: server --------------------------------------------------
- async def admin_get_server_config(self) -> tuple[ServerConfig, ServerProfile]:
- resp = await self._rpc(
- self._svc.admin_server.get_server_config(
- admin_server_pb2.GetServerConfigRequest(),
- headers=self._headers(),
- )
- )
- return (
- ServerConfig.parse(pb_to_dict(resp.config)),
- ServerProfile.parse(pb_to_dict(resp.public_profile)),
- )
- async def admin_update_server_config(
- self,
- *,
- server_name: str | None = None,
- description: str | None = None,
- motd: str | None = None,
- welcome_message: str | None = None,
- ) -> tuple[ServerConfig, ServerProfile]:
- req = admin_server_pb2.UpdateServerConfigRequest()
- if server_name is not None:
- req.server_name = server_name
- if description is not None:
- req.description = description
- if motd is not None:
- req.motd = motd
- if welcome_message is not None:
- req.welcome_message = welcome_message
- resp = await self._rpc(
- self._svc.admin_server.update_server_config(req, headers=self._headers())
- )
- return (
- ServerConfig.parse(pb_to_dict(resp.config)),
- ServerProfile.parse(pb_to_dict(resp.public_profile)),
- )
- async def admin_upload_server_logo(
- self,
- file_path: str | Path,
- *,
- content_type: str = "image/png",
- ) -> ServerProfile:
- p = Path(file_path)
- req = admin_server_pb2.UploadServerLogoRequest(
- image=common_pb2.ImageUpload(
- image=p.read_bytes(), filename=p.name, content_type=content_type
- )
- )
- resp = await self._rpc(
- self._svc.admin_server.upload_server_logo(req, headers=self._headers())
- )
- return ServerProfile.parse(pb_to_dict(resp.public_profile))
- async def admin_delete_server_logo(self) -> ServerProfile:
- resp = await self._rpc(
- self._svc.admin_server.delete_server_logo(
- admin_server_pb2.DeleteServerLogoRequest(),
- headers=self._headers(),
- )
- )
- return ServerProfile.parse(pb_to_dict(resp.public_profile))
- async def admin_upload_server_banner(
- self,
- file_path: str | Path,
- *,
- content_type: str = "image/png",
- ) -> ServerProfile:
- p = Path(file_path)
- req = admin_server_pb2.UploadServerBannerRequest(
- image=common_pb2.ImageUpload(
- image=p.read_bytes(), filename=p.name, content_type=content_type
- )
- )
- resp = await self._rpc(
- self._svc.admin_server.upload_server_banner(req, headers=self._headers())
- )
- return ServerProfile.parse(pb_to_dict(resp.public_profile))
- async def admin_delete_server_banner(self) -> ServerProfile:
- resp = await self._rpc(
- self._svc.admin_server.delete_server_banner(
- admin_server_pb2.DeleteServerBannerRequest(),
- headers=self._headers(),
- )
- )
- return ServerProfile.parse(pb_to_dict(resp.public_profile))
- async def admin_get_server_security_config(self) -> list[str]:
- resp = await self._rpc(
- self._svc.admin_server.get_server_security_config(
- admin_server_pb2.GetServerSecurityConfigRequest(),
- headers=self._headers(),
- )
- )
- return list(resp.blocked_usernames)
- async def admin_update_blocked_usernames(self, usernames: list[str]) -> list[str]:
- resp = await self._rpc(
- self._svc.admin_server.update_blocked_usernames(
- admin_server_pb2.UpdateBlockedUsernamesRequest(
- blocked_usernames=usernames
- ),
- headers=self._headers(),
- )
- )
- return list(resp.blocked_usernames)
- # --- Admin: room layout & sidebar links ---------------------------
- async def admin_list_room_groups(self) -> list[AdminRoomLayoutGroup]:
- resp = await self._rpc(
- self._svc.admin_room_layout.list_room_groups(
- room_layout_pb2.ListRoomGroupsRequest(),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- g
- for g in (
- AdminRoomLayoutGroup.parse(row) for row in data.get("groups") or []
- )
- if g is not None
- ]
- async def admin_create_room_group(
- self, name: str, description: str = ""
- ) -> AdminRoomLayoutGroup:
- resp = await self._rpc(
- self._svc.admin_room_layout.create_room_group(
- room_layout_pb2.CreateRoomGroupRequest(
- name=name, description=description
- ),
- headers=self._headers(),
- )
- )
- group = AdminRoomLayoutGroup.parse(pb_to_dict(resp.group))
- assert group is not None
- return group
- async def admin_update_room_group(
- self,
- group_id: str,
- *,
- name: str | None = None,
- description: str | None = None,
- ) -> AdminRoomLayoutGroup:
- req = room_layout_pb2.UpdateRoomGroupRequest(group_id=group_id)
- if name is not None:
- req.name = name
- if description is not None:
- req.description = description
- resp = await self._rpc(
- self._svc.admin_room_layout.update_room_group(
- req, headers=self._headers()
- )
- )
- group = AdminRoomLayoutGroup.parse(pb_to_dict(resp.group))
- assert group is not None
- return group
- async def admin_delete_room_group(self, group_id: str) -> bool:
- resp = await self._rpc(
- self._svc.admin_room_layout.delete_room_group(
- room_layout_pb2.DeleteRoomGroupRequest(group_id=group_id),
- headers=self._headers(),
- )
- )
- return resp.deleted
- async def admin_reorder_room_groups(
- self, ordered_group_ids: list[str]
- ) -> list[AdminRoomLayoutGroup]:
- resp = await self._rpc(
- self._svc.admin_room_layout.reorder_room_groups(
- room_layout_pb2.ReorderRoomGroupsRequest(
- ordered_group_ids=ordered_group_ids
- ),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- g
- for g in (
- AdminRoomLayoutGroup.parse(row) for row in data.get("groups") or []
- )
- if g is not None
- ]
- async def admin_move_room_to_group(self, room_id: str, group_id: str) -> Room:
- resp = await self._rpc(
- self._svc.admin_room_layout.move_room_to_group(
- room_layout_pb2.MoveRoomToGroupRequest(
- room_id=room_id, group_id=group_id
- ),
- headers=self._headers(),
- )
- )
- room = Room.parse(pb_to_dict(resp.room))
- assert room is not None
- return room
- async def admin_reorder_sidebar_items_in_group(
- self,
- group_id: str,
- items: list[tuple[AdminRoomLayoutItemKind, str]],
- ) -> AdminRoomLayoutGroup:
- req = room_layout_pb2.ReorderSidebarItemsInGroupRequest(group_id=group_id)
- for kind, item_id in items:
- item = req.items.add()
- item.kind = kind.value
- item.id = item_id
- resp = await self._rpc(
- self._svc.admin_room_layout.reorder_sidebar_items_in_group(
- req, headers=self._headers()
- )
- )
- group = AdminRoomLayoutGroup.parse(pb_to_dict(resp.group))
- assert group is not None
- return group
- async def admin_create_sidebar_link(
- self, group_id: str, label: str, url: str
- ) -> SidebarLink | None:
- resp = await self._rpc(
- self._svc.admin_room_layout.create_sidebar_link(
- room_layout_pb2.CreateSidebarLinkRequest(
- group_id=group_id, label=label, url=url
- ),
- headers=self._headers(),
- )
- )
- return SidebarLink.parse(pb_to_dict(resp).get("sidebarLink"))
- async def admin_update_sidebar_link(
- self,
- link_id: str,
- *,
- label: str | None = None,
- url: str | None = None,
- ) -> SidebarLink | None:
- req = room_layout_pb2.UpdateSidebarLinkRequest(link_id=link_id)
- if label is not None:
- req.label = label
- if url is not None:
- req.url = url
- resp = await self._rpc(
- self._svc.admin_room_layout.update_sidebar_link(
- req, headers=self._headers()
- )
- )
- return SidebarLink.parse(pb_to_dict(resp).get("sidebarLink"))
- async def admin_delete_sidebar_link(self, link_id: str) -> bool:
- resp = await self._rpc(
- self._svc.admin_room_layout.delete_sidebar_link(
- room_layout_pb2.DeleteSidebarLinkRequest(link_id=link_id),
- headers=self._headers(),
- )
- )
- return resp.deleted
- async def admin_move_sidebar_link_to_group(
- self, link_id: str, group_id: str
- ) -> SidebarLink | None:
- resp = await self._rpc(
- self._svc.admin_room_layout.move_sidebar_link_to_group(
- room_layout_pb2.MoveSidebarLinkToGroupRequest(
- link_id=link_id, group_id=group_id
- ),
- headers=self._headers(),
- )
- )
- return SidebarLink.parse(pb_to_dict(resp).get("sidebarLink"))
- # --- Admin: users --------------------------------------------------
- async def admin_list_members(
- self,
- *,
- search: str = "",
- limit: int | None = None,
- offset: int | None = None,
- ) -> tuple[list[AdminMember], list[Role], Page]:
- req = admin_members_pb2.ListMembersRequest(search=search)
- page = _page_pb(limit, offset)
- if page is not None:
- req.page.CopyFrom(page)
- resp = await self._rpc(
- self._svc.admin_users.list_members(req, headers=self._headers())
- )
- data = pb_to_dict(resp)
- members = [
- m
- for m in (AdminMember.parse(row) for row in data.get("members") or [])
- if m is not None
- ]
- roles = [
- r for r in (Role.parse(row) for row in data.get("roles") or []) if r is not None
- ]
- return members, roles, Page.parse(data.get("page"))
- async def admin_get_member(
- self,
- *,
- user_id: str | None = None,
- login: str | None = None,
- ) -> AdminMemberDetail:
- if bool(user_id) == bool(login):
- raise ValueError("admin_get_member requires exactly one of user_id or login")
- req = admin_members_pb2.GetMemberRequest()
- if user_id:
- req.user_id = user_id
- else:
- assert login is not None
- req.login = login
- resp = await self._rpc(
- self._svc.admin_users.get_member(req, headers=self._headers())
- )
- return AdminMemberDetail.parse(pb_to_dict(resp))
- async def admin_batch_get_members(self, user_ids: list[str]) -> list[AdminMember]:
- resp = await self._rpc(
- self._svc.admin_users.batch_get_members(
- admin_members_pb2.BatchGetMembersRequest(user_ids=user_ids),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- m
- for m in (AdminMember.parse(row) for row in data.get("members") or [])
- if m is not None
- ]
- async def admin_assign_role(
- self, user_id: str, role_name: str
- ) -> AdminMember | None:
- resp = await self._rpc(
- self._svc.admin_users.assign_role(
- admin_members_pb2.AssignRoleRequest(
- user_id=user_id, role_name=role_name
- ),
- headers=self._headers(),
- )
- )
- return AdminMember.parse(pb_to_dict(resp.member))
- async def admin_revoke_role(
- self, user_id: str, role_name: str
- ) -> AdminMember | None:
- resp = await self._rpc(
- self._svc.admin_users.revoke_role(
- admin_members_pb2.RevokeRoleRequest(
- user_id=user_id, role_name=role_name
- ),
- headers=self._headers(),
- )
- )
- return AdminMember.parse(pb_to_dict(resp.member))
- async def admin_update_user(
- self,
- user_id: str,
- *,
- display_name: str | None = None,
- login: str | None = None,
- ) -> tuple[User | None, AdminMember | None]:
- req = admin_members_pb2.UpdateUserRequest(user_id=user_id)
- if display_name is not None:
- req.display_name = display_name
- if login is not None:
- req.login = login
- resp = await self._rpc(
- self._svc.admin_users.update_user(req, headers=self._headers())
- )
- return (
- User.parse(pb_to_dict(resp.user)),
- AdminMember.parse(pb_to_dict(resp.member)),
- )
- async def admin_update_user_password(
- self, user_id: str, password: str
- ) -> AdminMember | None:
- resp = await self._rpc(
- self._svc.admin_users.update_user_password(
- admin_members_pb2.UpdateUserPasswordRequest(
- user_id=user_id, password=password
- ),
- headers=self._headers(),
- )
- )
- return AdminMember.parse(pb_to_dict(resp.member))
- async def admin_clear_username_cooldown(self, user_id: str) -> bool:
- resp = await self._rpc(
- self._svc.admin_users.clear_username_cooldown(
- admin_members_pb2.ClearUsernameCooldownRequest(user_id=user_id),
- headers=self._headers(),
- )
- )
- return resp.cleared
- async def admin_delete_user(
- self, user_id: str, *, current_password: str = ""
- ) -> bool:
- resp = await self._rpc(
- self._svc.admin_users.delete_user(
- admin_members_pb2.DeleteUserRequest(
- user_id=user_id, current_password=current_password
- ),
- headers=self._headers(),
- )
- )
- return resp.deleted
- # --- Admin: roles --------------------------------------------------
- async def admin_list_roles(self) -> list[AdminRole]:
- resp = await self._rpc(
- self._svc.admin_roles.list_roles(
- admin_roles_pb2.ListRolesRequest(), headers=self._headers()
- )
- )
- data = pb_to_dict(resp)
- return [
- r
- for r in (AdminRole.parse(row) for row in data.get("roles") or [])
- if r is not None
- ]
- async def admin_get_role(self, name: str) -> AdminRoleDetail:
- resp = await self._rpc(
- self._svc.admin_roles.get_role(
- admin_roles_pb2.GetRoleRequest(name=name), headers=self._headers()
- )
- )
- return AdminRoleDetail.parse(pb_to_dict(resp))
- async def admin_create_role(
- self,
- name: str,
- *,
- display_name: str = "",
- description: str = "",
- pingable: bool = False,
- ) -> AdminRole | None:
- resp = await self._rpc(
- self._svc.admin_roles.create_role(
- admin_roles_pb2.CreateRoleRequest(
- name=name,
- display_name=display_name,
- description=description,
- pingable=pingable,
- ),
- headers=self._headers(),
- )
- )
- return AdminRole.parse(pb_to_dict(resp.role))
- async def admin_update_role(
- self,
- name: str,
- *,
- display_name: str | None = None,
- description: str | None = None,
- pingable: bool | None = None,
- ) -> AdminRole | None:
- req = admin_roles_pb2.UpdateRoleRequest(name=name)
- if display_name is not None:
- req.display_name = display_name
- if description is not None:
- req.description = description
- if pingable is not None:
- req.pingable = pingable
- resp = await self._rpc(
- self._svc.admin_roles.update_role(req, headers=self._headers())
- )
- return AdminRole.parse(pb_to_dict(resp.role))
- async def admin_delete_role(self, name: str) -> bool:
- resp = await self._rpc(
- self._svc.admin_roles.delete_role(
- admin_roles_pb2.DeleteRoleRequest(name=name),
- headers=self._headers(),
- )
- )
- return resp.deleted
- async def admin_reorder_roles(self, role_names: list[str]) -> list[AdminRole]:
- resp = await self._rpc(
- self._svc.admin_roles.reorder_roles(
- admin_roles_pb2.ReorderRolesRequest(role_names=role_names),
- headers=self._headers(),
- )
- )
- data = pb_to_dict(resp)
- return [
- r
- for r in (AdminRole.parse(row) for row in data.get("roles") or [])
- if r is not None
- ]
- # --- Admin: event log / diagnostics / permissions ----------------
- async def admin_list_events(
- self,
- *,
- limit: int | None = None,
- before: str | None = None,
- event_type: str | None = None,
- actor_id: str | None = None,
- ) -> AdminEventLogPage:
- """List durable EVT entries newest-first.
- Args:
- limit: Maximum entries to return (server clamps to its diagnostic
- limit).
- before: Exclusive sequence cursor; returned entries are older than
- this sequence.
- event_type: Filter by event payload type name (e.g.
- ``"MessagePostedEvent"``).
- actor_id: Filter by actor user ID.
- """
- req = event_log_pb2.ListEventsRequest()
- if limit is not None:
- req.limit = limit
- if before is not None:
- req.before = before
- if event_type or actor_id:
- fltr = event_log_pb2.AdminEventLogFilter()
- if event_type:
- fltr.event_type = event_type
- if actor_id:
- fltr.actor_id = actor_id
- req.filter.CopyFrom(fltr)
- resp = await self._rpc(
- self._svc.admin_event_log.list_events(req, headers=self._headers())
- )
- return AdminEventLogPage.parse(pb_to_dict(resp))
- async def admin_list_event_types(self) -> list[str]:
- resp = await self._rpc(
- self._svc.admin_event_log.list_event_types(
- event_log_pb2.ListEventTypesRequest(),
- headers=self._headers(),
- )
- )
- return list(resp.event_types)
- async def admin_get_event(self, sequence: str) -> AdminEventLogEntry | None:
- """Read one durable EVT entry by its stream sequence.
- Args:
- sequence: EVT stream sequence string (e.g. ``"42"``).
- Returns:
- The matching ``AdminEventLogEntry``, or ``None`` when the proto
- ``entry`` field is absent.
- """
- resp = await self._rpc(
- self._svc.admin_event_log.get_event(
- event_log_pb2.GetEventRequest(sequence=sequence),
- headers=self._headers(),
- )
- )
- return AdminEventLogEntry.parse(pb_to_dict(resp).get("entry"))
- async def admin_get_system_info(self) -> AdminSystemInfoSnapshot:
- from chattolib._pb.chatto.admin.v1 import diagnostics_pb2
- resp = await self._rpc(
- self._svc.admin_diagnostics.get_system_info(
- diagnostics_pb2.GetSystemInfoRequest(),
- headers=self._headers(),
- )
- )
- return AdminSystemInfoSnapshot.parse(pb_to_dict(resp))
- async def admin_get_role_permission_matrix(
- self, role_name: str
- ) -> RolePermissionMatrix | None:
- resp = await self._rpc(
- self._svc.admin_permissions.get_role_permission_matrix(
- admin_permissions_pb2.GetRolePermissionMatrixRequest(
- role_name=role_name
- ),
- headers=self._headers(),
- )
- )
- return RolePermissionMatrix.parse(pb_to_dict(resp).get("matrix"))
- async def admin_list_role_permission_decisions(
- self, role_name: str
- ) -> RolePermissionDecisions:
- resp = await self._rpc(
- self._svc.admin_permissions.list_role_permission_decisions(
- admin_permissions_pb2.ListRolePermissionDecisionsRequest(
- role_name=role_name
- ),
- headers=self._headers(),
- )
- )
- return RolePermissionDecisions.parse(pb_to_dict(resp))
- async def admin_get_user_permission_matrix(
- self, user_id: str
- ) -> UserPermissionMatrix | None:
- resp = await self._rpc(
- self._svc.admin_permissions.get_user_permission_matrix(
- admin_permissions_pb2.GetUserPermissionMatrixRequest(user_id=user_id),
- headers=self._headers(),
- )
- )
- return UserPermissionMatrix.parse(pb_to_dict(resp).get("matrix"))
- async def admin_list_user_permission_decisions(
- self, user_id: str
- ) -> UserPermissionDecisions:
- resp = await self._rpc(
- self._svc.admin_permissions.list_user_permission_decisions(
- admin_permissions_pb2.ListUserPermissionDecisionsRequest(
- user_id=user_id
- ),
- headers=self._headers(),
- )
- )
- return UserPermissionDecisions.parse(pb_to_dict(resp))
|