Skip to main content

tuwunel_service/sending/sender/
split.rs

1use std::pin::pin;
2
3use futures::{StreamExt, future::OptionFuture};
4use ruma::ServerName;
5use tuwunel_core::{
6	Error, debug_info, extract_variant, implement,
7	matrix::ShortRoomId,
8	smallvec::SmallVec,
9	utils::{IterStream, ReadyExt},
10	warn,
11};
12
13use super::PDU_LIMIT;
14use crate::{
15	federation::is_content_rejection,
16	sending::{
17		Destination, SendingEvent, Service,
18		data::{Park, QueueItem},
19	},
20};
21
22/// Rooms of a rejected transaction still to send.
23///
24/// A rejected transaction rarely spans more than a couple of rooms.
25type Rooms = SmallVec<[ShortRoomId; 2]>;
26
27pub(super) type Slice = (Vec<QueueItem>, Split);
28
29/// Consecutive failures of a transaction before its rooms are sent apart.
30///
31/// The first few retries ride out a transient fault before the batch is blamed.
32const SPLIT_AFTER: u32 = 4;
33
34/// The rooms of a transaction a destination kept rejecting, sent one at a time.
35///
36/// The head room is the one in flight. Until a room is delivered, each failure
37/// moves the head last behind one more queued room brought in as a control. A
38/// room that fails after a delivery is parked. Boxed to keep transaction
39/// statuses small.
40#[derive(Debug, Eq, PartialEq)]
41pub(super) struct Split(Box<Rounds>);
42
43#[derive(Debug, Eq, PartialEq)]
44struct Rounds {
45	rooms: Rooms,
46	delivered: bool,
47}
48
49/// Choose how a failed federation transaction continues.
50///
51/// After repeated failures the transaction's PDU rows return to the queue and
52/// its rooms go out one per transaction; a room failing after a delivery is
53/// parked. Returns the split to keep and the failure count its retry backs off
54/// by, zero when the split advanced.
55#[implement(Service)]
56pub(super) async fn split_failure(
57	&self,
58	server: &ServerName,
59	error: &Error,
60	split: Option<Split>,
61	tries: u32,
62) -> (Option<Split>, u32) {
63	let implicated = implicates_content(error);
64
65	if split.is_none() && (!implicated || tries < SPLIT_AFTER) {
66		return (None, tries);
67	}
68
69	let dest = Destination::Federation(server.to_owned());
70	let active: Vec<_> = self.db.active_requests_for(&dest).collect().await;
71	let rejected = split
72		.as_ref()
73		.filter(|split| implicated && split.0.delivered)
74		.and_then(Split::head)
75		.map(|room| self.db.next_park(server, room));
76
77	let park = OptionFuture::from(rejected).await;
78
79	self.db
80		.demote(&active, park.map(|park| (server, park)));
81
82	if let Some(Park { room, until, count }) = park {
83		warn!(%server, room, until, count, "Parked a room the server keeps rejecting");
84	}
85
86	match split {
87		| Some(split) if !implicated => (Some(split), tries),
88		| Some(split) if split.0.delivered => (Some(split.skipped()), 0),
89		| Some(split) => {
90			let control = self.control(&dest, server, &split.0.rooms).await;
91
92			(Some(split.rotated(control)), tries)
93		},
94		| None => start(server, &active, tries),
95	}
96}
97
98/// Find the first queued room outside the split and the server's parks.
99///
100/// A control delivered while the split's rooms keep failing implicates those
101/// rooms rather than the destination.
102#[implement(Service)]
103async fn control(
104	&self,
105	dest: &Destination,
106	server: &ServerName,
107	rooms: &Rooms,
108) -> Option<ShortRoomId> {
109	let skip: Rooms = self
110		.db
111		.parks(server)
112		.map(|park| park.room)
113		.chain(rooms.iter().copied().stream())
114		.collect()
115		.await;
116
117	let skip = sorted(skip);
118
119	pin!(self.db.queued_except(dest, &skip))
120		.ready_find_map(|item| pdu_room(&item))
121		.await
122}
123
124/// Advance a split past its delivered head room, ending any park it had.
125///
126/// Returns the next room's transaction, or nothing once the split is done.
127#[implement(Service)]
128pub(super) async fn split_delivered(&self, dest: &Destination, split: Split) -> Option<Slice> {
129	if let (Destination::Federation(server), Some(room)) = (dest, split.head()) {
130		self.db.unpark(server, room);
131	}
132
133	let next = self.slice(dest, split.delivered()).await;
134
135	if next.is_none() {
136		debug_info!(?dest, "Finished sending rooms apart");
137	}
138
139	next
140}
141
142/// Promote the head room's queued rows as the split's next transaction.
143///
144/// The split ends when it has no rooms left, or its head room has no rows. The
145/// rejected transaction's EDU rows stay active and ride along until delivered.
146#[implement(Service)]
147pub(super) async fn slice(&self, dest: &Destination, split: Split) -> Option<Slice> {
148	let room = split.head()?;
149	let rows: Vec<_> = self
150		.db
151		.queued_room(dest, room)
152		.take(PDU_LIMIT)
153		.collect()
154		.await;
155
156	if rows.is_empty() {
157		return None;
158	}
159
160	self.db.mark_as_active(rows.iter());
161
162	let items = self.db.active_requests_for(dest).collect().await;
163
164	Some((items, split))
165}
166
167#[implement(Split)]
168fn new(rooms: Rooms) -> Self { Self(Box::new(Rounds { rooms, delivered: false })) }
169
170/// Retry a room whose park expired.
171///
172/// The probe counts as delivered, so a failure parks the room again at once.
173#[implement(Split)]
174pub(super) fn probe(room: ShortRoomId) -> Self {
175	Self(Box::new(Rounds {
176		rooms: Rooms::from_slice(&[room]),
177		delivered: true,
178	}))
179}
180
181#[implement(Split)]
182#[inline]
183fn head(&self) -> Option<ShortRoomId> { self.0.rooms.first().copied() }
184
185#[implement(Split)]
186fn delivered(mut self) -> Self {
187	self.0.delivered = true;
188	self.skipped()
189}
190
191#[implement(Split)]
192fn skipped(mut self) -> Self {
193	self.0.rooms.drain(..self.0.rooms.len().min(1));
194	self
195}
196
197#[implement(Split)]
198fn rotated(mut self, control: Option<ShortRoomId>) -> Self {
199	self.0.rooms.rotate_left(1);
200	self.0.rooms.insert_many(0, control);
201	self
202}
203
204/// Whether a failure answers for the transaction's content, not the path to the peer.
205///
206/// Content rejections and server errors come from the peer, and a timeout after
207/// connecting may be the peer stalling on the content. Connection failures and
208/// rate limits say nothing about it.
209fn implicates_content(error: &Error) -> bool {
210	match error {
211		| Error::Federation(_, response) =>
212			response.status_code.is_server_error() || is_content_rejection(error),
213		| Error::Reqwest(error) => error.is_timeout() && !error.is_connect(),
214		| _ => false,
215	}
216}
217
218fn start(server: &ServerName, active: &[QueueItem], tries: u32) -> (Option<Split>, u32) {
219	let rooms = sorted(active.iter().filter_map(pdu_room).collect());
220
221	if rooms.is_empty() {
222		return (None, tries);
223	}
224
225	debug_info!(%server, ?rooms, "Sending a rejected transaction's rooms apart");
226	(Some(Split::new(rooms)), 0)
227}
228
229fn sorted(mut rooms: Rooms) -> Rooms {
230	rooms.sort_unstable();
231	rooms.dedup();
232	rooms
233}
234
235pub(super) fn pdu_room((_, event): &QueueItem) -> Option<ShortRoomId> {
236	extract_variant!(event, SendingEvent::Pdu).map(|id| u64::from_be_bytes(id.shortroomid()))
237}
238
239#[cfg(test)]
240mod tests;