tuwunel_service/sending/sender/
wake.rs1#[cfg(test)]
2mod tests;
3
4use std::{
5 cmp::Reverse,
6 time::{Duration, SystemTime},
7};
8
9use ruma::OwnedServerName;
10use tokio::time::Instant;
11use tuwunel_core::{
12 Error, error,
13 error::error_chain,
14 implement, trace,
15 utils::{
16 BoolExt, exponential_backoff_remaining_secs, rand::secs as rand_secs, time::now_secs,
17 },
18 warn,
19};
20
21use super::{SendingFutures, TransactionStatus, TransactionStatuses, WakeQueue};
22use crate::{
23 federation::ShouldAttempt,
24 sending::{Destination, Msg, SendingEvent, Service},
25};
26
27const PUSH_FAILURE_STREAK: u32 = 4;
28const APPSERVICE_RETRY_BASE: u64 = 2;
29const APPSERVICE_RETRY_MAX_SECS: u64 = 512;
30const WAKE_OVERFLOW_DELAY: Duration = Duration::from_hours(365 * 24);
31
32#[implement(Service)]
33pub(super) async fn drain_due_wakes<'a>(
34 &'a self,
35 futures: &mut SendingFutures<'a>,
36 statuses: &mut TransactionStatuses,
37 wakes: &mut WakeQueue,
38) {
39 let now = Instant::now();
40
41 while wakes
42 .peek()
43 .is_some_and(|Reverse((due, _))| *due <= now)
44 {
45 let Reverse((_, dest)) = wakes.pop().expect("peeked entry");
46
47 self.handle_wake(dest, futures, statuses, wakes)
48 .await;
49 }
50}
51
52#[implement(Service)]
53#[tracing::instrument(name = "wake", level = "debug", skip_all)]
54async fn handle_wake<'a>(
55 &'a self,
56 dest: Destination,
57 futures: &mut SendingFutures<'a>,
58 statuses: &mut TransactionStatuses,
59 wakes: &mut WakeQueue,
60) {
61 let status = statuses.get(&dest);
62
63 if matches!(
64 status,
65 Some(TransactionStatus::Running { .. } | TransactionStatus::RunningForceRetry { .. })
66 ) {
67 return;
68 }
69
70 if matches!(
71 (&dest, status),
72 (Destination::Push(..), Some(TransactionStatus::Retrying { .. }))
73 ) {
74 trace!(?dest, "Dropping push wake while retry is in flight");
75 return;
76 }
77
78 if let (Destination::Push(..), Some(remaining)) = (&dest, self.push_backoff_remaining(status))
79 {
80 if is_armed(wakes, &dest) {
81 trace!(?dest, "Dropping stale push wake");
82 } else {
83 trace!(?dest, ?remaining, "Re-arming early push wake");
84 arm_wake_in(wakes, dest, remaining);
85 }
86
87 return;
88 }
89
90 match dest {
91 | dest @ (Destination::Appservice(_) | Destination::Push(..)) =>
92 self.handle_force_retry(dest, futures, statuses, wakes)
93 .await,
94 | Destination::Federation(server) =>
95 self.handle_federation_wake(server, futures, statuses, wakes)
96 .await,
97 }
98}
99
100pub(super) fn arm_appservice_wake(wakes: &mut WakeQueue, dest: Destination, tries: u32) {
101 if is_armed(wakes, &dest) {
102 return;
103 }
104
105 arm_wake_in(wakes, dest, appservice_delay(tries));
106}
107
108fn appservice_delay(tries: u32) -> Duration {
109 let exponent = tries.min(APPSERVICE_RETRY_MAX_SECS.ilog(APPSERVICE_RETRY_BASE));
110
111 Duration::from_secs(APPSERVICE_RETRY_BASE.pow(exponent))
112}
113
114#[implement(Service)]
115async fn handle_federation_wake<'a>(
116 &'a self,
117 server: OwnedServerName,
118 futures: &mut SendingFutures<'a>,
119 statuses: &mut TransactionStatuses,
120 wakes: &mut WakeQueue,
121) {
122 let should_attempt = self
123 .services
124 .federation
125 .should_attempt(&server)
126 .await;
127
128 let dest = Destination::Federation(server);
129
130 match should_attempt {
131 | ShouldAttempt::No { earliest_retry } => arm_wake(wakes, dest, earliest_retry),
132 | _ => {
133 let msg = Msg {
134 dest,
135 event: SendingEvent::Flush,
136 queue_id: Vec::new(),
137 };
138
139 self.handle_request(msg, futures, statuses, wakes)
140 .await;
141 },
142 }
143}
144
145#[implement(Service)]
146pub(super) async fn arm_federation_wake(
147 &self,
148 server: OwnedServerName,
149 tries: u32,
150 wakes: &mut WakeQueue,
151) {
152 let verdict = self
153 .services
154 .federation
155 .should_attempt(&server)
156 .await;
157
158 let dest = Destination::Federation(server);
159
160 match verdict {
161 | ShouldAttempt::No { earliest_retry } => arm_wake(wakes, dest, earliest_retry),
162 | _ => {
163 let delay = exponential_backoff_remaining_secs(
164 self.server.config.sender_timeout,
165 self.server.config.sender_retry_backoff_limit,
166 Duration::ZERO,
167 tries,
168 )
169 .unwrap_or_default();
170
171 arm_wake_in(wakes, dest, delay);
172 },
173 }
174}
175
176#[implement(Service)]
177pub(super) fn arm_push_wake(
178 &self,
179 dest: Destination,
180 error: &Error,
181 statuses: &TransactionStatuses,
182 wakes: &mut WakeQueue,
183) {
184 let Some(status @ TransactionStatus::Failed { tries, .. }) = statuses.get(&dest) else {
185 return;
186 };
187
188 let delay = self
189 .push_backoff_remaining(Some(status))
190 .unwrap_or_default();
191
192 let (deadline, retry_in) = wake_deadline(delay);
193
194 record_push_failure(&dest, error, *tries, retry_in);
195 wakes.push(Reverse((deadline, dest)));
196}
197
198#[implement(Service)]
199#[inline]
200pub(super) fn push_backoff_remaining(
201 &self,
202 status: Option<&TransactionStatus>,
203) -> Option<Duration> {
204 let Some(TransactionStatus::Failed { tries, last }) = status else {
205 return None;
206 };
207
208 exponential_backoff_remaining_secs(
209 self.server.config.sender_timeout,
210 self.server.config.sender_retry_backoff_limit,
211 last.elapsed(),
212 *tries,
213 )
214}
215
216fn record_push_failure(dest: &Destination, error: &Error, tries: u32, retry_in: Duration) {
217 let Destination::Push(user_id, pushkey) = dest else {
218 return;
219 };
220
221 match tries {
222 | PUSH_FAILURE_STREAK => error!(
223 %user_id,
224 %pushkey,
225 streak = tries,
226 retry_in_seconds = retry_in.as_secs(),
227 chain = %error_chain(error),
228 "Push notifications for this pusher are not being delivered",
229 ),
230 | _ => warn!(
231 %user_id,
232 %pushkey,
233 streak = tries,
234 retry_in_seconds = retry_in.as_secs(),
235 chain = %error_chain(error),
236 "Push transaction failed",
237 ),
238 }
239}
240
241pub(super) fn arm_park_wake(wakes: &mut WakeQueue, dest: Destination, until: u64) {
246 let deadline = park_deadline(until);
247
248 if wakes
249 .iter()
250 .any(|Reverse((due, armed))| armed == &dest && *due <= deadline)
251 .is_false()
252 {
253 wakes.push(Reverse((deadline, dest)));
254 }
255}
256
257fn park_deadline(until: u64) -> Instant {
258 let now = Instant::now();
259 let delay = Duration::from_secs(until.saturating_sub(now_secs()));
260
261 now.checked_add(delay)
262 .or_else(|| now.checked_add(WAKE_OVERFLOW_DELAY))
263 .unwrap_or(now)
264}
265
266pub(super) fn is_armed(wakes: &WakeQueue, dest: &Destination) -> bool {
267 wakes
268 .iter()
269 .any(|Reverse((_, armed))| armed == dest)
270}
271
272pub(super) fn arm_wake(wakes: &mut WakeQueue, dest: Destination, earliest_retry: SystemTime) {
273 let delay = earliest_retry
274 .duration_since(SystemTime::now())
275 .unwrap_or_default();
276
277 arm_wake_in(wakes, dest, delay);
278}
279
280pub(super) fn arm_wake_in(wakes: &mut WakeQueue, dest: Destination, delay: Duration) {
281 let (deadline, _) = wake_deadline(delay);
282
283 wakes.push(Reverse((deadline, dest)));
284}
285
286fn wake_deadline(delay: Duration) -> (Instant, Duration) {
287 let delay = delay.max(Duration::from_secs(1));
289
290 let jitter = rand_secs(0..delay.as_secs().max(3));
292 let now = Instant::now();
293 let scheduled = delay.saturating_add(jitter);
294 let deadline = now
295 .checked_add(scheduled)
296 .or_else(|| now.checked_add(WAKE_OVERFLOW_DELAY))
297 .unwrap_or(now);
298
299 let scheduled = deadline.saturating_duration_since(now);
300
301 (deadline, scheduled)
302}