Skip to main content

tuwunel_service/sending/sender/
wake.rs

1#[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
241/// Wake a destination holding only parked rooms when its earliest park expires.
242///
243/// Unlike a retry timer, it is not jittered, and it is skipped when a wake for
244/// the destination is already armed to fire no later.
245pub(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	// Floor the delay at 1s so clock steps and past deadlines wake promptly.
288	let delay = delay.max(Duration::from_secs(1));
289
290	// Jitter by up to another delay-width (3s floor) so a backoff tier trickles back.
291	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}