Skip to main content

dash_sdk/
sync.rs

1pub use dash_async::{block_on, AsyncError};
2
3use crate::error::Error;
4use rs_dapi_client::{
5    transport::sleep, update_address_ban_status, AddressList, CanRetry, ExecutionResult,
6    RequestSettings,
7};
8use std::future::Future;
9use std::time::Duration;
10
11impl From<AsyncError> for crate::Error {
12    fn from(error: AsyncError) -> Self {
13        Self::ContextProviderError(error.into())
14    }
15}
16
17/// Retry the provided closure.
18///
19/// This function is used to retry async code. It takes into account number of retries already executed by lower
20/// layers and stops retrying once the maximum number of retries is reached.
21///
22/// The `settings` should contain maximum number of retries that should be executed. In case of failure, total number of
23/// requests sent is expected to be at least `settings.retries + 1` (initial request + `retries` configured in settings).
24/// The actual number of requests sent can be higher, as the lower layers can retry the request multiple times.
25///
26/// `future_factory_fn` should be a `FnMut()` closure that returns a future that should be retried.
27/// It takes [`RequestSettings`] as an argument and returns [`ExecutionResult`].
28/// Retry mechanism can change [`RequestSettings`] between invocations of the `future_factory_fn` closure
29/// to limit the number of retries for lower layers.
30///
31/// ## Parameters
32///
33/// - `address_list` - list of addresses to be used for the requests.
34/// - `settings` - global settings with any request-specific settings overrides applied.
35/// - `future_factory_fn` - closure that returns a future that should be retried. It should take [`RequestSettings`] as
36///   an argument and return [`ExecutionResult`].
37///
38/// ## Returns
39///
40/// Returns future that resolves to [`ExecutionResult`].
41///
42/// ## Example
43///
44/// ```rust
45/// # use dash_sdk::RequestSettings;
46/// # use dash_sdk::error::{Error,StaleNodeError};
47/// # use rs_dapi_client::{ExecutionResult, ExecutionError};
48/// async fn retry_test_function(settings: RequestSettings) -> ExecutionResult<(), dash_sdk::Error> {
49/// // do something
50///     Err(ExecutionError {
51///         inner: Error::StaleNode(StaleNodeError::Height{
52///             expected_height: 10,
53///             received_height: 3,
54///             tolerance_blocks: 1,
55///         }),
56///        retries: 0,
57///       address: None,
58///    })
59/// }
60/// #[tokio::main]
61///     async fn main() {
62///     let address_list = rs_dapi_client::AddressList::default();
63///     let global_settings = RequestSettings::default();
64///     dash_sdk::sync::retry(&address_list, global_settings, retry_test_function).await.expect_err("should fail");
65/// }
66/// ```
67///
68/// ## Troubleshooting
69///
70/// Compiler error: `no method named retry found for closure`:
71/// - ensure returned value is [`ExecutionResult`].
72/// - consider adding `.await` at the end of the closure.
73pub async fn retry<Fut, FutureFactoryFn, R>(
74    address_list: &AddressList,
75    settings: RequestSettings,
76    future_factory_fn: FutureFactoryFn,
77) -> ExecutionResult<R, Error>
78where
79    Fut: Future<Output = ExecutionResult<R, Error>>,
80    FutureFactoryFn: FnMut(RequestSettings) -> Fut,
81    R: Send,
82{
83    retry_with_additional_error(address_list, settings, future_factory_fn, |_| false).await
84}
85
86/// Retry an operation-specific rejection only when its responding node can be
87/// excluded. This does not change the error's global retry classification.
88/// Callers must restrict the predicate to definitive, safe-to-repeat failures.
89pub(crate) async fn retry_with_additional_error<Fut, FutureFactoryFn, R, AdditionalError>(
90    address_list: &AddressList,
91    settings: RequestSettings,
92    mut future_factory_fn: FutureFactoryFn,
93    additional_error: AdditionalError,
94) -> ExecutionResult<R, Error>
95where
96    Fut: Future<Output = ExecutionResult<R, Error>>,
97    FutureFactoryFn: FnMut(RequestSettings) -> Fut,
98    R: Send,
99    AdditionalError: Fn(&Error) -> bool,
100{
101    let max_retries = settings.retries.unwrap_or_default();
102    let mut total_retries: usize = 0;
103    let mut current_settings = settings;
104
105    // Store the last meaningful error (not "no available addresses")
106    // so we can return it if we exhaust all addresses
107    let mut last_meaningful_error: Option<rs_dapi_client::ExecutionError<Error>> = None;
108
109    loop {
110        let result = future_factory_fn(current_settings).await;
111
112        // Ban or unban the address based on the result
113        update_address_ban_status(address_list, &result, &current_settings.finalize());
114
115        match result {
116            Ok(response) => return Ok(response),
117            Err(error) => {
118                // Check if this is a "no available addresses" error and we have a previous meaningful error
119                if error.is_no_available_addresses() {
120                    if let Some(prev_error) = last_meaningful_error.take() {
121                        tracing::error!(
122                            retry = total_retries,
123                            max_retries,
124                            error = ?prev_error,
125                            "no addresses available to retry"
126                        );
127                        // Wrap the last meaningful error in NoAvailableAddresses
128                        return Err(rs_dapi_client::ExecutionError {
129                            inner: Error::NoAvailableAddressesToRetry(Box::new(prev_error.inner)),
130                            retries: total_retries,
131                            address: prev_error.address,
132                        });
133                    }
134                    // No previous error, return the "no available addresses" error as-is
135                    return Err(error);
136                }
137
138                // Count requests sent in this attempt
139                let requests_sent = error.retries + 1;
140                total_retries += requests_sent;
141
142                let retry_additional_error = additional_error(&error.inner);
143                if !error.can_retry() && !retry_additional_error {
144                    // Non-retryable error, return immediately
145                    let mut final_error = error;
146                    final_error.retries = total_retries;
147                    return Err(final_error);
148                }
149
150                if total_retries > max_retries {
151                    // Exceeded max retries
152                    tracing::error!(
153                        retry = total_retries,
154                        max_retries,
155                        error = ?error,
156                        "no more retries left, giving up"
157                    );
158                    let mut final_error = error;
159                    final_error.retries = total_retries;
160                    return Err(final_error);
161                }
162
163                if retry_additional_error {
164                    // This rejection does not establish a health failure. Use
165                    // a short flat exclusion, never the exponential health ladder.
166                    // A node can retain rejected transaction hashes. Never
167                    // resend this rejection to the same node, including when
168                    // the caller disabled banning or the address is unknown.
169                    // Only exclude it when a retry to another node is possible;
170                    // a single-node client must remain usable after failure.
171                    let excluded = current_settings.finalize().ban_failed_address
172                        && error.address.as_ref().is_some_and(|address| {
173                            address_list
174                                .get_live_addresses()
175                                .iter()
176                                .any(|candidate| candidate != address)
177                                && address_list.ban_for(
178                                    address,
179                                    Duration::from_secs(2),
180                                    Some(error.to_string()),
181                                )
182                        });
183                    if !excluded || address_list.get_live_addresses().is_empty() {
184                        tracing::debug!(node = ?error.address, "DPNS failover stopped: no safely excluded alternative");
185                        let mut final_error = error;
186                        final_error.retries = total_retries;
187                        return Err(final_error);
188                    }
189                }
190
191                // Log retry decision (matches original `when()` callback)
192                tracing::warn!(
193                    retry = total_retries,
194                    max_retries,
195                    error = ?error,
196                    "retrying request"
197                );
198
199                // Update settings for next retry - limit retries for lower layer
200                current_settings.retries = Some(max_retries.saturating_sub(total_retries));
201
202                // Small delay to avoid spamming (we use different server, so no real delay needed)
203                // Log before sleep (matches original `notify()` callback)
204                let delay = Duration::from_millis(10);
205                tracing::warn!(duration = ?delay, error = ?error, "request failed, retrying");
206
207                // Store this as the last meaningful error before retrying
208                last_meaningful_error = Some(error);
209
210                sleep(delay).await;
211            }
212        }
213    }
214}
215
216#[cfg(test)]
217mod test {
218    use super::*;
219    use rs_dapi_client::ExecutionError;
220    use std::sync::{
221        atomic::{AtomicUsize, Ordering},
222        Arc,
223    };
224
225    use crate::error::StaleNodeError;
226    use rs_dapi_client::DapiClientError;
227
228    async fn retry_test_function(
229        settings: RequestSettings,
230        counter: Arc<AtomicUsize>,
231    ) -> ExecutionResult<(), Error> {
232        // num or retries increases with each call
233        let retries = counter.load(Ordering::Relaxed);
234        let retries = if settings.retries.unwrap_or_default() < retries {
235            settings.retries.unwrap_or_default()
236        } else {
237            retries
238        };
239
240        // we sent 1 initial request plus `retries` retries
241        counter.fetch_add(1 + retries, Ordering::Relaxed);
242
243        Err(ExecutionError {
244            inner: Error::StaleNode(StaleNodeError::Height {
245                expected_height: 100,
246                received_height: 50,
247                tolerance_blocks: 1,
248            }),
249            retries,
250            address: Some("http://localhost".parse().expect("valid address")),
251        })
252    }
253
254    #[test_case::test_matrix([1,2,3,5,7,8,10,11,23,49, usize::MAX])]
255    #[tokio::test]
256    async fn test_retry(expected_requests: usize) {
257        for _ in 0..1 {
258            let counter = Arc::new(AtomicUsize::new(0));
259
260            let address_list = AddressList::default();
261
262            // we retry 5 times, and expect 5 retries + 1 initial request
263            let mut global_settings = RequestSettings::default();
264            global_settings.retries = Some(expected_requests - 1);
265
266            let closure = |s| {
267                let counter = counter.clone();
268                retry_test_function(s, counter)
269            };
270
271            retry(&address_list, global_settings, closure)
272                .await
273                .expect_err("should fail");
274
275            assert_eq!(
276                counter.load(Ordering::Relaxed),
277                expected_requests,
278                "test failed for expected {} requests",
279                expected_requests
280            );
281        }
282    }
283
284    /// Test that when we get "no available addresses" error, we return the last meaningful error
285    /// wrapped in NoAvailableAddresses.
286    #[tokio::test]
287    async fn test_retry_returns_last_meaningful_error_on_no_addresses() {
288        let call_count = Arc::new(AtomicUsize::new(0));
289        let address_list = AddressList::default();
290
291        let mut settings = RequestSettings::default();
292        settings.retries = Some(5);
293
294        let call_count_clone = call_count.clone();
295        let closure = move |_settings: RequestSettings| {
296            let count = call_count_clone.fetch_add(1, Ordering::Relaxed);
297            async move {
298                if count == 0 {
299                    Err(ExecutionError {
300                        inner: Error::StaleNode(StaleNodeError::Height {
301                            expected_height: 100,
302                            received_height: 50,
303                            tolerance_blocks: 1,
304                        }),
305                        retries: 0,
306                        address: Some("http://localhost:1".parse().unwrap()),
307                    })
308                } else {
309                    Err(ExecutionError {
310                        inner: Error::DapiClientError(DapiClientError::NoAvailableAddresses),
311                        retries: 0,
312                        address: None,
313                    })
314                }
315            }
316        };
317
318        let result: ExecutionResult<(), Error> = retry(&address_list, settings, closure).await;
319
320        let error = result.expect_err("should fail");
321        match &error.inner {
322            Error::NoAvailableAddressesToRetry(inner) => {
323                assert!(
324                    matches!(**inner, Error::StaleNode(_)),
325                    "inner error should be StaleNode, got: {:?}",
326                    inner
327                );
328            }
329            _ => panic!(
330                "expected NoAvailableAddresses error, got: {:?}",
331                error.inner
332            ),
333        }
334        assert_eq!(
335            call_count.load(Ordering::Relaxed),
336            2,
337            "should have called twice"
338        );
339    }
340
341    /// Test that if we get "no available addresses" on the first call (no previous error),
342    /// we still return it as-is (not wrapped).
343    #[tokio::test]
344    async fn test_retry_returns_no_addresses_if_no_previous_error() {
345        let address_list = AddressList::default();
346
347        let mut settings = RequestSettings::default();
348        settings.retries = Some(5);
349
350        let closure = move |_settings: RequestSettings| async move {
351            Err(ExecutionError {
352                inner: Error::DapiClientError(DapiClientError::NoAvailableAddresses),
353                retries: 0,
354                address: None,
355            })
356        };
357
358        let result: ExecutionResult<(), Error> = retry(&address_list, settings, closure).await;
359
360        let error = result.expect_err("should fail");
361        assert!(
362            matches!(
363                error.inner,
364                Error::DapiClientError(DapiClientError::NoAvailableAddresses)
365            ),
366            "should return 'no available addresses' when there's no previous meaningful error, got: {:?}",
367            error.inner
368        );
369    }
370}