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, ¤t_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}