1use std::{
2 collections::BTreeMap,
3 sync::{
4 Mutex, MutexGuard,
5 atomic::{AtomicU64, Ordering},
6 },
7 thread::panicking,
8 time::{Duration, Instant},
9};
10
11use ruma::{OwnedEventId, OwnedRoomId, OwnedServerName};
12use tuwunel_core::{
13 Error, Result, debug_info, implement,
14 utils::{MutexExt, math::fetch_add_usize},
15 warn,
16};
17
18pub use self::history::{PrevWalkPass, PrevWalkRoom};
19use super::{Service, fetch_prev::PrevFetch, handle_prev_pdu::PrevUpgrade};
20
21mod history;
22#[cfg(test)]
23mod tests;
24
25type Walks = BTreeMap<u64, InFlightWalk>;
26
27#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
35pub struct PrevWalkMetrics {
36 pub entered: u64,
39
40 pub gapped: u64,
42
43 pub held: u64,
45
46 pub closed: u64,
49
50 pub fetch_failed: u64,
52
53 pub fetch_cancelled: u64,
56
57 pub walked: u64,
59
60 pub walked_prevs: u64,
62
63 pub capped: u64,
65
66 pub appended: u64,
68
69 pub not_appended: u64,
71
72 pub failed: u64,
74
75 pub cancelled: u64,
77
78 pub unprocessed_prevs: u64,
81}
82
83#[derive(Clone, Debug)]
87pub struct InFlightWalk {
88 pub room_id: OwnedRoomId,
90
91 pub event_id: OwnedEventId,
93
94 pub origin: OwnedServerName,
96
97 pub started: Instant,
99
100 pub walk: Option<Walk>,
102}
103
104#[derive(Clone, Copy, Debug)]
110pub struct Walk {
111 pub fetched: Instant,
113
114 pub prevs: usize,
116
117 pub capped: bool,
119}
120
121#[derive(Default)]
122pub(super) struct PrevWalkCounters {
123 entered: AtomicU64,
124 gapped: AtomicU64,
125 held: AtomicU64,
126 closed: AtomicU64,
127 fetch_failed: AtomicU64,
128 fetch_cancelled: AtomicU64,
129 walked: AtomicU64,
130 walked_prevs: AtomicU64,
131 capped: AtomicU64,
132 appended: AtomicU64,
133 not_appended: AtomicU64,
134 failed: AtomicU64,
135 cancelled: AtomicU64,
136 unprocessed_prevs: AtomicU64,
137}
138
139#[derive(Default)]
146pub(super) struct InFlightWalks {
147 next: AtomicU64,
148 walks: Mutex<Walks>,
149}
150
151#[clippy::has_significant_drop]
160#[must_use]
161pub(super) struct PrevWalk<'a> {
162 service: &'a Service,
163 id: u64,
164 upgrade: &'a PrevUpgrade<'a>,
165 started: Instant,
166 walk: Option<Walk>,
167 settled: Option<(Outcome, usize)>,
168}
169
170#[derive(Clone, Copy)]
175struct Pass {
176 outcome: Outcome,
177 prevs: usize,
178 unprocessed: usize,
179 capped: bool,
180 fetch: Duration,
181 upgrade: Duration,
182}
183
184#[derive(Clone, Copy, Debug, Eq, PartialEq)]
189pub enum Outcome {
190 Held = 0,
192
193 Closed = 1,
195
196 FetchFailed = 2,
198
199 FetchCancelled = 3,
201
202 Appended = 4,
204
205 NotAppended = 5,
207
208 Failed = 6,
210
211 Cancelled = 7,
213}
214
215#[implement(Service)]
219#[inline]
220#[must_use]
221pub fn prev_walk_metrics(&self) -> PrevWalkMetrics { self.prev_walk.snapshot() }
222
223#[implement(PrevWalkCounters)]
224fn snapshot(&self) -> PrevWalkMetrics {
225 PrevWalkMetrics {
226 entered: self.entered.load(Ordering::Relaxed),
227 gapped: self.gapped.load(Ordering::Relaxed),
228 held: self.held.load(Ordering::Relaxed),
229 closed: self.closed.load(Ordering::Relaxed),
230 fetch_failed: self.fetch_failed.load(Ordering::Relaxed),
231 fetch_cancelled: self.fetch_cancelled.load(Ordering::Relaxed),
232 walked: self.walked.load(Ordering::Relaxed),
233 walked_prevs: self.walked_prevs.load(Ordering::Relaxed),
234 capped: self.capped.load(Ordering::Relaxed),
235 appended: self.appended.load(Ordering::Relaxed),
236 not_appended: self.not_appended.load(Ordering::Relaxed),
237 failed: self.failed.load(Ordering::Relaxed),
238 cancelled: self.cancelled.load(Ordering::Relaxed),
239 unprocessed_prevs: self.unprocessed_prevs.load(Ordering::Relaxed),
240 }
241}
242
243#[implement(Service)]
249#[inline]
250#[must_use]
251pub fn prev_walks_in_flight(&self) -> impl ExactSizeIterator<Item = InFlightWalk> + Send + use<> {
252 self.prev_walks_in_flight.snapshot()
253}
254
255#[implement(Service)]
259#[inline]
260#[must_use]
261pub fn prev_walks_in_flight_count(&self) -> usize { self.prev_walks_in_flight.len() }
262
263#[implement(InFlightWalks)]
264fn snapshot(&self) -> impl ExactSizeIterator<Item = InFlightWalk> + Send + use<> {
265 self.lock()
266 .values()
267 .cloned()
268 .collect::<Vec<_>>()
269 .into_iter()
270}
271
272#[implement(InFlightWalks)]
277#[inline]
278fn lock(&self) -> MutexGuard<'_, Walks> { self.walks.lock_adopting() }
279
280#[implement(InFlightWalks)]
281#[inline]
282pub(super) fn len(&self) -> usize { self.lock().len() }
283
284#[implement(PrevWalkCounters)]
285pub(super) fn enter(&self, gapped: bool) {
286 self.entered.fetch_add(1, Ordering::Relaxed);
287
288 if gapped {
289 self.gapped.fetch_add(1, Ordering::Relaxed);
290 }
291}
292
293#[implement(PrevWalk, generics = "<'a>", params = "<'a>")]
294pub(super) fn start(service: &'a Service, upgrade: &'a PrevUpgrade<'a>) -> Self {
295 let started = Instant::now();
296 let id = service
297 .prev_walks_in_flight
298 .insert(upgrade, started);
299
300 Self {
301 service,
302 id,
303 upgrade,
304 started,
305 walk: None,
306 settled: None,
307 }
308}
309
310#[implement(InFlightWalks)]
311fn insert(&self, upgrade: &PrevUpgrade<'_>, started: Instant) -> u64 {
312 let PrevUpgrade { origin, room_id, event_id, .. } = *upgrade;
313 let id = self.next.fetch_add(1, Ordering::Relaxed);
314 let entry = InFlightWalk {
315 room_id: room_id.to_owned(),
316 event_id: event_id.to_owned(),
317 origin: origin.to_owned(),
318 started,
319 walk: None,
320 };
321
322 self.lock().insert(id, entry);
323
324 id
325}
326
327#[implement(PrevWalk, params = "<'_>")]
328pub(super) fn hold(self) { self.end(Outcome::Held, 0); }
329
330#[implement(PrevWalk, params = "<'_>")]
331fn end(mut self, outcome: Outcome, unprocessed: usize) {
332 self.settled = Some((outcome, unprocessed));
333}
334
335#[implement(PrevWalk, params = "<'_>")]
341pub(super) fn fetched(self, fetch: Result<&PrevFetch, &Error>, stopping: bool) -> Option<Self> {
342 let outcome = match fetch {
343 | Err(error) => Outcome::fetch_error(error, stopping),
344 | Ok(fetch) if fetch.sorted.is_empty() => Outcome::Closed,
345 | Ok(fetch) => return Some(self.begin_walk(fetch)),
346 };
347
348 self.end(outcome, 0);
349
350 None
351}
352
353#[implement(Outcome)]
354fn fetch_error(error: &Error, stopping: bool) -> Self {
355 if Self::cut_off(error, stopping) {
356 Self::FetchCancelled
357 } else {
358 Self::FetchFailed
359 }
360}
361
362#[implement(Outcome)]
367fn cut_off(error: &Error, stopping: bool) -> bool { stopping || error.is_interrupted() }
368
369#[implement(PrevWalk, params = "<'_>")]
370fn begin_walk(self, fetch: &PrevFetch) -> Self {
371 let prevs = fetch.pdus.len();
372 let walk = Walk {
373 fetched: Instant::now(),
374 prevs,
375 capped: fetch.capped,
376 };
377
378 self.service.prev_walk.start_walk(prevs);
379 self.service
380 .prev_walks_in_flight
381 .walking(self.id, walk);
382
383 self.with_walk(walk)
384}
385
386#[implement(PrevWalkCounters)]
387fn start_walk(&self, prevs: usize) {
388 self.walked.fetch_add(1, Ordering::Relaxed);
389 fetch_add_usize(&self.walked_prevs, prevs, Ordering::Relaxed);
390}
391
392#[implement(InFlightWalks)]
393fn walking(&self, id: u64, walk: Walk) {
394 self.lock()
395 .entry(id)
396 .and_modify(|in_flight| in_flight.walk = Some(walk));
397}
398
399#[implement(PrevWalk, params = "<'_>")]
400fn with_walk(mut self, walk: Walk) -> Self {
401 self.walk = Some(walk);
402 self
403}
404
405#[implement(PrevWalk, params = "<'_>")]
406pub(super) fn settle(self, appended: Result<bool, &Error>, upgraded: usize, stopping: bool) {
407 let prevs = self.walk.map_or(0, |walk| walk.prevs);
408
409 debug_assert!(upgraded <= prevs, "upgraded more previous events than were collected");
410
411 let outcome = Outcome::new(appended, stopping);
412 let unprocessed = match outcome {
413 | Outcome::Cancelled => 0,
414 | _ => prevs.saturating_sub(upgraded),
415 };
416
417 self.end(outcome, unprocessed);
418}
419
420#[implement(Outcome)]
421fn new(appended: Result<bool, &Error>, stopping: bool) -> Self {
422 match appended {
423 | Ok(true) => Self::Appended,
424 | Ok(false) => Self::NotAppended,
425 | Err(error) if Self::cut_off(error, stopping) => Self::Cancelled,
426 | Err(_) => Self::Failed,
427 }
428}
429
430impl Drop for PrevWalk<'_> {
431 fn drop(&mut self) {
432 let Self {
433 service,
434 id,
435 upgrade,
436 started,
437 walk,
438 settled,
439 } = *self;
440
441 let pass = Pass::new(started, walk, settled, Instant::now());
442
443 service.prev_walks_in_flight.remove(id);
444 service.prev_walk.settle_pass(&pass);
445
446 if !pass.outcome.reported() {
447 return;
448 }
449
450 match walk {
451 | None => log_fetch_end(upgrade, &pass),
452 | Some(_) => log_walk_end(upgrade, &pass),
453 }
454
455 if !panicking() {
457 service.record_pass(upgrade, &pass);
458 }
459 }
460}
461
462#[implement(Pass)]
468fn new(
469 started: Instant,
470 walk: Option<Walk>,
471 settled: Option<(Outcome, usize)>,
472 ended: Instant,
473) -> Self {
474 let unsettled = walk.map_or(Outcome::FetchCancelled, |_| Outcome::Cancelled);
475 let (outcome, unprocessed) = settled.unwrap_or((unsettled, 0));
476 let fetched = walk.map_or(ended, |walk| walk.fetched);
477
478 Self {
479 outcome,
480 prevs: walk.map_or(0, |walk| walk.prevs),
481 unprocessed,
482 capped: walk.is_some_and(|walk| walk.capped),
483 fetch: fetched.saturating_duration_since(started),
484 upgrade: ended.saturating_duration_since(fetched),
485 }
486}
487
488#[implement(InFlightWalks)]
489fn remove(&self, id: u64) { self.lock().remove(&id); }
490
491#[implement(PrevWalkCounters)]
492fn settle_pass(&self, &Pass { outcome, capped, unprocessed, .. }: &Pass) {
493 let counter = match outcome {
494 | Outcome::Held => &self.held,
495 | Outcome::Closed => &self.closed,
496 | Outcome::FetchFailed => &self.fetch_failed,
497 | Outcome::FetchCancelled => &self.fetch_cancelled,
498 | Outcome::Appended => &self.appended,
499 | Outcome::NotAppended => &self.not_appended,
500 | Outcome::Failed => &self.failed,
501 | Outcome::Cancelled => &self.cancelled,
502 };
503
504 counter.fetch_add(1, Ordering::Relaxed);
505
506 if capped {
507 self.capped.fetch_add(1, Ordering::Relaxed);
508 }
509
510 fetch_add_usize(&self.unprocessed_prevs, unprocessed, Ordering::Relaxed);
511}
512
513#[implement(Outcome)]
515fn reported(self) -> bool {
516 match self {
517 | Self::Held | Self::Closed => false,
518 | Self::FetchFailed
519 | Self::FetchCancelled
520 | Self::Appended
521 | Self::NotAppended
522 | Self::Failed
523 | Self::Cancelled => true,
524 }
525}
526
527fn log_fetch_end(upgrade: &PrevUpgrade<'_>, pass: &Pass) {
528 let PrevUpgrade { origin, room_id, event_id, .. } = *upgrade;
529 let fetch_ms = pass.fetch.as_millis();
530
531 warn!(
532 %room_id,
533 %event_id,
534 %origin,
535 outcome = pass.outcome.name(),
536 fetch_ms,
537 "Prev walk ended."
538 );
539}
540
541fn log_walk_end(upgrade: &PrevUpgrade<'_>, pass: &Pass) {
542 let PrevUpgrade { origin, room_id, event_id, .. } = *upgrade;
543 let Pass { outcome, prevs, unprocessed, capped, .. } = *pass;
544 let fetch_ms = pass.fetch.as_millis();
545 let upgrade_ms = pass.upgrade.as_millis();
546
547 if capped || matches!(outcome, Outcome::Failed | Outcome::Cancelled) {
548 warn!(
549 %room_id,
550 %event_id,
551 %origin,
552 outcome = outcome.name(),
553 prevs,
554 unprocessed,
555 capped,
556 fetch_ms,
557 upgrade_ms,
558 "Prev walk ended."
559 );
560 } else {
561 debug_info!(
562 %room_id,
563 %event_id,
564 %origin,
565 outcome = outcome.name(),
566 prevs,
567 unprocessed,
568 capped,
569 fetch_ms,
570 upgrade_ms,
571 "Prev walk ended."
572 );
573 }
574}
575
576#[implement(Outcome)]
581#[must_use]
582pub fn name(self) -> &'static str {
583 match self {
584 | Self::Held => "held",
585 | Self::Closed => "closed",
586 | Self::FetchFailed => "fetch_failed",
587 | Self::FetchCancelled => "fetch_cancelled",
588 | Self::Appended => "appended",
589 | Self::NotAppended => "not_appended",
590 | Self::Failed => "failed",
591 | Self::Cancelled => "cancelled",
592 }
593}