Skip to main content

tuwunel_service/rooms/event_handler/
prev_walk.rs

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/// Snapshot of process-lifetime totals for incoming previous-event walks.
28///
29/// Every gapped event ends as exactly one of `held`, `closed`, `fetch_failed`,
30/// `fetch_cancelled` or `walked`, and every walk as exactly one of `appended`,
31/// `not_appended`, `failed` or `cancelled`. Each counter is loaded
32/// independently, so totals within one snapshot need not reconcile while
33/// passes are in flight.
34#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
35pub struct PrevWalkMetrics {
36	/// Top-level incoming timeline events that reached the gap check, whether from
37	/// a federation transaction, a join, an invite or backfill.
38	pub entered: u64,
39
40	/// Events with a previous event missing from the timeline.
41	pub gapped: u64,
42
43	/// Gapped events withheld by a backoff hold.
44	pub held: u64,
45
46	/// Gapped events whose previous events all reached the timeline during the
47	/// backward fetch, leaving nothing to walk.
48	pub closed: u64,
49
50	/// Gapped events whose backward fetch returned an error.
51	pub fetch_failed: u64,
52
53	/// Gapped events dropped before their walk began, or whose backward fetch
54	/// was interrupted.
55	pub fetch_cancelled: u64,
56
57	/// Passes that walked previous events missing from the timeline.
58	pub walked: u64,
59
60	/// Previous events collected for upgrade by walking passes.
61	pub walked_prevs: u64,
62
63	/// Passes whose walk hit the `max_fetch_prev_events` cap.
64	pub capped: u64,
65
66	/// Passes that appended the incoming event.
67	pub appended: u64,
68
69	/// Passes that finished without appending the incoming event.
70	pub not_appended: u64,
71
72	/// Passes that returned an error.
73	pub failed: u64,
74
75	/// Passes dropped before they settled or interrupted by shutdown.
76	pub cancelled: u64,
77
78	/// Collected previous events left without an upgrade by passes that were
79	/// not cancelled.
80	pub unprocessed_prevs: u64,
81}
82
83/// A gapped incoming event whose pass is in flight.
84///
85/// The event is listed from its gap check until its pass settles or is dropped.
86#[derive(Clone, Debug)]
87pub struct InFlightWalk {
88	/// Room of the incoming event.
89	pub room_id: OwnedRoomId,
90
91	/// The gapped incoming event.
92	pub event_id: OwnedEventId,
93
94	/// Server the incoming event arrived from.
95	pub origin: OwnedServerName,
96
97	/// When the pass started, right after the gap check.
98	pub started: Instant,
99
100	/// The walking phase, absent while the backoff lookup or backward fetch runs.
101	pub walk: Option<Walk>,
102}
103
104/// The walking phase of a pass, entered when the backward fetch returns
105/// previous events.
106///
107/// A pass dropped before its walk starts counts as a cancelled fetch rather
108/// than a cancelled walk.
109#[derive(Clone, Copy, Debug)]
110pub struct Walk {
111	/// When the backward fetch returned.
112	pub fetched: Instant,
113
114	/// Previous events the fetch collected for upgrade.
115	pub prevs: usize,
116
117	/// Whether the fetch hit the `max_fetch_prev_events` cap.
118	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/// Registry of the passes in flight, keyed by ids taken in registration order.
140///
141/// Concurrent passes in one room would collide on a room key, so each pass
142/// takes its own id. An entry lives exactly as long as its pass's guard, whose
143/// drop removes it, so the registry never outgrows the gapped passes running at
144/// once. The lock is never held across an await.
145#[derive(Default)]
146pub(super) struct InFlightWalks {
147	next: AtomicU64,
148	walks: Mutex<Walks>,
149}
150
151/// Drop guard over one gapped incoming event, from its backoff lookup through
152/// the backward fetch and the upgrade of its previous events.
153///
154/// The event is listed in flight while the guard lives. Dropping the guard
155/// records the event under its outcome, or as a cancelled fetch or walk when
156/// dropped unsettled. A walk, a failed fetch and a cancelled fetch each log one
157/// line and record one row in the room's history, the row skipped while the
158/// thread unwinds; a hold and a fetch that leaves nothing to walk do neither.
159#[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/// The end of one pass, computed once when its guard drops.
171///
172/// The counters, the log line and the room's history row all read this one
173/// record, so they agree on the outcome and the timings.
174#[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/// How the pass of a gapped incoming event ended.
185///
186/// The first four end a pass before its walk began, and the last four end a
187/// walk.
188#[derive(Clone, Copy, Debug, Eq, PartialEq)]
189pub enum Outcome {
190	/// Withheld by a backoff hold.
191	Held = 0,
192
193	/// The backward fetch left nothing to walk.
194	Closed = 1,
195
196	/// The backward fetch returned an error.
197	FetchFailed = 2,
198
199	/// Dropped before the walk began, or the backward fetch was interrupted.
200	FetchCancelled = 3,
201
202	/// The walk appended the incoming event.
203	Appended = 4,
204
205	/// The walk finished without appending the incoming event.
206	NotAppended = 5,
207
208	/// The walk returned an error.
209	Failed = 6,
210
211	/// The walk was dropped before it settled, or interrupted by shutdown.
212	Cancelled = 7,
213}
214
215/// Read a snapshot of incoming previous-event walk totals.
216///
217/// Reading does not reset the totals.
218#[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/// Read the gapped incoming events whose passes are in flight, in registration
244/// order.
245///
246/// The entries are copied out under the registry lock, since an iterator over
247/// the registry itself could not outlive the lock.
248#[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/// Count the gapped incoming events whose passes are in flight.
256///
257/// Counting copies nothing out of the registry.
258#[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/// Take the registry lock, adopting a poisoned one.
273///
274/// Nothing under the lock can panic, so a poisoned lock still guards a
275/// consistent map. Adopting it keeps `PrevWalk`'s drop from panicking.
276#[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/// Record the end of the backward fetch.
336///
337/// Returns `None` once the pass is settled, when the fetch failed or left
338/// nothing to walk. Otherwise hands the guard back for the walk, which `settle`
339/// ends.
340#[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/// Whether an error cut the pass off rather than failing it.
363///
364/// An error while the server stops, or an interruption, says nothing about the
365/// events themselves.
366#[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		// A panic from the write while unwinding would abort the process.
456		if !panicking() {
457			service.record_pass(upgrade, &pass);
458		}
459	}
460}
461
462/// Compute how a pass ended as of `ended`.
463///
464/// An unsettled pass counts as cancelled, as a fetch before its walk began and
465/// as a walk after. The fetch phase runs until the walk began, or until the end
466/// for a pass that never walked.
467#[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/// Whether a pass ending this way logs and records its end.
514#[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/// The outcome's name.
577///
578/// The same name labels the outcome in the `Prev walk ended.` log line and in
579/// the per-room history view.
580#[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}