Skip to main content

tuwunel_service/rooms/state_accessor/
room_state.rs

1//! Resolves the current state snapshot for a room.
2//!
3//! These adapters obtain a room's current short-state hash and delegate to the
4//! historical snapshot readers. Stream variants surface a missing room snapshot
5//! while retaining the delegated readers' best-effort item behavior.
6
7use 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/// Deserializes one current state event's content.
16///
17/// The event is selected by `(event_type, state_key)`. Snapshot lookup,
18/// timeline lookup, and content errors are returned to the caller.
19#[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/// Streams current state events of one type.
35///
36/// Failure to resolve the room's current snapshot is yielded as an error.
37/// Missing reverse mappings and unavailable PDUs are skipped by the delegated
38/// best-effort stream.
39#[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/// Streams the room's full current state with type and state keys.
58///
59/// Failure to resolve the current snapshot is yielded as an error. Entries
60/// whose IDs, PDUs, or state keys cannot be resolved are skipped.
61#[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/// Streams every resolvable PDU in the room's current state.
76///
77/// Failure to resolve the current snapshot is yielded as an error. Individual
78/// state entries with missing reverse mappings or PDUs are skipped.
79#[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/// Returns the event ID for one current state tuple.
94///
95/// The room's current snapshot must exist, and both short-ID mappings must
96/// resolve for `(event_type, state_key)`.
97#[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/// Streams state keys and event IDs for one current state event type.
113///
114/// Failure to resolve the current snapshot is yielded as an error. Individual
115/// short-ID mapping failures are omitted by the best-effort inner stream.
116#[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/// Streams state keys for one current state event type.
135///
136/// Failure to resolve the current snapshot is yielded as an error. Individual
137/// state-key mapping failures are omitted by the best-effort inner stream.
138#[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/// Returns one current state PDU.
157///
158/// The event is selected by `(event_type, state_key)`. Snapshot, short-ID, and
159/// timeline lookup failures are returned to the caller.
160#[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}