Skip to main content

tuwunel_service/sending/sender/
response.rs

1use std::time::Instant;
2
3use tuwunel_core::{Error, debug, debug_info, error::error_chain, implement, warn};
4
5use super::{
6	NewEvents, RetryAction, SendingFutures, TransactionStatus, TransactionStatuses, WakeQueue,
7	dispatch::{Completion, SendingError},
8	select::Selection,
9	split::Split,
10	wake::arm_appservice_wake,
11};
12use crate::{
13	federation::is_content_rejection,
14	sending::{Destination, Service},
15};
16
17#[implement(Service)]
18#[tracing::instrument(name = "response", level = "debug", skip_all)]
19pub(super) async fn handle_response<'a>(
20	&'a self,
21	Completion { result, keys, split }: Completion,
22	futures: &mut SendingFutures<'a>,
23	statuses: &mut TransactionStatuses,
24	wakes: &mut WakeQueue,
25) {
26	match result {
27		| Ok(dest) => {
28			let _cork = self.db.db.cork();
29
30			self.db.delete_active_requests(&keys);
31			self.handle_response_ok(dest, split, futures, statuses, wakes)
32				.await;
33		},
34		| Err(error) =>
35			self.handle_response_err(error, split, futures, statuses, wakes)
36				.await,
37	}
38}
39
40#[implement(Service)]
41async fn handle_response_ok<'a>(
42	&'a self,
43	dest: Destination,
44	split: Option<Split>,
45	futures: &mut SendingFutures<'a>,
46	statuses: &mut TransactionStatuses,
47	wakes: &mut WakeQueue,
48) {
49	log_recovery(&dest, statuses);
50
51	if let Some(split) = split
52		&& let Some((items, split)) = self.split_delivered(&dest, split).await
53	{
54		run_status(&dest, statuses);
55		futures.push(self.send_events(dest, items, Some(split)));
56		return;
57	}
58
59	let next = match &dest {
60		| Destination::Federation(server) =>
61			self.federation_batch(&dest, server, NewEvents::new())
62				.await,
63		| _ => Selection::Events(self.resume_queued(&dest, &[]).await),
64	};
65
66	run_status(&dest, statuses);
67	self.schedule_events(dest, next, futures, statuses, wakes);
68}
69
70/// Log a delivery that ends a destination's streak of failed transactions.
71///
72/// Reads the streak before `run_status` resets it, so it runs first. Push,
73/// whose retry in flight is `Retrying`, reports its own failures instead.
74fn log_recovery(dest: &Destination, statuses: &TransactionStatuses) {
75	if let Some(
76		&(TransactionStatus::Running { tries } | TransactionStatus::RunningForceRetry { tries }),
77	) = statuses.get(dest)
78		&& tries > 0
79	{
80		debug_info!(?dest, streak = tries, "Transaction delivered after failures");
81	}
82}
83
84/// Mark a destination's transaction as running, if one is tracked.
85fn run_status(dest: &Destination, statuses: &mut TransactionStatuses) {
86	if let Some(status) = statuses.get_mut(dest) {
87		*status = TransactionStatus::Running { tries: 0 };
88	}
89}
90
91#[implement(Service)]
92async fn handle_response_err<'a>(
93	&'a self,
94	(dest, error): SendingError,
95	split: Option<Split>,
96	futures: &mut SendingFutures<'a>,
97	statuses: &mut TransactionStatuses,
98	wakes: &mut WakeQueue,
99) {
100	let retry_action = fail_status(&dest, statuses);
101
102	log_failure(&dest, &error, statuses);
103
104	let tries = match statuses.get(&dest) {
105		| Some(TransactionStatus::Retrying { tries }) => *tries,
106		| _ => 0,
107	};
108
109	match dest {
110		| dest @ Destination::Push(..) => self.arm_push_wake(dest, &error, statuses, wakes),
111		| dest @ Destination::Appservice(_) => {
112			let forced = matches!(retry_action, RetryAction::Force).then(|| dest.clone());
113
114			arm_appservice_wake(wakes, dest, tries);
115			if let Some(dest) = forced {
116				self.handle_force_retry(dest, futures, statuses, wakes)
117					.await;
118			}
119		},
120		| Destination::Federation(server) => {
121			if tries > 0 {
122				self.stalled
123					.lock()
124					.expect("locked")
125					.insert(server.clone(), Some(Instant::now()));
126			}
127
128			let (split, tries) = self
129				.split_failure(&server, &error, split, tries)
130				.await;
131
132			if let Some(split) = split {
133				let status = TransactionStatus::Splitting { tries, split };
134
135				statuses.insert(Destination::Federation(server.clone()), status);
136			}
137
138			self.arm_federation_wake(server, tries, wakes)
139				.await;
140		},
141	}
142}
143
144/// Advance a destination's status after a failed transaction.
145///
146/// Reports whether a forced retry was requested while the transaction ran.
147fn fail_status(dest: &Destination, statuses: &mut TransactionStatuses) -> RetryAction {
148	// Push records its local clock; other destinations wait for a replay trigger.
149	let push = matches!(dest, Destination::Push(..));
150
151	let Some(status) = statuses.get_mut(dest) else {
152		return RetryAction::None;
153	};
154
155	let (tries, retry_action) = match status {
156		| TransactionStatus::Pending => (1, RetryAction::None),
157		| TransactionStatus::RunningForceRetry { tries } =>
158			(tries.saturating_add(1), RetryAction::Force),
159		| TransactionStatus::Running { tries }
160		| TransactionStatus::Failed { tries, .. }
161		| TransactionStatus::Splitting { tries, .. }
162		| TransactionStatus::Retrying { tries } => (tries.saturating_add(1), RetryAction::None),
163	};
164
165	*status = if push {
166		TransactionStatus::Failed { tries, last: Instant::now() }
167	} else {
168		TransactionStatus::Retrying { tries }
169	};
170
171	retry_action
172}
173
174/// Log a failed transaction, at warn when a content rejection starts a
175/// destination's streak.
176///
177/// A content rejection records no peer backoff and its batch replays unchanged,
178/// so without the warning a stuck destination is silent at default levels.
179/// Call it after `fail_status`, whose advanced status it reads.
180fn log_failure(dest: &Destination, error: &Error, statuses: &TransactionStatuses) {
181	match statuses.get(dest) {
182		| Some(TransactionStatus::Retrying { tries: 1 }) if is_content_rejection(error) =>
183			warn!(?dest, chain = %error_chain(error), "Transaction failed"),
184		| _ => debug!(?dest, chain = %error_chain(error), "Transaction failed"),
185	}
186}
187
188#[implement(Service)]
189pub(super) async fn handle_force_retry<'a>(
190	&'a self,
191	dest: Destination,
192	futures: &mut SendingFutures<'a>,
193	statuses: &mut TransactionStatuses,
194	wakes: &mut WakeQueue,
195) {
196	let Ok(selection) = self
197		.select_events(&dest, NewEvents::new(), statuses)
198		.await
199	else {
200		return;
201	};
202
203	self.schedule_events(dest, selection, futures, statuses, wakes);
204}