Skip to main content

tuwunel_admin/query/feds/
ping.rs

1use std::{
2	borrow::Cow,
3	fmt::{Display, Formatter, Result as FmtResult, Write as _},
4	iter::once,
5	num::NonZeroUsize,
6	time::{Duration, Instant},
7};
8
9use futures::{StreamExt, stream::iter};
10use ruma::{OwnedRoomOrAliasId, api::federation::discovery::get_server_version::v1::Request};
11use tuwunel_core::{
12	Result, implement,
13	itertools::Itertools,
14	utils::time::{Elapsed, now_secs},
15};
16use tuwunel_service::federation::{
17	PeerBackoff,
18	feds::{Fault, Outcome},
19};
20
21use super::{
22	Backoffs, ListMode, Sort, SweepArgs, count_results, fault_message, markdown_cell,
23	partition_backoffs, prepare, render_totals, retry_after, sorted, write_cell,
24	write_elapsed_cell,
25};
26use crate::admin_command;
27
28pub(super) const WIDTH_DEFAULT: NonZeroUsize = NonZeroUsize::new(192).expect("192 is nonzero");
29
30const OPEN_BUCKET: Duration = Duration::from_secs(5);
31
32/// Exclusive upper bounds of the latency histogram buckets.
33///
34/// The last bound opens the final bucket, which collects every slower
35/// response.
36const BUCKETS: &[Duration] = &[
37	Duration::from_millis(10),
38	Duration::from_millis(25),
39	Duration::from_millis(50),
40	Duration::from_millis(100),
41	Duration::from_millis(250),
42	Duration::from_millis(500),
43	Duration::from_secs(1),
44	Duration::from_millis(2_500),
45	OPEN_BUCKET,
46];
47
48type PingOutcome = Outcome<()>;
49
50/// Settlement class of one destination.
51///
52/// The classes partition every outcome, so their counts sum to the
53/// destination count.
54#[derive(Clone, Copy, Eq, PartialEq)]
55enum Disposition {
56	Responded,
57	Failed,
58	TimedOut,
59	BackingOff,
60	NotAttempted,
61}
62
63impl Disposition {
64	const ALL: [Self; 5] = [
65		Self::Responded,
66		Self::Failed,
67		Self::TimedOut,
68		Self::BackingOff,
69		Self::NotAttempted,
70	];
71}
72
73/// Distribution of one latency population.
74///
75/// Percentiles use the nearest rank and the deviation is the population
76/// standard deviation, both computed in integer nanoseconds.
77struct Stats {
78	count: usize,
79	min: Duration,
80	p50: Duration,
81	mean: Duration,
82	p90: Duration,
83	p99: Duration,
84	max: Duration,
85	stddev: Duration,
86}
87
88/// Fraction of a population.
89///
90/// Renders as a percentage truncated to one decimal place.
91struct Share {
92	count: usize,
93	total: usize,
94}
95
96#[admin_command]
97pub(super) async fn feds_ping(
98	&self,
99	room: OwnedRoomOrAliasId,
100	list: bool,
101	list_all: bool,
102	list_errors: bool,
103	sort: Sort,
104	sweep: SweepArgs,
105) -> Result {
106	let prepared = prepare(self, &room, sweep, WIDTH_DEFAULT).await?;
107	let backoffs = self.services.federation.peer_backoffs().await;
108	let now = now_secs();
109	let (eligible, outcomes) = partition_backoffs(self, &prepared, &backoffs, now).await;
110
111	let started = Instant::now();
112	let responses = self
113		.services
114		.federation
115		.fanout_to(iter(eligible), |_| Request::new(), prepared.opts)
116		.map(|outcome| Outcome {
117			origin: outcome.origin,
118			elapsed: outcome.elapsed,
119			result: outcome.result.map(drop),
120		});
121
122	let outcomes: Vec<_> = iter(outcomes).chain(responses).collect().await;
123	let total = started.elapsed();
124	let list_mode = ListMode::new(list, list_all, list_errors);
125	let outcomes = sorted(outcomes, sort, fault_cell);
126	let output = render(&outcomes, &backoffs, now, total, list_mode);
127
128	self.write_str(&output).await
129}
130
131fn fault_cell(outcome: &PingOutcome) -> Cow<'static, str> {
132	outcome
133		.result
134		.as_ref()
135		.err()
136		.map(fault_message)
137		.unwrap_or_default()
138}
139
140fn render(
141	outcomes: &[PingOutcome],
142	backoffs: &Backoffs,
143	now: u64,
144	total: Duration,
145	list_mode: ListMode,
146) -> String {
147	let mut output = String::new();
148
149	render_into(&mut output, outcomes, backoffs, now, total, list_mode)
150		.expect("writing to a String cannot fail");
151
152	output
153}
154
155fn render_into(
156	output: &mut String,
157	outcomes: &[PingOutcome],
158	backoffs: &Backoffs,
159	now: u64,
160	total: Duration,
161	list_mode: ListMode,
162) -> FmtResult {
163	render_dispositions(output, outcomes)?;
164
165	let responded = sorted_elapsed(outcomes, |disposition| disposition == Disposition::Responded);
166	let failed = sorted_elapsed(outcomes, |disposition| disposition == Disposition::Failed);
167	let attempted = sorted_elapsed(outcomes, Disposition::attempted);
168
169	render_latencies(output, &responded, &failed, &attempted)?;
170	render_histogram(output, &responded)?;
171
172	if !matches!(list_mode, ListMode::None) {
173		render_listing(output, outcomes, backoffs, now, list_mode)?;
174	}
175
176	render_totals(output, count_results(outcomes), total)
177}
178
179fn render_dispositions(output: &mut String, outcomes: &[PingOutcome]) -> FmtResult {
180	writeln!(output, "| disposition | servers | share |")?;
181	writeln!(output, "| :--- | ---: | ---: |")?;
182
183	let total = outcomes.len();
184	for disposition in Disposition::ALL {
185		let count = outcomes
186			.iter()
187			.filter(|outcome| Disposition::of(outcome) == disposition)
188			.count();
189
190		writeln!(output, "| {} | {count} | {} |", disposition.label(), Share { count, total })?;
191	}
192
193	Ok(())
194}
195
196#[implement(Disposition)]
197fn of(outcome: &PingOutcome) -> Self {
198	match &outcome.result {
199		| Ok(()) => Self::Responded,
200		| Err(Fault::Error(_)) => Self::Failed,
201		| Err(Fault::Elapsed) => Self::TimedOut,
202		| Err(Fault::Backoff { .. }) => Self::BackingOff,
203		| Err(Fault::NotAttempted) => Self::NotAttempted,
204	}
205}
206
207#[implement(Disposition)]
208fn label(self) -> &'static str {
209	match self {
210		| Self::Responded => "responded",
211		| Self::Failed => "failed",
212		| Self::TimedOut => "timed out",
213		| Self::BackingOff => "backing off",
214		| Self::NotAttempted => "not attempted",
215	}
216}
217
218impl Display for Share {
219	fn fmt(&self, formatter: &mut Formatter<'_>) -> FmtResult {
220		let permille = self
221			.count
222			.saturating_mul(1000)
223			.checked_div(self.total)
224			.unwrap_or_default();
225
226		write!(formatter, "{}.{}%", permille / 10, permille % 10)
227	}
228}
229
230/// Collects the ascending elapsed times of the selected dispositions.
231///
232/// The statistics and histogram index the population, which forces the
233/// buffer.
234fn sorted_elapsed(
235	outcomes: &[PingOutcome],
236	select: impl Fn(Disposition) -> bool,
237) -> Vec<Duration> {
238	outcomes
239		.iter()
240		.filter(|outcome| select(Disposition::of(outcome)))
241		.map(|outcome| outcome.elapsed)
242		.sorted_unstable()
243		.collect()
244}
245
246/// Whether the destination's elapsed time measures a real request.
247#[implement(Disposition)]
248fn attempted(self) -> bool { !matches!(self, Self::BackingOff | Self::NotAttempted) }
249
250fn render_latencies(
251	output: &mut String,
252	responded: &[Duration],
253	failed: &[Duration],
254	attempted: &[Duration],
255) -> FmtResult {
256	let populations = [
257		("responded", Stats::new(responded)),
258		("failed", Stats::new(failed)),
259		("attempted", Stats::new(attempted)),
260	];
261
262	if populations
263		.iter()
264		.all(|(_, stats)| stats.is_none())
265	{
266		return Ok(());
267	}
268
269	writeln!(output, "\n| population | count | min | p50 | mean | p90 | p99 | max | stddev |")?;
270	writeln!(output, "| :--- | ---: | ---: | ---: | ---: | ---: | ---: | ---: | ---: |")?;
271
272	for (label, stats) in populations
273		.iter()
274		.filter_map(|(label, stats)| stats.as_ref().map(|stats| (label, stats)))
275	{
276		writeln!(
277			output,
278			"| {label} | {} | {} | {} | {} | {} | {} | {} | {} |",
279			stats.count,
280			Elapsed::from(stats.min),
281			Elapsed::from(stats.p50),
282			Elapsed::from(stats.mean),
283			Elapsed::from(stats.p90),
284			Elapsed::from(stats.p99),
285			Elapsed::from(stats.max),
286			Elapsed::from(stats.stddev),
287		)?;
288	}
289
290	Ok(())
291}
292
293impl Stats {
294	/// Summarizes an ascending population, `None` when it is empty.
295	///
296	/// Callers sort first because the percentiles index the slice by rank.
297	fn new(sorted: &[Duration]) -> Option<Self> {
298		let (&min, &max) = sorted.first().zip(sorted.last())?;
299		let count = sorted.len();
300		let divisor = u128::try_from(count).unwrap_or(u128::MAX);
301		let mean = sorted
302			.iter()
303			.map(Duration::as_nanos)
304			.sum::<u128>()
305			.checked_div(divisor)
306			.unwrap_or_default();
307
308		let variance = sorted
309			.iter()
310			.map(Duration::as_nanos)
311			.map(|elapsed| elapsed.abs_diff(mean))
312			.map(|deviation| deviation.saturating_mul(deviation))
313			.sum::<u128>()
314			.checked_div(divisor)
315			.unwrap_or_default();
316
317		Some(Self {
318			count,
319			min,
320			p50: percentile(sorted, 50),
321			mean: from_nanos(mean),
322			p90: percentile(sorted, 90),
323			p99: percentile(sorted, 99),
324			max,
325			stddev: from_nanos(variance.isqrt()),
326		})
327	}
328}
329
330/// Selects the nearest-rank percentile of an ascending population.
331///
332/// The rank is the ceiling of `percent` times the population size; an empty
333/// population yields zero.
334fn percentile(sorted: &[Duration], percent: usize) -> Duration {
335	let index = sorted
336		.len()
337		.saturating_mul(percent)
338		.div_ceil(100)
339		.saturating_sub(1);
340
341	sorted.get(index).copied().unwrap_or_default()
342}
343
344fn from_nanos(nanos: u128) -> Duration {
345	Duration::from_nanos(u64::try_from(nanos).unwrap_or(u64::MAX))
346}
347
348fn render_histogram(output: &mut String, responded: &[Duration]) -> FmtResult {
349	if responded.is_empty() {
350		return Ok(());
351	}
352
353	writeln!(output, "\n| latency | servers | cumulative |")?;
354	writeln!(output, "| :--- | ---: | ---: |")?;
355
356	let total = responded.len();
357	let cumulative = BUCKETS
358		.iter()
359		.map(|bound| responded.partition_point(|elapsed| elapsed < bound))
360		.chain(once(total));
361
362	let labels = BUCKETS
363		.iter()
364		.map(|bound| ("<", *bound))
365		.chain(once((">=", OPEN_BUCKET)));
366
367	for ((previous, count), (relation, bound)) in once(0)
368		.chain(cumulative)
369		.tuple_windows()
370		.zip(labels)
371	{
372		let servers = count.saturating_sub(previous);
373
374		writeln!(output, "| {relation} {} | {servers} | {} |", Elapsed::from(bound), Share {
375			count,
376			total
377		})?;
378	}
379
380	Ok(())
381}
382
383fn render_listing(
384	output: &mut String,
385	outcomes: &[PingOutcome],
386	backoffs: &Backoffs,
387	now: u64,
388	list_mode: ListMode,
389) -> FmtResult {
390	writeln!(output, "\n| origin | elapsed | class | newest | oldest | retry | fault |")?;
391	writeln!(output, "| :--- | ---: | :--- | ---: | ---: | ---: | :--- |")?;
392
393	for outcome in outcomes
394		.iter()
395		.filter(|outcome| list_mode.includes(outcome))
396	{
397		render_row(output, outcome, backoffs.get(&outcome.origin), now)?;
398	}
399
400	Ok(())
401}
402
403fn render_row(
404	output: &mut String,
405	outcome: &PingOutcome,
406	backoff: Option<&PeerBackoff>,
407	now: u64,
408) -> FmtResult {
409	write!(output, "| {} |", outcome.origin)?;
410	write_elapsed_cell(output, outcome)?;
411
412	match backoff {
413		| None => write!(output, " | | | |")?,
414		| Some(backoff) => write_backoff_cells(output, backoff, now)?,
415	}
416
417	let fault = fault_cell(outcome);
418	let fault = markdown_cell(&fault);
419
420	write_cell(output, &fault)?;
421	writeln!(output)
422}
423
424fn write_backoff_cells(output: &mut String, backoff: &PeerBackoff, now: u64) -> FmtResult {
425	let newest = Duration::from_secs(now.saturating_sub(backoff.anchor_secs));
426	let oldest = Duration::from_secs(now.saturating_sub(backoff.oldest_secs));
427
428	write!(
429		output,
430		" {:?} | {} | {} |",
431		backoff.class,
432		Elapsed::from(newest),
433		Elapsed::from(oldest)
434	)?;
435
436	match retry_after(backoff, now) {
437		| None => write!(output, " |"),
438		| Some(retry) => write!(output, " {} |", Elapsed::from(retry)),
439	}
440}
441
442#[cfg(test)]
443mod tests {
444	use ruma::{ServerName, server_name};
445	use tuwunel_core::{Error, err};
446	use tuwunel_service::federation::Classification;
447
448	use super::*;
449
450	#[test]
451	fn summary_reports_dispositions_latency_statistics_and_histogram() {
452		let outcomes = vec![
453			responded(server_name!("fast.example"), 10),
454			responded(server_name!("medium.example"), 20),
455			responded(server_name!("slow.example"), 60),
456			failed(server_name!("broken.example"), 40, Fault::Error(error())),
457			failed(server_name!("late.example"), 1_000, Fault::Elapsed),
458			failed(server_name!("skipped.example"), 0, Fault::NotAttempted),
459		];
460
461		let output =
462			rendered(outcomes, &Backoffs::new(), 0, Duration::ZERO, ListMode::None, Sort::Origin);
463
464		assert!(output.starts_with("| disposition | servers | share |\n"));
465		assert!(output.contains("| responded | 3 | 50.0% |\n"));
466		assert!(output.contains("| failed | 1 | 16.6% |\n"));
467		assert!(output.contains("| timed out | 1 | 16.6% |\n"));
468		assert!(output.contains("| backing off | 0 | 0.0% |\n"));
469		assert!(output.contains("| not attempted | 1 | 16.6% |\n"));
470
471		let header = "| population | count | min | p50 | mean | p90 | p99 | max | stddev |\n";
472		let responded = "| responded | 3 | 10ms | 20ms | 30ms | 60ms | 60ms | 60ms | 21.6ms |\n";
473		let failed = "| failed | 1 | 40ms | 40ms | 40ms | 40ms | 40ms | 40ms | 0ns |\n";
474
475		assert!(output.contains(header));
476		assert!(output.contains(responded));
477		assert!(output.contains(failed));
478		assert!(output.contains("| attempted | 5 | 10ms | 40ms | 226ms |"));
479
480		assert!(output.contains("| latency | servers | cumulative |\n"));
481		assert!(output.contains("| < 10ms | 0 | 0.0% |\n"));
482		assert!(output.contains("| < 25ms | 2 | 66.6% |\n"));
483		assert!(output.contains("| < 100ms | 1 | 100.0% |\n"));
484		assert!(output.contains("| >= 5s | 0 | 100.0% |\n"));
485		assert!(!output.contains("| origin |"));
486		assert!(output.ends_with("\n3 results in 0ns.\n"));
487	}
488
489	#[test]
490	fn empty_populations_render_no_latency_tables() {
491		let outcomes = vec![failed(server_name!("skipped.example"), 0, Fault::NotAttempted)];
492
493		let output =
494			rendered(outcomes, &Backoffs::new(), 0, Duration::ZERO, ListMode::None, Sort::Origin);
495
496		assert!(output.contains("| not attempted | 1 | 100.0% |\n"));
497		assert!(!output.contains("| population |"));
498		assert!(!output.contains("| latency |"));
499		assert!(output.ends_with("\n0 results in 0ns.\n"));
500	}
501
502	#[test]
503	fn listing_reports_peer_status_beside_each_origin() {
504		let now = 1_000;
505		let backoffs = Backoffs::from([
506			(server_name!("recovered.example").to_owned(), PeerBackoff {
507				class: Classification::Transient,
508				anchor_secs: 900,
509				oldest_secs: 600,
510				delay_secs: 60,
511			}),
512			(server_name!("held.example").to_owned(), PeerBackoff {
513				class: Classification::Permanent,
514				anchor_secs: 990,
515				oldest_secs: 990,
516				delay_secs: 3_600,
517			}),
518		]);
519
520		let outcomes = || {
521			vec![
522				responded(server_name!("clean.example"), 10),
523				responded(server_name!("recovered.example"), 20),
524				failed(server_name!("held.example"), 0, Fault::Backoff {
525					class: Classification::Permanent,
526					age: Duration::from_secs(10),
527					retry: Duration::from_secs(3_590),
528				}),
529			]
530		};
531
532		let all =
533			rendered(outcomes(), &backoffs, now, Duration::ZERO, ListMode::All, Sort::Origin);
534
535		assert!(all.contains("| origin | elapsed | class | newest | oldest | retry | fault |\n"));
536		assert!(all.contains("| clean.example | 10ms | | | | | |\n"));
537		assert!(all.contains("| recovered.example | 20ms | Transient | 100s | 400s | | |\n"));
538		assert!(all.contains(
539			"| held.example | | Permanent | 10s | 10s | 3590s | peer backoff (Permanent, age \
540			 10s, retry 3590s) |\n"
541		));
542
543		let successes = rendered(
544			outcomes(),
545			&backoffs,
546			now,
547			Duration::ZERO,
548			ListMode::Successes,
549			Sort::Origin,
550		);
551
552		assert!(successes.contains("clean.example"));
553		assert!(!successes.contains("held.example"));
554
555		let errors =
556			rendered(outcomes(), &backoffs, now, Duration::ZERO, ListMode::Errors, Sort::Origin);
557
558		assert!(!errors.contains("clean.example"));
559		assert!(errors.contains("held.example"));
560	}
561
562	#[test]
563	fn listing_sorts_by_elapsed_with_origin_as_tie_breaker() {
564		let outcomes = vec![
565			responded(server_name!("slow.example"), 50),
566			responded(server_name!("b.example"), 10),
567			responded(server_name!("a.example"), 10),
568		];
569
570		let output =
571			rendered(outcomes, &Backoffs::new(), 0, Duration::ZERO, ListMode::All, Sort::Elapsed);
572
573		let (_, listing) = output
574			.split_once("| retry | fault |\n")
575			.expect("detail listing should be rendered");
576
577		let origins: Vec<_> = listing
578			.lines()
579			.skip(1)
580			.take_while(|line| line.starts_with('|'))
581			.filter_map(|line| line.split('|').nth(1))
582			.map(str::trim)
583			.collect();
584
585		assert_eq!(origins, ["a.example", "b.example", "slow.example"]);
586	}
587
588	#[test]
589	fn percentiles_use_nearest_rank() {
590		let sorted: Vec<_> = (1..=10).map(Duration::from_millis).collect();
591
592		assert_eq!(percentile(&sorted, 0), Duration::from_millis(1));
593		assert_eq!(percentile(&sorted, 50), Duration::from_millis(5));
594		assert_eq!(percentile(&sorted, 90), Duration::from_millis(9));
595		assert_eq!(percentile(&sorted, 99), Duration::from_millis(10));
596		assert_eq!(percentile(&[], 50), Duration::ZERO);
597	}
598
599	fn rendered(
600		outcomes: Vec<PingOutcome>,
601		backoffs: &Backoffs,
602		now: u64,
603		total: Duration,
604		list_mode: ListMode,
605		sort: Sort,
606	) -> String {
607		let outcomes = sorted(outcomes, sort, fault_cell);
608
609		render(&outcomes, backoffs, now, total, list_mode)
610	}
611
612	fn responded(origin: &ServerName, elapsed_ms: u64) -> PingOutcome {
613		Outcome {
614			origin: origin.to_owned(),
615			elapsed: Duration::from_millis(elapsed_ms),
616			result: Ok(()),
617		}
618	}
619
620	fn failed(origin: &ServerName, elapsed_ms: u64, fault: Fault) -> PingOutcome {
621		Outcome {
622			origin: origin.to_owned(),
623			elapsed: Duration::from_millis(elapsed_ms),
624			result: Err(fault),
625		}
626	}
627
628	fn error() -> Error { err!(Request(Unknown("connection refused"))) }
629}