Skip to main content

tuwunel_service/rooms/pdu_metadata/
purge.rs

1use futures::StreamExt;
2use ruma::{
3	EventId, RoomId,
4	events::{relation::RelationType, room::encrypted::Relation},
5};
6use tuwunel_core::{
7	PduId, Result,
8	arrayvec::ArrayVec,
9	implement,
10	matrix::{Event, Pdu, PduCount, RawPduId},
11	utils::{
12		stream::{ReadyExt, TryIgnore, automatic_width},
13		u64_from_u8,
14	},
15};
16
17use super::{ExtractRelatesTo, Service};
18use crate::rooms::short::ShortRoomId;
19
20type Prefix = ArrayVec<u8, 16>;
21
22/// Purges one event's metadata during a history purge.
23///
24/// Relation rows keyed by this event as parent or target are removed. Rows
25/// keyed by it as a surviving event's child remain dangling because relation
26/// reads discard IDs that no longer resolve. Soft-fail and policy decisions are
27/// cleared after the relation indexes.
28#[implement(Service)]
29pub async fn purge_event_relations(
30	&self,
31	shortroomid: ShortRoomId,
32	parent: PduCount,
33	room_id: &RoomId,
34	event_id: &EventId,
35) {
36	let target = parent.to_be_bytes();
37
38	self.db
39		.tofrom_relation
40		.raw_keys_from(target.as_slice())
41		.ignore_err()
42		.ready_take_while(move |key| key.starts_with(&target))
43		.ready_for_each(|key| self.db.tofrom_relation.remove(key))
44		.await;
45
46	let mut prefix = Prefix::new();
47
48	prefix.extend(shortroomid.to_be_bytes());
49	prefix.extend(parent.to_be_bytes());
50
51	self.db
52		.relatesto_typed
53		.raw_keys_from(prefix.as_slice())
54		.ignore_err()
55		.ready_take_while(move |key| key.starts_with(&prefix))
56		.ready_for_each(|key| self.db.relatesto_typed.remove(key))
57		.await;
58
59	self.db.referencedevents.del((room_id, event_id));
60
61	self.services
62		.event_handler
63		.clear_policy_signature_state(event_id);
64
65	self.clear_event_soft_failed(event_id);
66}
67
68/// Rebuild `relatesto_typed` from every stored PDU. Run once at startup behind
69/// a `global` marker, and on demand from the admin command. Clears first so a
70/// partial or stale index is replaced wholesale.
71#[implement(Service)]
72pub async fn rebuild_typed_relations(&self) -> Result {
73	self.db.relatesto_typed.clear().await;
74
75	let pdus = self.services.db["pduid_pdu"].clone();
76
77	pdus.raw_stream()
78		.ignore_err()
79		.ready_filter_map(|(key, value)| {
80			let raw_pdu_id = RawPduId::from(key);
81			let pdu_id = PduId {
82				shortroomid: u64_from_u8(&raw_pdu_id.shortroomid()),
83				count: raw_pdu_id.pdu_count(),
84			};
85			let pdu = serde_json::from_slice::<Pdu>(value).ok()?;
86
87			Some((pdu_id, pdu))
88		})
89		.for_each_concurrent(automatic_width(), async |(pdu_id, pdu)| {
90			self.index_pdu_relations(pdu_id, &pdu).await;
91		})
92		.await;
93
94	Ok(())
95}
96
97#[implement(Service)]
98async fn index_pdu_relations(&self, pdu_id: PduId, pdu: &Pdu) {
99	let Ok(content) = pdu.get_content::<ExtractRelatesTo>() else {
100		return;
101	};
102
103	let (rel_type, parent) = match content.relates_to {
104		| Relation::Replacement(replacement) => (RelationType::Replacement, replacement.event_id),
105		| Relation::Reference(reference) => (RelationType::Reference, reference.event_id),
106		| _ => return,
107	};
108
109	self.add_typed_relation(pdu_id.shortroomid, pdu_id.count, &parent, pdu, rel_type)
110		.await;
111}