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
32const 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#[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
73struct 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
88struct 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
230fn 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#[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 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
330fn 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}