tuwunel_service/rooms/event_handler/prev_walk/
history.rs1use 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
17pub(super) type PassKey<'a> = (&'a RoomId, u64, &'a EventId);
22
23pub(super) type PassVal<'a> = (u8, u64, u64, u8, u64, u64, &'a str);
29
30pub(super) type PassRow<'a> = (PassKey<'a>, PassVal<'a>);
31
32#[derive(Clone, Debug, Eq, PartialEq)]
38pub struct PrevWalkRoom {
39 pub room_id: OwnedRoomId,
41
42 pub passes: u64,
44
45 pub appended: u64,
47
48 pub not_appended: u64,
50
51 pub failed: u64,
53
54 pub cancelled: u64,
56
57 pub fetch_failed: u64,
59
60 pub fetch_cancelled: u64,
62
63 pub capped: u64,
65
66 pub prevs: u64,
68
69 pub unprocessed: u64,
72
73 pub fetch: Duration,
75
76 pub upgrade: Duration,
78}
79
80#[derive(Clone, Debug, Eq, PartialEq)]
86pub struct PrevWalkPass {
87 pub ended: SystemTime,
89
90 pub event_id: OwnedEventId,
92
93 pub origin: Option<OwnedServerName>,
95
96 pub outcome: Option<Outcome>,
98
99 pub prevs: u64,
101
102 pub unprocessed: u64,
105
106 pub capped: bool,
108
109 pub fetch: Duration,
112
113 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#[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#[implement(super::super::Service)]
181#[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
189pub(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#[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
226pub(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#[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
299pub(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}