tuwunel_service/rooms/state_accessor/
room_state.rs1use futures::{Stream, StreamExt, TryFutureExt};
8use ruma::{OwnedEventId, RoomId, events::StateEventType};
9use serde::Deserialize;
10use tuwunel_core::{
11 Result, err, implement,
12 matrix::{Event, Pdu, StateKey},
13};
14
15#[implement(super::Service)]
20pub async fn room_state_get_content<T>(
21 &self,
22 room_id: &RoomId,
23 event_type: &StateEventType,
24 state_key: &str,
25) -> Result<T>
26where
27 T: for<'de> Deserialize<'de> + Send,
28{
29 self.room_state_get(room_id, event_type, state_key)
30 .await
31 .and_then(|event| event.get_content())
32}
33
34#[implement(super::Service)]
40#[tracing::instrument(skip(self), level = "debug")]
41pub fn room_state_type_pdus<'a>(
42 &'a self,
43 room_id: &'a RoomId,
44 event_type: &'a StateEventType,
45) -> impl Stream<Item = Result<impl Event>> + Send + 'a {
46 self.services
47 .state
48 .get_room_shortstatehash(room_id)
49 .map_ok(|shortstatehash| {
50 self.state_type_pdus(shortstatehash, event_type)
51 .map(Ok)
52 })
53 .map_err(move |e| err!(Database("Missing state for {room_id:?}: {e:?}")))
54 .try_flatten_stream()
55}
56
57#[implement(super::Service)]
62#[tracing::instrument(skip(self), level = "debug")]
63pub fn room_state_full<'a>(
64 &'a self,
65 room_id: &'a RoomId,
66) -> impl Stream<Item = Result<((StateEventType, StateKey), impl Event)>> + Send + 'a {
67 self.services
68 .state
69 .get_room_shortstatehash(room_id)
70 .map_ok(|shortstatehash| self.state_full(shortstatehash).map(Ok))
71 .map_err(move |e| err!(Database("Missing state for {room_id:?}: {e:?}")))
72 .try_flatten_stream()
73}
74
75#[implement(super::Service)]
80#[tracing::instrument(skip(self), level = "debug")]
81pub fn room_state_full_pdus<'a>(
82 &'a self,
83 room_id: &'a RoomId,
84) -> impl Stream<Item = Result<impl Event>> + Send + 'a {
85 self.services
86 .state
87 .get_room_shortstatehash(room_id)
88 .map_ok(|shortstatehash| self.state_full_pdus(shortstatehash).map(Ok))
89 .map_err(move |e| err!(Database("Missing state for {room_id:?}: {e:?}")))
90 .try_flatten_stream()
91}
92
93#[implement(super::Service)]
98#[tracing::instrument(skip(self), level = "debug")]
99pub async fn room_state_get_id(
100 &self,
101 room_id: &RoomId,
102 event_type: &StateEventType,
103 state_key: &str,
104) -> Result<OwnedEventId> {
105 self.services
106 .state
107 .get_room_shortstatehash(room_id)
108 .and_then(|shortstatehash| self.state_get_id(shortstatehash, event_type, state_key))
109 .await
110}
111
112#[implement(super::Service)]
117#[tracing::instrument(skip(self), level = "debug")]
118pub fn room_state_keys_with_ids<'a>(
119 &'a self,
120 room_id: &'a RoomId,
121 event_type: &'a StateEventType,
122) -> impl Stream<Item = Result<(StateKey, OwnedEventId)>> + Send + 'a {
123 self.services
124 .state
125 .get_room_shortstatehash(room_id)
126 .map_ok(|shortstatehash| {
127 self.state_keys_with_ids(shortstatehash, event_type)
128 .map(Ok)
129 })
130 .map_err(move |e| err!(Database("Missing state for {room_id:?}: {e:?}")))
131 .try_flatten_stream()
132}
133
134#[implement(super::Service)]
139#[tracing::instrument(skip(self), level = "debug")]
140pub fn room_state_keys<'a>(
141 &'a self,
142 room_id: &'a RoomId,
143 event_type: &'a StateEventType,
144) -> impl Stream<Item = Result<StateKey>> + Send + 'a {
145 self.services
146 .state
147 .get_room_shortstatehash(room_id)
148 .map_ok(|shortstatehash| {
149 self.state_keys(shortstatehash, event_type)
150 .map(Ok)
151 })
152 .map_err(move |e| err!(Database("Missing state for {room_id:?}: {e:?}")))
153 .try_flatten_stream()
154}
155
156#[implement(super::Service)]
161#[tracing::instrument(skip(self), level = "debug")]
162pub async fn room_state_get(
163 &self,
164 room_id: &RoomId,
165 event_type: &StateEventType,
166 state_key: &str,
167) -> Result<Pdu> {
168 self.services
169 .state
170 .get_room_shortstatehash(room_id)
171 .and_then(|shortstatehash| self.state_get(shortstatehash, event_type, state_key))
172 .await
173}