Skip to main content

tuwunel_service/rooms/timeline/
purge.rs

1use futures::TryStreamExt;
2use ruma::{RoomId, api::Direction, events::TimelineEventType};
3use tuwunel_core::{
4	Result, implement,
5	matrix::{
6		Event,
7		pdu::{PduCount, PduEvent},
8	},
9	trace,
10	utils::stream::TryReadyExt,
11};
12
13use super::{ExtractBody, RawPduId, bias_count};
14
15/// Selectively purges room history strictly before `until` in stream order,
16/// returning the number of events removed. State events are always preserved,
17/// and locally-sent events are kept unless `delete_local_events`. Forward
18/// extremities are never touched, so the room stays live. tuwunel orders by
19/// per-room stream position rather than topological depth, so "same-depth
20/// events retained" becomes "strictly earlier in stream order".
21#[implement(super::Service)]
22pub async fn purge_history(
23	&self,
24	room_id: &RoomId,
25	until: PduCount,
26	delete_local_events: bool,
27) -> Result<usize> {
28	let shortroomid = self
29		.services
30		.short
31		.get_shortroomid(room_id)
32		.await?;
33
34	let start = self
35		.count_to_id(room_id, PduCount::min(), Direction::Forward)
36		.await?;
37
38	let prefix = start.shortroomid();
39
40	self.db
41		.pduid_pdu
42		.raw_stream_from(&start)
43		.ready_try_take_while(move |kv| {
44			let (key, _) = *kv;
45			Ok(key.starts_with(&prefix) && RawPduId::from(key).pdu_count() < until)
46		})
47		.try_fold(0_usize, async |purged, (key, value)| {
48			let pdu = serde_json::from_slice::<PduEvent>(value)?;
49
50			if pdu.state_key.is_some()
51				|| (!delete_local_events && self.services.globals.user_is_local(&pdu.sender))
52			{
53				return Ok(purged);
54			}
55
56			let mut txn = self.db.db.txn();
57
58			let raw_id = RawPduId::from(key);
59			let count = raw_id.pdu_count();
60			let event_id = pdu.event_id.clone();
61			let ts: u64 = pdu.origin_server_ts.into();
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_id = (room_id, ts, bias_count(raw_id.count()));
68			txn.del(&self.db.roomid_tscount_pducount, room_id_ts_id);
69
70			txn.execute();
71
72			if pdu.kind == TimelineEventType::RoomMessage
73				&& let Ok(ExtractBody { body: Some(body) }) = pdu.get_content()
74			{
75				self.services
76					.search
77					.deindex_pdu(shortroomid, &raw_id, &body);
78			}
79
80			self.services
81				.pdu_metadata
82				.purge_event_relations(shortroomid, count, room_id, &event_id)
83				.await;
84
85			self.services.retention.purge_original(&event_id);
86
87			trace!(?event_id, ?room_id, "Purged");
88
89			Ok(purged.saturating_add(1))
90		})
91		.await
92}