Skip to main content

tuwunel_service/rooms/event_handler/prev_walk/
history.rs

1use std::time::{Duration, SystemTime, UNIX_EPOCH};
2
3use futures::{Stream, StreamExt};
4use ruma::{EventId, OwnedEventId, OwnedRoomId, OwnedServerName, RoomId, ServerName};
5use tuwunel_core::{
6	implement,
7	utils::{
8		math::{u64_from_u128_saturating, u64_from_usize_saturating},
9		stream::{ReadyExt, TryIgnore},
10		time::{now_millis, timepoint_from_epoch},
11	},
12};
13use tuwunel_database::Interfix;
14
15use super::{Outcome, Pass, PrevUpgrade};
16
17/// Key of a pass row.
18///
19/// In order: the room, the wall clock in milliseconds since the Unix epoch when
20/// the pass ended, and the incoming event. Rows sort by room, then by end.
21pub(super) type PassKey<'a> = (&'a RoomId, u64, &'a EventId);
22
23/// Value of a pass row.
24///
25/// In order: the outcome code, the collected prevs, the unprocessed prevs,
26/// whether the fetch was capped as 0 or 1, the fetch and upgrade milliseconds,
27/// and the origin.
28pub(super) type PassVal<'a> = (u8, u64, u64, u8, u64, u64, &'a str);
29
30pub(super) type PassRow<'a> = (PassKey<'a>, PassVal<'a>);
31
32/// Recorded prev walk totals for one room.
33///
34/// Every recorded pass counts in `passes`. A pass whose outcome this build
35/// does not know, written by a newer one, counts there and in the other totals
36/// but in none of the outcome counts.
37#[derive(Clone, Debug, Eq, PartialEq)]
38pub struct PrevWalkRoom {
39	/// The room the passes ran in.
40	pub room_id: OwnedRoomId,
41
42	/// Recorded passes, whatever their outcome.
43	pub passes: u64,
44
45	/// Passes that appended the incoming event.
46	pub appended: u64,
47
48	/// Passes that finished without appending the incoming event.
49	pub not_appended: u64,
50
51	/// Passes whose walk returned an error.
52	pub failed: u64,
53
54	/// Passes whose walk was dropped or interrupted.
55	pub cancelled: u64,
56
57	/// Passes whose backward fetch returned an error.
58	pub fetch_failed: u64,
59
60	/// Passes dropped or interrupted before their walk began.
61	pub fetch_cancelled: u64,
62
63	/// Passes whose backward fetch hit the `max_fetch_prev_events` cap.
64	pub capped: u64,
65
66	/// Previous events the passes collected for upgrade.
67	pub prevs: u64,
68
69	/// Collected previous events left without an upgrade by passes that were
70	/// not cancelled.
71	pub unprocessed: u64,
72
73	/// Time the passes spent before their walks began.
74	pub fetch: Duration,
75
76	/// Time the passes spent walking.
77	pub upgrade: Duration,
78}
79
80/// One recorded pass.
81///
82/// Decoded from one row of the room's history. The origin is absent when the
83/// recorded one does not parse, and the outcome when a newer build wrote a code
84/// this one does not know.
85#[derive(Clone, Debug, Eq, PartialEq)]
86pub struct PrevWalkPass {
87	/// When the pass ended, by the wall clock.
88	pub ended: SystemTime,
89
90	/// The gapped incoming event.
91	pub event_id: OwnedEventId,
92
93	/// Server the incoming event arrived from, absent when unparsable.
94	pub origin: Option<OwnedServerName>,
95
96	/// How the pass ended, absent when written by a newer build.
97	pub outcome: Option<Outcome>,
98
99	/// Previous events the pass collected for upgrade.
100	pub prevs: u64,
101
102	/// Collected previous events the pass left without an upgrade, none when it
103	/// was cancelled.
104	pub unprocessed: u64,
105
106	/// Whether the backward fetch hit the `max_fetch_prev_events` cap.
107	pub capped: bool,
108
109	/// Time before the walk began, or until the end for a pass that never
110	/// walked.
111	pub fetch: Duration,
112
113	/// Time spent walking.
114	pub upgrade: Duration,
115}
116
117impl From<Outcome> for u8 {
118	#[inline]
119	fn from(outcome: Outcome) -> Self {
120		match outcome {
121			| Outcome::Held => 0,
122			| Outcome::Closed => 1,
123			| Outcome::FetchFailed => 2,
124			| Outcome::FetchCancelled => 3,
125			| Outcome::Appended => 4,
126			| Outcome::NotAppended => 5,
127			| Outcome::Failed => 6,
128			| Outcome::Cancelled => 7,
129		}
130	}
131}
132
133impl TryFrom<u8> for Outcome {
134	type Error = u8;
135
136	#[inline]
137	fn try_from(code: u8) -> Result<Self, Self::Error> {
138		match code {
139			| 0 => Ok(Self::Held),
140			| 1 => Ok(Self::Closed),
141			| 2 => Ok(Self::FetchFailed),
142			| 3 => Ok(Self::FetchCancelled),
143			| 4 => Ok(Self::Appended),
144			| 5 => Ok(Self::NotAppended),
145			| 6 => Ok(Self::Failed),
146			| 7 => Ok(Self::Cancelled),
147			| unknown => Err(unknown),
148		}
149	}
150}
151
152/// Totals for a room with no recorded passes.
153///
154/// Every count and time is zero, the base the room's rows are counted onto.
155#[implement(PrevWalkRoom)]
156#[must_use]
157pub fn empty(room_id: &RoomId) -> Self {
158	Self {
159		room_id: room_id.to_owned(),
160		passes: 0,
161		appended: 0,
162		not_appended: 0,
163		failed: 0,
164		cancelled: 0,
165		fetch_failed: 0,
166		fetch_cancelled: 0,
167		capped: 0,
168		prevs: 0,
169		unprocessed: 0,
170		fetch: Duration::ZERO,
171		upgrade: Duration::ZERO,
172	}
173}
174
175/// Record a reported pass in its room's history.
176///
177/// The row is keyed by the wall clock at the write, and passes ending in the
178/// same millisecond stay apart by event. The write is synchronous and panics on
179/// a read-only database, like the backoff attempt recorded on the same path.
180#[implement(super::super::Service)]
181// keeps the key and value buffers out of the guard's drop frame
182#[inline(never)]
183pub(super) fn record_pass(&self, upgrade: &PrevUpgrade<'_>, pass: &Pass) {
184	let (key, val) = pass_row(upgrade, pass, now_millis());
185
186	self.db.roomtseventid_prevwalk.put(key, val);
187}
188
189/// Encode a pass as the row recording it.
190///
191/// The end is given in milliseconds since the Unix epoch: the wall clock when
192/// recording, a fixed time in a test.
193pub(super) fn pass_row<'a>(upgrade: &PrevUpgrade<'a>, pass: &Pass, ended_ms: u64) -> PassRow<'a> {
194	let PrevUpgrade { origin, room_id, event_id, .. } = *upgrade;
195	let val = (
196		u8::from(pass.outcome),
197		u64_from_usize_saturating(pass.prevs),
198		u64_from_usize_saturating(pass.unprocessed),
199		u8::from(pass.capped),
200		u64_from_u128_saturating(pass.fetch.as_millis()),
201		u64_from_u128_saturating(pass.upgrade.as_millis()),
202		origin.as_str(),
203	);
204
205	((room_id, ended_ms, event_id), val)
206}
207
208/// Read the recorded prev walk totals of every room, grouped by room in key
209/// order.
210///
211/// The whole history is swept once, so the cost grows with the passes it keeps,
212/// at most three days of them and fewer under load. A row that fails to decode
213/// is skipped, or fails an assertion in a debug build.
214#[implement(super::super::Service)]
215#[tracing::instrument(level = "debug", skip_all)]
216pub async fn prev_walk_rooms(&self) -> impl ExactSizeIterator<Item = PrevWalkRoom> + Send {
217	self.db
218		.roomtseventid_prevwalk
219		.stream()
220		.ignore_err()
221		.ready_fold(Vec::new(), tally)
222		.await
223		.into_iter()
224}
225
226/// Fold one row into the rooms read so far.
227///
228/// Rows arrive grouped by room, so a row either extends the last room or
229/// starts the next.
230pub(super) fn tally(
231	mut rooms: Vec<PrevWalkRoom>,
232	((room_id, ..), pass): PassRow<'_>,
233) -> Vec<PrevWalkRoom> {
234	match rooms.last_mut() {
235		| Some(room) if room.room_id == room_id => room.count(pass),
236		| _ => rooms.push(PrevWalkRoom::new(room_id, pass)),
237	}
238
239	rooms
240}
241
242#[implement(PrevWalkRoom)]
243fn new(room_id: &RoomId, pass: PassVal<'_>) -> Self {
244	let mut room = Self::empty(room_id);
245
246	room.count(pass);
247	room
248}
249
250#[implement(PrevWalkRoom)]
251fn count(&mut self, (code, prevs, unprocessed, capped, fetch_ms, upgrade_ms, _): PassVal<'_>) {
252	let counter = Outcome::try_from(code)
253		.ok()
254		.and_then(|outcome| match outcome {
255			| Outcome::Held | Outcome::Closed => None,
256			| Outcome::Appended => Some(&mut self.appended),
257			| Outcome::NotAppended => Some(&mut self.not_appended),
258			| Outcome::Failed => Some(&mut self.failed),
259			| Outcome::Cancelled => Some(&mut self.cancelled),
260			| Outcome::FetchFailed => Some(&mut self.fetch_failed),
261			| Outcome::FetchCancelled => Some(&mut self.fetch_cancelled),
262		});
263
264	if let Some(counter) = counter {
265		*counter = counter.saturating_add(1);
266	}
267
268	self.passes = self.passes.saturating_add(1);
269	self.capped = self.capped.saturating_add(u64::from(capped != 0));
270	self.prevs = self.prevs.saturating_add(prevs);
271	self.unprocessed = self.unprocessed.saturating_add(unprocessed);
272	self.fetch = self
273		.fetch
274		.saturating_add(Duration::from_millis(fetch_ms));
275
276	self.upgrade = self
277		.upgrade
278		.saturating_add(Duration::from_millis(upgrade_ms));
279}
280
281/// Stream the recorded passes of one room, latest first.
282///
283/// A row that fails to decode is skipped, or fails an assertion in a debug
284/// build.
285#[implement(super::super::Service)]
286pub fn prev_walk_passes<'a>(
287	&'a self,
288	room_id: &'a RoomId,
289) -> impl Stream<Item = PrevWalkPass> + Send + 'a {
290	let rows = self
291		.db
292		.roomtseventid_prevwalk
293		.rev_stream_from(&(room_id, u64::MAX, Interfix))
294		.ignore_err();
295
296	room_passes(rows, room_id)
297}
298
299/// Take the leading rows of one room from a stream of rows, as passes.
300///
301/// The stream ends at the first row of another room without reading past it.
302pub(super) fn room_passes<'a, S>(
303	rows: S,
304	room_id: &'a RoomId,
305) -> impl Stream<Item = PrevWalkPass> + Send + 'a
306where
307	S: Stream<Item = PassRow<'a>> + Send + 'a,
308{
309	rows.ready_take_while(move |((room, ..), _)| *room == room_id)
310		.map(PrevWalkPass::from_row)
311}
312
313#[implement(PrevWalkPass)]
314fn from_row(((_, ended_ms, event_id), val): PassRow<'_>) -> Self {
315	let (code, prevs, unprocessed, capped, fetch_ms, upgrade_ms, origin) = val;
316
317	Self {
318		ended: timepoint_from_epoch(Duration::from_millis(ended_ms)).unwrap_or(UNIX_EPOCH),
319		event_id: event_id.to_owned(),
320		origin: ServerName::parse(origin).ok(),
321		outcome: Outcome::try_from(code).ok(),
322		prevs,
323		unprocessed,
324		capped: capped != 0,
325		fetch: Duration::from_millis(fetch_ms),
326		upgrade: Duration::from_millis(upgrade_ms),
327	}
328}