tuwunel_service/rooms/retention/
mod.rs1use std::{sync::Arc, time::Duration};
8
9use async_trait::async_trait;
10use futures::{Stream, TryStreamExt};
11use ruma::{CanonicalJsonObject, EventId};
12use tuwunel_core::{
13 Result, debug_info, expected, implement,
14 matrix::pdu::PduEvent,
15 utils::{TryReadyExt, time::now},
16};
17use tuwunel_database::{Deserialized, Json, Map};
18
19use crate::rooms::timeline::RoomMutexGuard;
20
21#[cfg(test)]
22mod tests;
23
24pub struct Service {
30 services: Arc<crate::services::OnceServices>,
31 eventid_originalpdu: Arc<Map>,
32 timeredacted_eventid: Arc<Map>,
33}
34
35#[async_trait]
36impl crate::Service for Service {
37 fn build(args: &crate::Args<'_>) -> Result<Arc<Self>> {
38 Ok(Arc::new(Self {
39 services: args.services.clone(),
40 eventid_originalpdu: args.db["eventid_originalpdu"].clone(),
41 timeredacted_eventid: args.db["timeredacted_eventid"].clone(),
42 }))
43 }
44
45 async fn worker(self: Arc<Self>) -> Result {
46 loop {
47 let retention_seconds = self.services.config.redaction_retention_seconds;
48
49 if retention_seconds != 0 {
50 debug_info!("Cleaning up retained events");
51
52 let now = now().as_secs();
53 let count = self
54 .timeredacted_eventid
55 .keys::<(u64, &EventId)>()
56 .ready_try_take_while(|(time_redacted, _)| {
57 let time_redacted = *time_redacted;
58 Ok(expected!(time_redacted + retention_seconds) < now)
59 })
60 .ready_try_fold_default(|count: usize, (time_redacted, event_id)| {
61 self.eventid_originalpdu.remove(event_id);
62 self.timeredacted_eventid
63 .del((time_redacted, event_id));
64 Ok(count.saturating_add(1))
65 })
66 .await?;
67
68 debug_info!(?count, "Finished cleaning up retained events");
69 }
70
71 tokio::select! {
72 () = tokio::time::sleep(Duration::from_hours(1)) => {},
73 () = self.services.server.until_shutdown() => return Ok(())
74 };
75 }
76 }
77
78 fn name(&self) -> &str { crate::service::make_name(std::module_path!()) }
79}
80
81#[implement(Service)]
85pub async fn get_original_pdu(&self, event_id: &EventId) -> Result<PduEvent> {
86 self.eventid_originalpdu
87 .get(event_id)
88 .await?
89 .deserialized()
90}
91
92#[implement(Service)]
96pub async fn get_original_pdu_json(&self, event_id: &EventId) -> Result<CanonicalJsonObject> {
97 self.eventid_originalpdu
98 .get(event_id)
99 .await?
100 .deserialized()
101}
102
103#[implement(Service)]
110pub async fn save_original_pdu(
111 &self,
112 event_id: &EventId,
113 pdu: &CanonicalJsonObject,
114 _state_lock: &RoomMutexGuard,
115) {
116 if !self.services.config.save_unredacted_events {
117 return;
118 }
119
120 if self
121 .eventid_originalpdu
122 .exists(event_id)
123 .await
124 .is_ok()
125 {
126 return;
127 }
128
129 self.insert_original(event_id, pdu, now().as_secs());
130}
131
132#[implement(Service)]
136fn insert_original(&self, event_id: &EventId, pdu: &CanonicalJsonObject, time: u64) {
137 let mut txn = self.services.db.txn();
138
139 txn.raw_put(&self.eventid_originalpdu, event_id, Json(pdu));
140 txn.put_raw(&self.timeredacted_eventid, (time, event_id), []);
141 txn.execute();
142}
143
144#[implement(Service)]
149pub fn retained_pdus_raw(&self) -> impl Stream<Item = Result<&[u8]>> + Send {
150 self.eventid_originalpdu
151 .raw_stream()
152 .map_ok(|x| x.1)
153}
154
155#[implement(Service)]
160pub fn purge_original(&self, event_id: &EventId) { self.eventid_originalpdu.remove(event_id); }