tuwunel_service/sending/sender/
split.rs1use 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
22type Rooms = SmallVec<[ShortRoomId; 2]>;
26
27pub(super) type Slice = (Vec<QueueItem>, Split);
28
29const SPLIT_AFTER: u32 = 4;
33
34#[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#[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#[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#[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#[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#[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
204fn 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;