Skip to main content

tuwunel_admin/federation/
incoming_federation.rs

1use std::{
2	fmt::{Result as FmtResult, Write as _},
3	time::Instant,
4};
5
6use ruma::OwnedRoomId;
7use tuwunel_core::{
8	Result,
9	itertools::Itertools,
10	utils::{
11		string::{markdown_cell, plural},
12		time::Elapsed,
13	},
14};
15use tuwunel_service::rooms::event_handler::{InFlightWalk, Walk};
16
17use crate::admin_command;
18
19#[admin_command]
20pub(super) async fn incoming_federation(&self) -> Result {
21	let event_handler = &self.services.event_handler;
22	let locked = event_handler.mutex_federation.keys();
23	let walks: Vec<_> = event_handler.prev_walks_in_flight().collect();
24	let now = Instant::now();
25	let output = render(&walks, unwalked(locked, &walks), now);
26
27	self.write_str(&output).await
28}
29
30fn unwalked(
31	locked: impl IntoIterator<Item = OwnedRoomId>,
32	walks: &[InFlightWalk],
33) -> impl ExactSizeIterator<Item = OwnedRoomId> {
34	locked
35		.into_iter()
36		.filter(|room_id| walks.iter().all(|walk| walk.room_id.ne(room_id)))
37		.sorted_unstable()
38}
39
40fn render(
41	walks: &[InFlightWalk],
42	unwalked: impl ExactSizeIterator<Item = OwnedRoomId>,
43	now: Instant,
44) -> String {
45	let mut output = String::new();
46
47	render_into(&mut output, walks, unwalked, now).expect("writing to a String cannot fail");
48
49	output
50}
51
52fn render_into(
53	output: &mut String,
54	walks: &[InFlightWalk],
55	unwalked: impl ExactSizeIterator<Item = OwnedRoomId>,
56	now: Instant,
57) -> FmtResult {
58	write_walks(output, walks, now)?;
59	write_unwalked(output, unwalked)
60}
61
62fn write_walks(output: &mut String, walks: &[InFlightWalk], now: Instant) -> FmtResult {
63	let noun = plural(walks.len(), "prev walk", "prev walks");
64
65	writeln!(output, "{} {noun} in flight.", walks.len())?;
66	if walks.is_empty() {
67		return Ok(());
68	}
69
70	writeln!(output, "\n| room | event | origin | phase | elapsed | fetch | prevs | capped |")?;
71	writeln!(output, "| :--- | :--- | :--- | :--- | ---: | ---: | ---: | :--- |")?;
72	for walk in walks {
73		write_walk(output, walk, now)?;
74	}
75
76	Ok(())
77}
78
79fn write_walk(output: &mut String, in_flight: &InFlightWalk, now: Instant) -> FmtResult {
80	let InFlightWalk { room_id, event_id, origin, started, walk } = in_flight;
81	let elapsed = Elapsed::from(now.saturating_duration_since(*started));
82
83	write!(
84		output,
85		"| {} | {} | {} |",
86		markdown_cell(room_id.as_str()),
87		markdown_cell(event_id.as_str()),
88		markdown_cell(origin.as_str()),
89	)?;
90
91	match walk {
92		| None => writeln!(output, " fetch | {elapsed} | | | |"),
93		| Some(Walk { fetched, prevs, capped }) => {
94			let fetch = Elapsed::from(fetched.saturating_duration_since(*started));
95			let capped = if *capped { "yes" } else { "no" };
96
97			writeln!(output, " walk | {elapsed} | {fetch} | {prevs} | {capped} |")
98		},
99	}
100}
101
102fn write_unwalked(
103	output: &mut String,
104	unwalked: impl ExactSizeIterator<Item = OwnedRoomId>,
105) -> FmtResult {
106	let rooms = unwalked.len();
107	let noun = plural(rooms, "room", "rooms");
108
109	writeln!(output, "\n{rooms} {noun} holding the federation lock with no walk.")?;
110	if rooms == 0 {
111		return Ok(());
112	}
113
114	writeln!(output, "\n| room |\n| :--- |")?;
115	for room_id in unwalked {
116		writeln!(output, "| {} |", markdown_cell(room_id.as_str()))?;
117	}
118
119	Ok(())
120}
121
122#[cfg(test)]
123mod tests {
124	use std::time::Duration;
125
126	use ruma::{owned_event_id, owned_room_id, owned_server_name};
127
128	use super::*;
129
130	#[test]
131	fn render_lists_walks_then_rooms_without_one() {
132		let started = Instant::now();
133		let walk = Walk {
134			fetched: after(started, 850),
135			prevs: 37,
136			capped: true,
137		};
138
139		let walks = [
140			InFlightWalk {
141				room_id: owned_room_id!("!fetching:example.org"),
142				event_id: owned_event_id!("$fetch"),
143				origin: owned_server_name!("other.example"),
144				started: after(started, 12_830),
145				walk: None,
146			},
147			InFlightWalk {
148				room_id: owned_room_id!("!walking:example.org"),
149				event_id: owned_event_id!("$walk"),
150				origin: owned_server_name!("remote.example"),
151				started,
152				walk: Some(walk),
153			},
154		];
155
156		let locked =
157			[owned_room_id!("!walking:example.org"), owned_room_id!("!idle:example.org")];
158
159		let output = render(&walks, unwalked(locked, &walks), after(started, 14_030));
160		let (_, rooms) = output
161			.split_once("holding the federation lock")
162			.expect("the rooms without a walk are counted");
163
164		assert!(rooms.contains("| !idle:example.org |"), "an idle locked room must be listed");
165		assert!(
166			!rooms.contains("!walking:example.org"),
167			"a room with a walk must not be listed again"
168		);
169
170		assert_eq!(output.lines().collect::<Vec<_>>(), [
171			"2 prev walks in flight.",
172			"",
173			"| room | event | origin | phase | elapsed | fetch | prevs | capped |",
174			"| :--- | :--- | :--- | :--- | ---: | ---: | ---: | :--- |",
175			"| !fetching:example.org | $fetch | other.example | fetch | 1.2s | | | |",
176			"| !walking:example.org | $walk | remote.example | walk | 14.03s | 850ms | 37 | yes \
177			 |",
178			"",
179			"1 room holding the federation lock with no walk.",
180			"",
181			"| room |",
182			"| :--- |",
183			"| !idle:example.org |",
184		]);
185	}
186
187	#[test]
188	fn render_omits_empty_tables() {
189		let output = render(&[], unwalked([], &[]), Instant::now());
190
191		assert_eq!(
192			output,
193			"0 prev walks in flight.\n\n0 rooms holding the federation lock with no walk.\n"
194		);
195	}
196
197	fn after(started: Instant, millis: u64) -> Instant {
198		started
199			.checked_add(Duration::from_millis(millis))
200			.expect("the test instant is in range")
201	}
202}