tuwunel_admin/federation/
incoming_federation.rs1use 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}