heddle_thread_api/
reopen.rs1use std::future::Future;
4
5use api::{
6 heddle::api::common::{CallFailure, CallFailureCode},
7 v2::client::ClientError,
8};
9
10use crate::transport::{Error, RemoteFailure};
11
12pub(crate) const ATTEMPTS: u32 = 5;
14
15pub fn is_reopen_retryable(failure: &CallFailure) -> bool {
18 retryable(failure.code(), &failure.message)
19}
20
21pub(crate) fn remote_failure_is_reopen_retryable(failure: &RemoteFailure) -> bool {
22 retryable(
23 CallFailureCode::try_from(failure.code).unwrap_or_default(),
24 &failure.message,
25 )
26}
27
28pub(crate) fn error_is_reopen_retryable(error: &Error) -> bool {
29 match error {
30 Error::Remote(failure) => remote_failure_is_reopen_retryable(failure),
31 _ => false,
32 }
33}
34
35pub(crate) fn client_error_is_reopen_retryable(error: &ClientError<Error>) -> bool {
36 match error {
37 ClientError::Transport(error) => error_is_reopen_retryable(error),
38 _ => false,
39 }
40}
41
42pub(crate) trait ReopenRetryable {
43 fn is_reopen_retryable(&self) -> bool;
44}
45
46impl ReopenRetryable for Error {
47 fn is_reopen_retryable(&self) -> bool {
48 error_is_reopen_retryable(self)
49 }
50}
51
52impl ReopenRetryable for ClientError<Error> {
53 fn is_reopen_retryable(&self) -> bool {
54 client_error_is_reopen_retryable(self)
55 }
56}
57
58pub(crate) async fn retry<T, E, F, Fut>(mut op: F) -> Result<T, E>
59where
60 F: FnMut() -> Fut,
61 Fut: Future<Output = Result<T, E>>,
62 E: ReopenRetryable,
63{
64 let mut attempt = 0;
65 loop {
66 match op().await {
67 Ok(value) => return Ok(value),
68 Err(error) if error.is_reopen_retryable() && attempt + 1 < ATTEMPTS => {
69 attempt += 1;
70 backoff(attempt).await;
71 }
72 Err(error) => return Err(error),
73 }
74 }
75}
76
77pub(crate) async fn backoff(attempt: u32) {
78 let millis = 20u64.saturating_mul(u64::from(attempt));
79 #[cfg(feature = "iroh")]
80 {
81 tokio::time::sleep(std::time::Duration::from_millis(millis)).await;
82 }
83 #[cfg(not(feature = "iroh"))]
84 {
85 let _ = millis;
86 }
87}
88
89fn retryable(code: CallFailureCode, message: &str) -> bool {
90 let message = message.to_ascii_lowercase();
91 let authority = message.contains("authority changed")
92 || message.contains("reopen exact selection")
93 || message.contains("authorization changed")
94 || message.contains("is no longer valid");
95 let pool = message.contains("pool timed out");
96 match code {
97 CallFailureCode::Aborted if authority || pool => true,
98 CallFailureCode::Unavailable | CallFailureCode::ResourceExhausted if pool => true,
99 _ => false,
100 }
101}
102
103#[cfg(test)]
104mod tests {
105 use super::*;
106
107 fn failure(code: CallFailureCode, message: &str) -> CallFailure {
108 CallFailure {
109 code: code as i32,
110 message: message.into(),
111 ..Default::default()
112 }
113 }
114
115 #[test]
116 fn classifies_weft_v2_authority_change_and_pool_timeout_signals() {
117 assert!(is_reopen_retryable(&failure(
118 CallFailureCode::Aborted,
119 "material authority changed; reopen exact selection",
120 )));
121 assert!(is_reopen_retryable(&failure(
122 CallFailureCode::Aborted,
123 "authorization changed",
124 )));
125 assert!(is_reopen_retryable(&failure(
126 CallFailureCode::Aborted,
127 "authorization is no longer valid",
128 )));
129 assert!(is_reopen_retryable(&failure(
130 CallFailureCode::Unavailable,
131 "pool timed out while waiting for an open connection",
132 )));
133 assert!(is_reopen_retryable(&failure(
134 CallFailureCode::ResourceExhausted,
135 "pool timed out",
136 )));
137 assert!(is_reopen_retryable(&failure(
138 CallFailureCode::Aborted,
139 "pool timed out",
140 )));
141 }
142
143 #[test]
144 fn non_transient_failures_are_not_reopen_retryable() {
145 assert!(!is_reopen_retryable(&failure(
146 CallFailureCode::PermissionDenied,
147 "material authority changed; reopen exact selection",
148 )));
149 assert!(!is_reopen_retryable(&failure(
150 CallFailureCode::Aborted,
151 "review required",
152 )));
153 assert!(!is_reopen_retryable(&failure(
154 CallFailureCode::NotFound,
155 "selected source unavailable",
156 )));
157 assert!(!is_reopen_retryable(&failure(
158 CallFailureCode::Unavailable,
159 "endpoint draining",
160 )));
161 assert!(!is_reopen_retryable(&failure(
162 CallFailureCode::Internal,
163 "pool timed out",
164 )));
165 }
166}