tuwunel_service/sending/sender/
response.rs1use 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
70fn 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
84fn 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
144fn fail_status(dest: &Destination, statuses: &mut TransactionStatuses) -> RetryAction {
148 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
174fn 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}