Skip to main content

heddle_thread_api/
reopen.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Bounded reopen of a stale exact selection after weft v2 retryable signals.
3use std::future::Future;
4
5use api::{
6    heddle::api::common::{CallFailure, CallFailureCode},
7    v2::client::ClientError,
8};
9
10use crate::transport::{Error, RemoteFailure};
11
12/// Initial attempt plus four reopens of a fresh exact selection.
13pub(crate) const ATTEMPTS: u32 = 5;
14
15/// Weft fail-fasts a stale exact selection with these retryable signals.
16/// A bare resend of the same selection would observe the same `Aborted`.
17pub 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}