tuwunel_service/rooms/pdu_metadata/
purge.rs1use 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#[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#[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}