Skip to main content

tuwunel_service/rooms/timeline/
pdus.rs

1//! Streams, indexes, and deletes accepted room timeline rows.
2//!
3//! Directional scans use encoded room-local counts, while a secondary index
4//! supports timestamp lookup. Decoded streams apply requester-specific event
5//! presentation without changing which rows are selected.
6
7use 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
23/// Standard item yielded by decoded room timeline streams.
24///
25/// The count is the event's room-local pagination token. The producer
26/// determines whether the decoded PDU receives presentation adjustments.
27pub type PdusIterItem = (PduCount, PduEvent);
28
29/// Converts a signed PDU-count encoding into ordered offset-binary form.
30///
31/// The transformation preserves signed numeric ordering in unsigned database
32/// keys, placing negative backfilled counts before positive normal counts.
33#[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/// Deletes every accepted timeline row belonging to a room.
41///
42/// Each event's accepted row, ID mapping, outlier copy, and timestamp index are
43/// removed in one transaction. An error stops the scan after any prior events
44/// were already removed, and metadata owned by other services is not purged.
45#[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(&current)
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/// Streams decoded room events from the timestamp index in one direction.
80///
81/// The timestamp boundary is inclusive. `user_id` controls sender-only
82/// transaction metadata, while event age is always updated. It does not filter
83/// which events are yielded.
84#[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/// Streams timestamp and PDU-ID pairs from the room timestamp index.
103///
104/// Forward scans begin at the first row at or after the timestamp, while
105/// backward scans begin at the first row at or before it. Only rows for the
106/// requested room are yielded.
107#[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/// Streams all accepted PDUs in a room in forward order.
149///
150/// The requesting user receives sender-only transaction metadata where
151/// applicable. Unknown rooms and all per-item storage or decoding errors are
152/// suppressed, producing only successfully presented events.
153#[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/// Streams accepted room events after an optional count in forward order.
165///
166/// The count boundary is exclusive and defaults to the minimum count. The
167/// optional user controls presentation only; stream, storage, and decoding
168/// errors remain visible to the caller.
169#[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(&current)
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/// Streams accepted room events before an optional count in reverse order.
191///
192/// The count boundary is exclusive and defaults to the maximum count. The
193/// optional user controls presentation only; stream, storage, and decoding
194/// errors remain visible to the caller.
195#[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(&current)
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/// Streams raw JSON values from all accepted timeline rows.
217///
218/// Values borrow the database cursor and must be owned before they are retained
219/// across another poll. Storage errors remain in the stream.
220#[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/// Streams raw JSON values from all outlier rows.
226///
227/// Values borrow the database cursor and must be owned before they are retained
228/// across another poll. Storage errors remain in the stream.
229#[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}