Skip to main content

tuwunel_service/rooms/retention/
mod.rs

1//! Retains original event JSON when accepted timeline events are redacted.
2//!
3//! Originals are indexed by redaction time for periodic expiry. Saving and
4//! scheduled cleanup are controlled independently so operators can retain
5//! originals indefinitely when expiry is disabled.
6
7use 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
24/// Stores and expires the unredacted originals of redacted events.
25///
26/// Each retained event has a primary row and a time-ordered expiry index row.
27/// The background worker removes entries once their configured retention
28/// interval has elapsed.
29pub 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/// Returns a retained original decoded as a PDU.
82///
83/// Missing rows and malformed stored JSON are reported to the caller.
84#[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/// Returns a retained original as canonical event JSON.
93///
94/// Missing rows and malformed stored JSON are reported to the caller.
95#[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/// Retains an event's original JSON before it is redacted.
104///
105/// Saving is skipped when saving originals is disabled or an original already
106/// exists.
107/// The primary row and its redaction-time expiry index are committed in one
108/// transaction while the caller retains the room timeline guard.
109#[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/// Writes a retained original and its expiry index atomically.
133///
134/// The supplied timestamp is the redaction time used by the cleanup worker.
135#[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/// Streams the raw JSON values of all retained originals.
145///
146/// Values borrow the database cursor and must be owned before they are retained
147/// across another poll. Storage errors remain in the stream.
148#[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/// Drops the retained original of a purged event.
156///
157/// The paired expiry-index row is deliberately left for the retention worker
158/// to reap at its scheduled time.
159#[implement(Service)]
160pub fn purge_original(&self, event_id: &EventId) { self.eventid_originalpdu.remove(event_id); }