tuwunel_service/rooms/timeline/
pdus.rs1use futures::{Stream, StreamExt, TryFutureExt, TryStreamExt};
8use ruma::{MilliSecondsSinceUnixEpoch, RoomId, UInt, UserId, api::Direction};
9use tuwunel_core::{
10 Result, at, err, implement,
11 matrix::pdu::{PduCount, PduEvent},
12 trace,
13 utils::{
14 result::LogErr,
15 stream::{TryIgnore, TryReadyExt, TryWidebandExt},
16 },
17 warn,
18};
19use tuwunel_database::{KeyVal, keyval::Val};
20
21use super::{PduId, RawPduId};
22
23pub type PdusIterItem = (PduCount, PduEvent);
28
29#[must_use]
34pub fn bias_count(count: [u8; 8]) -> u64 {
35 i64::from_be_bytes(count)
36 .wrapping_sub(i64::MIN)
37 .cast_unsigned()
38}
39
40#[implement(super::Service)]
46pub async fn delete_pdus(&self, room_id: &RoomId) -> Result {
47 let current = self
48 .count_to_id(room_id, PduCount::min(), Direction::Forward)
49 .await?;
50
51 let prefix = current.shortroomid();
52 self.db
53 .pduid_pdu
54 .raw_stream_from(¤t)
55 .ready_try_take_while(move |(key, _)| Ok(key.starts_with(&prefix)))
56 .ready_try_for_each(move |(key, value)| {
57 let pdu = serde_json::from_slice::<PduEvent>(value)?;
58 let ts: u64 = pdu.origin_server_ts.into();
59 let event_id = &pdu.event_id;
60
61 let mut txn = self.db.db.txn();
62
63 txn.del_raw(&self.db.pduid_pdu, key);
64 txn.del_raw(&self.db.eventid_pduid, event_id);
65 txn.del_raw(&self.db.eventid_outlierpdu, event_id);
66
67 let room_id_ts_key = (room_id, ts, bias_count(RawPduId::from(key).count()));
68 txn.del(&self.db.roomid_tscount_pducount, room_id_ts_key);
69
70 txn.execute();
71
72 trace!(?event_id, ?room_id, ?ts, ?key, "Removed");
73
74 Ok(())
75 })
76 .await
77}
78
79#[implement(super::Service)]
85pub fn pdus_near_ts(
86 &self,
87 user_id: Option<&UserId>,
88 room_id: &RoomId,
89 ts: MilliSecondsSinceUnixEpoch,
90 dir: Direction,
91) -> impl Stream<Item = Result<PdusIterItem>> + Send {
92 self.pdu_ids_near_ts(room_id, ts, dir)
93 .map_ok(|(ts, pdu_id)| (ts, pdu_id.into()))
94 .wide_and_then(async |(_, pdu_id): (_, RawPduId)| {
95 self.get_pdu_from_id(&pdu_id)
96 .map_ok(|pdu| (pdu_id, pdu))
97 .await
98 })
99 .ready_and_then(move |item| Self::each_pdu(item, user_id))
100}
101
102#[implement(super::Service)]
108pub fn pdu_ids_near_ts(
109 &self,
110 room_id: &RoomId,
111 ts: MilliSecondsSinceUnixEpoch,
112 dir: Direction,
113) -> impl Stream<Item = Result<(MilliSecondsSinceUnixEpoch, PduId)>> + Send {
114 use Direction::{Backward, Forward};
115
116 type KeyVal<'a> = ((&'a RoomId, UInt, u64), i64);
117
118 let ts: u64 = ts.get().into();
119
120 self.services
121 .short
122 .get_shortroomid(room_id)
123 .map_err(|e| err!(Request(NotFound("Room not found: {e:?}"))))
124 .map_ok(move |shortroomid| {
125 match dir {
126 | Forward => self
127 .db
128 .roomid_tscount_pducount
129 .stream_from(&(room_id, ts, u64::MIN))
130 .left_stream(),
131
132 | Backward => self
133 .db
134 .roomid_tscount_pducount
135 .rev_stream_from(&(room_id, ts, u64::MAX))
136 .right_stream(),
137 }
138 .ready_try_take_while(
139 move |((room_id_, ..), _): &KeyVal<'_>| Ok(room_id == *room_id_),
140 )
141 .map_ok(move |((_, ts, _), count)| {
142 (MilliSecondsSinceUnixEpoch(ts), PduId { shortroomid, count: count.into() })
143 })
144 })
145 .try_flatten_stream()
146}
147
148#[implement(super::Service)]
154#[inline]
155pub fn all_pdus<'a>(
156 &'a self,
157 user_id: &'a UserId,
158 room_id: &'a RoomId,
159) -> impl Stream<Item = PdusIterItem> + Send + 'a {
160 self.pdus(Some(user_id), room_id, None)
161 .ignore_err()
162}
163
164#[implement(super::Service)]
170#[tracing::instrument(skip(self), level = "debug")]
171pub fn pdus<'a>(
172 &'a self,
173 user_id: Option<&'a UserId>,
174 room_id: &'a RoomId,
175 from: Option<PduCount>,
176) -> impl Stream<Item = Result<PdusIterItem>> + Send + 'a {
177 let from = from.unwrap_or_else(PduCount::min);
178 self.count_to_id(room_id, from, Direction::Forward)
179 .map_ok(move |current| {
180 let prefix = current.shortroomid();
181 self.db
182 .pduid_pdu
183 .raw_stream_from(¤t)
184 .ready_try_take_while(move |(key, _)| Ok(key.starts_with(&prefix)))
185 .ready_and_then(move |item| Self::each_slice(item, user_id))
186 })
187 .try_flatten_stream()
188}
189
190#[implement(super::Service)]
196#[tracing::instrument(skip(self), level = "debug")]
197pub fn pdus_rev<'a>(
198 &'a self,
199 user_id: Option<&'a UserId>,
200 room_id: &'a RoomId,
201 until: Option<PduCount>,
202) -> impl Stream<Item = Result<PdusIterItem>> + Send + 'a {
203 let until = until.unwrap_or_else(PduCount::max);
204 self.count_to_id(room_id, until, Direction::Backward)
205 .map_ok(move |current| {
206 let prefix = current.shortroomid();
207 self.db
208 .pduid_pdu
209 .rev_raw_stream_from(¤t)
210 .ready_try_take_while(move |(key, _)| Ok(key.starts_with(&prefix)))
211 .ready_and_then(move |item| Self::each_slice(item, user_id))
212 })
213 .try_flatten_stream()
214}
215
216#[implement(super::Service)]
221pub fn pdus_raw(&self) -> impl Stream<Item = Result<Val<'_>>> + Send {
222 self.db.pduid_pdu.raw_stream().map_ok(at!(1))
223}
224
225#[implement(super::Service)]
230pub fn outlier_pdus_raw(&self) -> impl Stream<Item = Result<Val<'_>>> + Send {
231 self.db
232 .eventid_outlierpdu
233 .raw_stream()
234 .map_ok(at!(1))
235}
236
237#[implement(super::Service)]
238fn each_slice((pdu_id, pdu): KeyVal<'_>, user_id: Option<&UserId>) -> Result<PdusIterItem> {
239 let pdu_id: RawPduId = pdu_id.into();
240 let pdu = serde_json::from_slice::<PduEvent>(pdu)?;
241
242 Self::each_pdu((pdu_id, pdu), user_id)
243}
244
245#[implement(super::Service)]
246fn each_pdu(
247 (pdu_id, mut pdu): (RawPduId, PduEvent),
248 user_id: Option<&UserId>,
249) -> Result<PdusIterItem> {
250 pdu.remove_transaction_id_unless_sender(user_id)?;
251 pdu.add_age().log_err().ok();
252
253 Ok((pdu_id.pdu_count(), pdu))
254}