Skip to main content

rs_dapi_client/
dapi_client.rs

1//! [DapiClient] definition.
2
3use dapi_grpc::mock::Mockable;
4use dapi_grpc::tonic::async_trait;
5#[cfg(not(target_arch = "wasm32"))]
6use dapi_grpc::tonic::transport::Certificate;
7use std::fmt::{Debug, Display};
8use std::time::Duration;
9use tracing::Instrument;
10
11use crate::address_list::AddressListError;
12use crate::connection_pool::ConnectionPool;
13use crate::request_settings::AppliedRequestSettings;
14use crate::transport::{self, TransportError};
15use crate::{
16    transport::{TransportClient, TransportRequest},
17    AddressList, CanRetry, DapiRequestExecutor, ExecutionError, ExecutionResponse, ExecutionResult,
18    RequestSettings,
19};
20
21/// Intended minimum for the Envoy-advertised `RateLimit-Reset` ban duration.
22/// Note: the `> 0` filter applied before the clamp already rejects 0 → `None`,
23/// so this constant never actively clamps the lower bound — it documents intent
24/// (the smallest meaningful reset is 1 s) and acts as the `.clamp(MIN, MAX)`
25/// lower argument for clarity.
26pub(crate) const MIN_RATE_LIMIT_BAN_SECS: u64 = 1;
27/// Ceiling for the Envoy-advertised `RateLimit-Reset` ban duration.
28/// Prevents a misconfigured or hostile header from parking a healthy node for
29/// an unreasonably long time.
30pub(crate) const MAX_RATE_LIMIT_BAN_SECS: u64 = 600;
31
32/// General DAPI request error type.
33#[derive(Debug, thiserror::Error, Clone)]
34#[cfg_attr(feature = "mocks", derive(serde::Serialize, serde::Deserialize))]
35pub enum DapiClientError {
36    /// The error happened on transport layer
37    #[error("transport error: {0}")]
38    Transport(
39        #[cfg_attr(feature = "mocks", serde(with = "dapi_grpc::mock::serde_mockable"))]
40        TransportError,
41    ),
42    /// There are no valid DAPI addresses to use.
43    #[error("no available addresses to use")]
44    NoAvailableAddresses,
45    /// All available addresses have been exhausted (banned due to errors).
46    /// Contains the last meaningful error that caused addresses to be banned.
47    #[error("no available addresses to retry, last error: {0}")]
48    NoAvailableAddressesToRetry(
49        #[cfg_attr(feature = "mocks", serde(with = "dapi_grpc::mock::serde_mockable"))]
50        Box<TransportError>,
51    ),
52    /// [AddressListError] errors
53    #[error("address list error: {0}")]
54    AddressList(AddressListError),
55
56    #[cfg(feature = "mocks")]
57    #[error("mock error: {0}")]
58    /// Error happened in mock client
59    Mock(#[from] crate::mock::MockError),
60}
61
62impl CanRetry for DapiClientError {
63    fn can_retry(&self) -> bool {
64        use DapiClientError::*;
65        match self {
66            NoAvailableAddresses => false,
67            NoAvailableAddressesToRetry(_) => false,
68            Transport(transport_error) => transport_error.can_retry(),
69            AddressList(_) => false,
70            #[cfg(feature = "mocks")]
71            Mock(_) => false,
72        }
73    }
74
75    fn is_no_available_addresses(&self) -> bool {
76        matches!(
77            self,
78            DapiClientError::NoAvailableAddresses | DapiClientError::NoAvailableAddressesToRetry(_)
79        )
80    }
81
82    fn rate_limit_ban_duration(&self) -> Option<Duration> {
83        match self {
84            DapiClientError::Transport(te) => te.rate_limit_ban_duration(),
85            _ => None,
86        }
87    }
88}
89
90/// Serialization of [DapiClientError].
91///
92/// We need to do manual serialization because of the generic type parameter which doesn't support serde derive.
93impl Mockable for DapiClientError {
94    #[cfg(feature = "mocks")]
95    fn mock_serialize(&self) -> Option<Vec<u8>> {
96        Some(serde_json::to_vec(self).expect("serialize DAPI client error"))
97    }
98
99    #[cfg(feature = "mocks")]
100    fn mock_deserialize(data: &[u8]) -> Option<Self> {
101        Some(serde_json::from_slice(data).expect("deserialize DAPI client error"))
102    }
103}
104
105/// Access point to DAPI.
106#[derive(Debug, Clone)]
107pub struct DapiClient {
108    address_list: AddressList,
109    settings: RequestSettings,
110    pool: ConnectionPool,
111    #[cfg(not(target_arch = "wasm32"))]
112    /// Certificate Authority certificate to use for verifying the server's certificate.
113    pub ca_certificate: Option<Certificate>,
114    #[cfg(feature = "dump")]
115    pub(crate) dump_dir: Option<std::path::PathBuf>,
116}
117
118impl DapiClient {
119    /// Initialize new [DapiClient] and optionally override default settings.
120    pub fn new(address_list: AddressList, settings: RequestSettings) -> Self {
121        // multiply by 3 as we need to store core and platform addresses, and we want some spare capacity just in case
122        let address_count = 3 * address_list.len();
123
124        Self {
125            address_list,
126            settings,
127            pool: ConnectionPool::new(address_count),
128            #[cfg(feature = "dump")]
129            dump_dir: None,
130            #[cfg(not(target_arch = "wasm32"))]
131            ca_certificate: None,
132        }
133    }
134
135    /// Set CA certificate to use when verifying the server's certificate.
136    ///
137    /// # Arguments
138    ///
139    /// * `pem_ca_cert` - CA certificate in PEM format.
140    ///
141    /// # Returns
142    /// [DapiClient] with CA certificate set.
143    #[cfg(not(target_arch = "wasm32"))]
144    pub fn with_ca_certificate(mut self, ca_cert: Certificate) -> Self {
145        self.ca_certificate = Some(ca_cert);
146
147        self
148    }
149
150    /// Return the [DapiClient] address list.
151    pub fn address_list(&self) -> &AddressList {
152        &self.address_list
153    }
154
155    /// Get all non-banned addresses from the address list.
156    ///
157    /// Returns a vector of addresses that are not currently banned or whose ban period has expired.
158    /// This is useful for diagnostics, monitoring, or when you need to know which DAPI nodes are
159    /// currently available for making requests.
160    ///
161    /// # Examples
162    ///
163    /// ```no_run
164    /// use rs_dapi_client::{DapiClient, AddressList, RequestSettings};
165    ///
166    /// let address_list = "http://127.0.0.1:3000,http://127.0.0.1:3001".parse().unwrap();
167    /// let client = DapiClient::new(address_list, RequestSettings::default());
168    ///
169    /// // Get all currently available (non-banned) addresses
170    /// let live_addresses = client.get_live_addresses();
171    /// println!("Available DAPI nodes: {}", live_addresses.len());
172    /// ```
173    pub fn get_live_addresses(&self) -> Vec<crate::Address> {
174        self.address_list.get_live_addresses()
175    }
176}
177
178/// Ban address in case of retryable error or unban it
179/// if it was banned, and the request was successful.
180pub fn update_address_ban_status<R, E>(
181    address_list: &AddressList,
182    result: &ExecutionResult<R, E>,
183    applied_settings: &AppliedRequestSettings,
184) where
185    E: CanRetry + Display + Debug,
186{
187    match &result {
188        Ok(response) => {
189            // Unban the address if it was banned and node responded successfully this time
190            if address_list.is_banned(&response.address) {
191                if address_list.unban(&response.address) {
192                    tracing::debug!(address = ?response.address, "unban successfully responded address {}", response.address);
193                } else {
194                    // The address might be already removed from the list
195                    // by background process (i.e., SML update), and it's fine.
196                    tracing::debug!(
197                        address = ?response.address,
198                        "unable to unban address {} because it's not in the list anymore",
199                        response.address
200                    );
201                }
202            }
203        }
204        Err(error) => {
205            if error.can_retry() {
206                if let Some(address) = error.address.as_ref() {
207                    if applied_settings.ban_failed_address {
208                        let reason = Some(error.to_string());
209                        let period_opt = error.rate_limit_ban_duration();
210                        let banned = match period_opt {
211                            // Envoy advertised a reset window: ban for exactly that period.
212                            // ban_count is set to max(ban_count,1) so diagnostics see the node
213                            // as banned, but the exponential ladder is not inflated.
214                            Some(period) => address_list.ban_for(address, period, reason),
215                            // No rate-limit hint: normal exponential health-ban ladder.
216                            None => address_list.ban_with_reason(address, reason),
217                        };
218                        if banned {
219                            if let Some(period) = period_opt {
220                                tracing::debug!(
221                                    ?address,
222                                    ban_secs = period.as_secs(),
223                                    "rate-limited (ResourceExhausted): banning {address} for {}s (from RateLimit-Reset header)",
224                                    period.as_secs()
225                                );
226                            }
227                            tracing::warn!(
228                                ?address,
229                                ?error,
230                                "ban address {address} due to error: {error}"
231                            );
232                        } else {
233                            // The address might be already removed from the list
234                            // by background process (i.e., SML update), and it's fine.
235                            tracing::debug!(
236                                ?address,
237                                ?error,
238                                "unable to ban address {address} because it's not in the list anymore"
239                            );
240                        }
241                    } else {
242                        // Banning is disabled for this request, but failover
243                        // must still move traffic away from the failing node:
244                        // drop it from the sticky rotation, ban state untouched.
245                        address_list.evict_from_rotation(address);
246                        tracing::debug!(
247                            ?error,
248                            ?address,
249                            "banning is disabled; evicted address {address} from rotation due to the error"
250                        );
251                    }
252                } else {
253                    tracing::debug!(
254                        ?error,
255                        "we should ban an address due to the error but address is absent"
256                    );
257                }
258            }
259        }
260    };
261}
262
263#[cfg(test)]
264#[allow(clippy::items_after_test_module)]
265mod tests {
266    use super::*;
267
268    fn mock_address() -> crate::Address {
269        "http://127.0.0.1:3000".parse().expect("valid address")
270    }
271
272    fn make_applied_settings(ban: bool) -> AppliedRequestSettings {
273        AppliedRequestSettings {
274            connect_timeout: None,
275            timeout: Duration::from_secs(10),
276            retries: 5,
277            ban_failed_address: ban,
278            max_decoding_message_size: None,
279            #[cfg(not(target_arch = "wasm32"))]
280            ca_certificate: None,
281        }
282    }
283
284    #[test]
285    fn test_can_retry_no_available_addresses() {
286        let err = DapiClientError::NoAvailableAddresses;
287        assert!(!err.can_retry());
288    }
289
290    #[test]
291    fn test_can_retry_no_available_addresses_to_retry() {
292        let transport_err = TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("gone"));
293        let err = DapiClientError::NoAvailableAddressesToRetry(Box::new(transport_err));
294        assert!(!err.can_retry());
295    }
296
297    #[test]
298    fn test_can_retry_transport_retryable() {
299        let transport_err =
300            TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("temporary"));
301        let err = DapiClientError::Transport(transport_err);
302        assert!(err.can_retry());
303    }
304
305    #[test]
306    fn test_can_retry_transport_non_retryable() {
307        let transport_err = TransportError::Grpc(dapi_grpc::tonic::Status::not_found("permanent"));
308        let err = DapiClientError::Transport(transport_err);
309        assert!(!err.can_retry());
310    }
311
312    #[test]
313    fn test_can_retry_address_list_error() {
314        let err =
315            DapiClientError::AddressList(AddressListError::InvalidAddressUri("bad".to_string()));
316        assert!(!err.can_retry());
317    }
318
319    /// `rate_limit_ban_duration` returns `Some` only when the `ratelimit-reset`
320    /// header is present and positive on a `ResourceExhausted` response, and
321    /// the value is clamped to `[MIN_RATE_LIMIT_BAN_SECS, MAX_RATE_LIMIT_BAN_SECS]`.
322    #[test]
323    fn test_rate_limit_ban_duration_header_parse() {
324        use dapi_grpc::tonic::metadata::MetadataValue;
325
326        // Helper: build a ResourceExhausted status with a ratelimit-reset header.
327        let make_rl_status = |header: Option<&str>| -> dapi_grpc::tonic::Status {
328            let mut status = dapi_grpc::tonic::Status::resource_exhausted("429");
329            if let Some(v) = header {
330                status
331                    .metadata_mut()
332                    .insert("ratelimit-reset", MetadataValue::try_from(v).unwrap());
333            }
334            status
335        };
336
337        // Normal header value: returned clamped.
338        let s = make_rl_status(Some("45"));
339        let dur = TransportError::Grpc(s).rate_limit_ban_duration();
340        assert_eq!(dur, Some(Duration::from_secs(45)));
341
342        // Value above MAX → clamped to MAX.
343        let s = make_rl_status(Some("9999"));
344        let dur = TransportError::Grpc(s).rate_limit_ban_duration();
345        assert_eq!(dur, Some(Duration::from_secs(MAX_RATE_LIMIT_BAN_SECS)));
346
347        // Clamp edge: exactly MIN (1) → 1 s (passes through unchanged).
348        let s = make_rl_status(Some("1"));
349        assert_eq!(
350            TransportError::Grpc(s).rate_limit_ban_duration(),
351            Some(Duration::from_secs(1))
352        );
353
354        // Clamp edge: exactly MAX (600) → 600 s (not clamped).
355        let s = make_rl_status(Some("600"));
356        assert_eq!(
357            TransportError::Grpc(s).rate_limit_ban_duration(),
358            Some(Duration::from_secs(600))
359        );
360
361        // One above MAX (601) → clamped to 600 s.
362        let s = make_rl_status(Some("601"));
363        assert_eq!(
364            TransportError::Grpc(s).rate_limit_ban_duration(),
365            Some(Duration::from_secs(600))
366        );
367
368        // Value below MIN (0) → filtered to None before clamp.
369        let s = make_rl_status(Some("0"));
370        assert!(TransportError::Grpc(s).rate_limit_ban_duration().is_none());
371
372        // Non-numeric → None.
373        let s = make_rl_status(Some("garbage"));
374        assert!(TransportError::Grpc(s).rate_limit_ban_duration().is_none());
375
376        // Header absent → None.
377        let s = make_rl_status(None);
378        assert!(TransportError::Grpc(s).rate_limit_ban_duration().is_none());
379
380        // Non-ResourceExhausted code → None regardless of header.
381        let mut unavail = dapi_grpc::tonic::Status::unavailable("down");
382        unavail
383            .metadata_mut()
384            .insert("ratelimit-reset", MetadataValue::try_from("30").unwrap());
385        assert!(TransportError::Grpc(unavail)
386            .rate_limit_ban_duration()
387            .is_none());
388    }
389
390    /// When `ResourceExhausted` carries a valid `ratelimit-reset` header,
391    /// `update_address_ban_status` calls `ban_for` (exact period, no ladder
392    /// inflation); when the header is absent it falls through to `ban_with_reason`
393    /// (normal exponential ladder).
394    #[test]
395    fn test_update_address_ban_status_rate_limit_ban_path() {
396        use dapi_grpc::tonic::metadata::MetadataValue;
397
398        let mut address_list = AddressList::new();
399        let addr = mock_address();
400        address_list.add(addr.clone());
401
402        // Build a ResourceExhausted status with ratelimit-reset: 45.
403        let mut status = dapi_grpc::tonic::Status::resource_exhausted("429");
404        status
405            .metadata_mut()
406            .insert("ratelimit-reset", MetadataValue::try_from("45").unwrap());
407
408        let result: ExecutionResult<i32, DapiClientError> = Err(ExecutionError {
409            inner: DapiClientError::Transport(TransportError::Grpc(status)),
410            retries: 0,
411            address: Some(addr.clone()),
412        });
413        let before = chrono::Utc::now();
414        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
415        let after = chrono::Utc::now();
416
417        let info = address_list.ban_info();
418        let entry = info.iter().find(|i| i.uri == addr.to_string()).unwrap();
419
420        // Node is banned for ~45 s.
421        assert!(entry.banned, "rate-limited node must be banned");
422        assert_eq!(entry.ban_count, 1, "ban_count must be 1 after ban_for");
423        let until = entry.banned_until.expect("banned_until set");
424        let lo = (until - before).num_milliseconds() as f64 / 1000.0;
425        let hi = (until - after).num_milliseconds() as f64 / 1000.0;
426        assert!(
427            lo >= 44.9 && hi <= 45.1,
428            "ban window must be ~45 s, got lo={lo} hi={hi}"
429        );
430    }
431
432    /// When `ResourceExhausted` has NO `ratelimit-reset` header,
433    /// `update_address_ban_status` must fall back to the normal `ban_with_reason`
434    /// ladder (not produce a zero-second or panic ban).
435    #[test]
436    fn test_update_address_ban_status_rate_limit_no_header_uses_ladder() {
437        let mut address_list = AddressList::new();
438        let addr = mock_address();
439        address_list.add(addr.clone());
440
441        let result: ExecutionResult<i32, DapiClientError> = Err(ExecutionError {
442            inner: DapiClientError::Transport(TransportError::Grpc(
443                dapi_grpc::tonic::Status::resource_exhausted("429"),
444            )),
445            retries: 0,
446            address: Some(addr.clone()),
447        });
448        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
449
450        // The ban ladder is invoked: first ban → ban_count = 1, window = 60 s.
451        let info = address_list.ban_info();
452        let entry = info.iter().find(|i| i.uri == addr.to_string()).unwrap();
453        assert!(
454            entry.banned,
455            "node must be banned on ResourceExhausted without header"
456        );
457        assert_eq!(
458            entry.ban_count, 1,
459            "first health-ladder ban → ban_count = 1"
460        );
461    }
462
463    #[cfg(feature = "mocks")]
464    #[test]
465    fn test_can_retry_mock_error() {
466        let err = DapiClientError::Mock(crate::mock::MockError::MockExpectationNotFound(
467            "test".to_string(),
468        ));
469        assert!(!err.can_retry());
470    }
471
472    #[test]
473    fn test_is_no_available_addresses() {
474        assert!(DapiClientError::NoAvailableAddresses.is_no_available_addresses());
475
476        let transport_err = TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("gone"));
477        assert!(
478            DapiClientError::NoAvailableAddressesToRetry(Box::new(transport_err))
479                .is_no_available_addresses()
480        );
481
482        let transport_err =
483            TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("temporary"));
484        assert!(!DapiClientError::Transport(transport_err).is_no_available_addresses());
485    }
486
487    #[test]
488    fn test_update_address_ban_status_success_unbans() {
489        let mut address_list = AddressList::new();
490        let addr = mock_address();
491        address_list.add(addr.clone());
492        address_list.ban(&addr);
493        assert!(address_list.is_banned(&addr));
494
495        let result: ExecutionResult<i32, DapiClientError> = Ok(ExecutionResponse {
496            inner: 42,
497            retries: 0,
498            address: addr.clone(),
499        });
500
501        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
502
503        assert!(!address_list.is_banned(&addr));
504    }
505
506    #[test]
507    fn test_update_address_ban_status_success_on_unbanned_is_noop() {
508        let mut address_list = AddressList::new();
509        let addr = mock_address();
510        address_list.add(addr.clone());
511
512        let result: ExecutionResult<i32, DapiClientError> = Ok(ExecutionResponse {
513            inner: 42,
514            retries: 0,
515            address: addr.clone(),
516        });
517
518        // Should not panic or change anything
519        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
520        assert!(!address_list.is_banned(&addr));
521    }
522
523    #[test]
524    fn test_update_address_ban_status_retryable_error_bans_address() {
525        let mut address_list = AddressList::new();
526        let addr = mock_address();
527        address_list.add(addr.clone());
528
529        let transport_err =
530            TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("temporary"));
531        let result: ExecutionResult<i32, DapiClientError> = Err(ExecutionError {
532            inner: DapiClientError::Transport(transport_err),
533            retries: 0,
534            address: Some(addr.clone()),
535        });
536
537        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
538        assert!(address_list.is_banned(&addr));
539
540        // The ban reason must be propagated from the error via this call path,
541        // not just the ban itself.
542        let info = address_list.ban_info();
543        assert_eq!(info.len(), 1);
544        let reason = info[0].reason.as_deref().expect("ban reason recorded");
545        assert!(
546            reason.contains("temporary"),
547            "ban reason should carry the underlying error, got: {reason}"
548        );
549    }
550
551    #[test]
552    fn test_update_address_ban_status_retryable_error_ban_disabled() {
553        let mut address_list = AddressList::new();
554        let addr = mock_address();
555        address_list.add(addr.clone());
556
557        let transport_err =
558            TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("temporary"));
559        let result: ExecutionResult<i32, DapiClientError> = Err(ExecutionError {
560            inner: DapiClientError::Transport(transport_err),
561            retries: 0,
562            address: Some(addr.clone()),
563        });
564
565        update_address_ban_status(&address_list, &result, &make_applied_settings(false));
566        // With ban disabled, the address should NOT be banned
567        assert!(!address_list.is_banned(&addr));
568    }
569
570    #[test]
571    fn test_update_address_ban_status_non_retryable_error_does_not_ban() {
572        let mut address_list = AddressList::new();
573        let addr = mock_address();
574        address_list.add(addr.clone());
575
576        let result: ExecutionResult<i32, DapiClientError> = Err(ExecutionError {
577            inner: DapiClientError::NoAvailableAddresses,
578            retries: 0,
579            address: Some(addr.clone()),
580        });
581
582        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
583        assert!(!address_list.is_banned(&addr));
584    }
585
586    #[test]
587    fn test_update_address_ban_status_retryable_error_no_address() {
588        let address_list = AddressList::new();
589
590        let transport_err =
591            TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("temporary"));
592        let result: ExecutionResult<i32, DapiClientError> = Err(ExecutionError {
593            inner: DapiClientError::Transport(transport_err),
594            retries: 0,
595            address: None,
596        });
597
598        // Should not panic when address is None
599        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
600    }
601
602    #[test]
603    fn test_update_address_ban_status_unban_removed_address() {
604        let mut address_list = AddressList::new();
605        let addr = mock_address();
606        address_list.add(addr.clone());
607        address_list.ban(&addr);
608
609        // Remove the address
610        address_list.remove(&addr);
611
612        let result: ExecutionResult<i32, DapiClientError> = Ok(ExecutionResponse {
613            inner: 42,
614            retries: 0,
615            address: addr.clone(),
616        });
617
618        // Should not panic when trying to unban a removed address
619        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
620    }
621
622    #[test]
623    fn test_update_address_ban_status_ban_removed_address() {
624        let address_list = AddressList::new();
625        let addr = mock_address();
626
627        let transport_err =
628            TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("temporary"));
629        let result: ExecutionResult<i32, DapiClientError> = Err(ExecutionError {
630            inner: DapiClientError::Transport(transport_err),
631            retries: 0,
632            address: Some(addr),
633        });
634
635        // Should not panic when trying to ban an address not in the list
636        update_address_ban_status(&address_list, &result, &make_applied_settings(true));
637    }
638
639    #[test]
640    fn test_dapi_client_new() {
641        let address_list: AddressList = "http://127.0.0.1:3000,http://127.0.0.1:3001"
642            .parse()
643            .unwrap();
644        let client = DapiClient::new(address_list, RequestSettings::default());
645        assert_eq!(client.address_list().len(), 2);
646    }
647
648    #[test]
649    fn test_dapi_client_get_live_addresses() {
650        let address_list: AddressList = "http://127.0.0.1:3000,http://127.0.0.1:3001"
651            .parse()
652            .unwrap();
653        let client = DapiClient::new(address_list, RequestSettings::default());
654        let live = client.get_live_addresses();
655        assert_eq!(live.len(), 2);
656    }
657
658    #[cfg(not(target_arch = "wasm32"))]
659    #[test]
660    fn test_dapi_client_with_ca_certificate() {
661        let address_list: AddressList = "http://127.0.0.1:3000".parse().unwrap();
662        let client = DapiClient::new(address_list, RequestSettings::default());
663        let cert = dapi_grpc::tonic::transport::Certificate::from_pem("fake-pem-data");
664        let client = client.with_ca_certificate(cert);
665        assert!(client.ca_certificate.is_some());
666    }
667
668    #[cfg(feature = "mocks")]
669    #[test]
670    fn test_dapi_client_error_mock_serialize_deserialize() {
671        use dapi_grpc::mock::Mockable;
672
673        let err = DapiClientError::NoAvailableAddresses;
674        let serialized = err.mock_serialize().expect("should serialize");
675        let deserialized =
676            DapiClientError::mock_deserialize(&serialized).expect("should deserialize");
677        assert!(matches!(
678            deserialized,
679            DapiClientError::NoAvailableAddresses
680        ));
681    }
682
683    #[cfg(feature = "mocks")]
684    #[test]
685    fn test_dapi_client_error_transport_mock_roundtrip() {
686        use dapi_grpc::mock::Mockable;
687
688        let transport_err = TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("test"));
689        let err = DapiClientError::Transport(transport_err);
690        let serialized = err.mock_serialize().expect("should serialize");
691        let deserialized =
692            DapiClientError::mock_deserialize(&serialized).expect("should deserialize");
693        assert!(matches!(deserialized, DapiClientError::Transport(_)));
694    }
695
696    #[test]
697    fn test_dapi_client_error_display() {
698        let err = DapiClientError::NoAvailableAddresses;
699        let display = format!("{}", err);
700        assert!(display.contains("no available addresses"));
701
702        let transport_err = TransportError::Grpc(dapi_grpc::tonic::Status::unavailable("gone"));
703        let err = DapiClientError::NoAvailableAddressesToRetry(Box::new(transport_err));
704        let display = format!("{}", err);
705        assert!(display.contains("no available addresses to retry"));
706
707        let err =
708            DapiClientError::AddressList(AddressListError::InvalidAddressUri("bad".to_string()));
709        let display = format!("{}", err);
710        assert!(display.contains("address list error"));
711    }
712}
713
714#[async_trait]
715impl DapiRequestExecutor for DapiClient {
716    /// Execute the [DapiRequest](crate::DapiRequest).
717    async fn execute<R>(
718        &self,
719        request: R,
720        settings: RequestSettings,
721    ) -> ExecutionResult<R::Response, DapiClientError>
722    where
723        R: TransportRequest + Mockable,
724        R::Response: Mockable,
725        TransportError: Mockable,
726    {
727        // Join settings of different sources to get final version of the settings for this execution:
728        let applied_settings = self
729            .settings
730            .override_by(R::SETTINGS_OVERRIDES)
731            .override_by(settings)
732            .finalize();
733        #[cfg(not(target_arch = "wasm32"))]
734        let applied_settings = applied_settings.with_ca_certificate(self.ca_certificate.clone());
735
736        // Save dump dir for later use
737        #[cfg(feature = "dump")]
738        let dump_dir = self.dump_dir.clone();
739        #[cfg(feature = "dump")]
740        let dump_request = request.clone();
741
742        let max_retries = applied_settings.retries;
743        let retry_delay = Duration::from_millis(10);
744
745        let mut retries: usize = 0;
746        // Track the last transport error for when all addresses get exhausted
747        let mut last_transport_error: Option<TransportError> = None;
748
749        let result: ExecutionResult<R::Response, DapiClientError> = async {
750            loop {
751                // Try to get an address to initialize transport on:
752                let Some(address) = self.address_list.get_live_address() else {
753                    // No available addresses - wrap with last meaningful error if we have one
754                    let error = if let Some(transport_error) = last_transport_error.take() {
755                        tracing::debug!(
756                            "no addresses available, returning last transport error"
757                        );
758                        DapiClientError::NoAvailableAddressesToRetry(Box::new(
759                            transport_error,
760                        ))
761                    } else {
762                        DapiClientError::NoAvailableAddresses
763                    };
764
765                    return Err(ExecutionError {
766                        inner: error,
767                        retries,
768                        address: None,
769                    });
770                };
771
772                // Rec 3 — explicit trace event so the resolved DAPI endpoint
773                // appears in flat plain-text log output (not just the span context).
774                tracing::trace!(
775                    target: "dapi_client::dispatch",
776                    ?address,
777                    method = request.method_name(),
778                    request_type = request.request_name(),
779                    "dispatching request to DAPI endpoint"
780                );
781                tracing::trace!(
782                    ?request,
783                    "calling {} with {} request",
784                    request.method_name(),
785                    request.request_name(),
786                );
787
788                let transport_request = request.clone();
789                let response_name = request.response_name();
790
791                // Try to create transport client
792                let transport_client_result = R::Client::with_uri_and_settings(
793                    address.uri().clone(),
794                    &applied_settings,
795                    &self.pool,
796                );
797
798                let mut transport_client = match transport_client_result {
799                    Ok(client) => client,
800                    Err(transport_error) => {
801                        let can_retry_error = transport_error.can_retry();
802
803                        // Clone error before moving it
804                        let cloned_error = transport_error.clone();
805
806                        let execution_error = ExecutionError {
807                            inner: DapiClientError::Transport(transport_error),
808                            retries,
809                            address: Some(address.clone()),
810                        };
811
812                        update_address_ban_status::<R::Response, DapiClientError>(
813                            &self.address_list,
814                            &Err(execution_error.clone()),
815                            &applied_settings,
816                        );
817
818                        if can_retry_error && retries < max_retries {
819                            // Store last transport error
820                            last_transport_error = Some(cloned_error);
821
822                            retries += 1;
823                            tracing::warn!(
824                                error = ?execution_error,
825                                "retrying error with sleeping {} secs",
826                                retry_delay.as_secs_f32()
827                            );
828                            transport::sleep(retry_delay).await;
829                            continue;
830                        }
831
832                        return Err(execution_error);
833                    }
834                };
835
836                // Execute the transport request
837                let result = transport_request
838                    .execute_transport(&mut transport_client, &applied_settings)
839                    .instrument(tracing::trace_span!(
840                        "execute_request",
841                        ?address,
842                        settings = ?applied_settings,
843                        method = request.method_name(),
844                    ))
845                    .await;
846
847                let execution_result = match result {
848                    Ok(response) => {
849                        tracing::trace!(response = ?response, "received {} response", response_name);
850                        Ok(ExecutionResponse {
851                            inner: response,
852                            retries,
853                            address: address.clone(),
854                        })
855                    }
856                    Err(transport_error) => {
857                        tracing::debug!(error = ?transport_error, "received error: {transport_error}");
858                        Err(ExecutionError {
859                            inner: DapiClientError::Transport(transport_error),
860                            retries,
861                            address: Some(address.clone()),
862                        })
863                    }
864                };
865
866                update_address_ban_status::<R::Response, DapiClientError>(
867                    &self.address_list,
868                    &execution_result,
869                    &applied_settings,
870                );
871
872                match execution_result {
873                    Ok(response) => return Ok(response),
874                    Err(error) => {
875                        if error.can_retry() && retries < max_retries {
876                            // Store last transport error
877                            if let DapiClientError::Transport(ref te) = error.inner {
878                                last_transport_error = Some(te.clone());
879                            }
880
881                            retries += 1;
882                            tracing::warn!(
883                                ?error,
884                                "retrying error with sleeping {} secs",
885                                retry_delay.as_secs_f32()
886                            );
887                            transport::sleep(retry_delay).await;
888                            continue;
889                        }
890
891                        return Err(error);
892                    }
893                }
894            }
895        }
896        .instrument(tracing::info_span!("request routine"))
897        .await;
898
899        if let Err(error) = &result {
900            if !error.can_retry() {
901                tracing::error!(?error, "request failed");
902            }
903        }
904
905        // Dump request and response to disk if dump_dir is set:
906        #[cfg(feature = "dump")]
907        Self::dump_request_response(&dump_request, &result, dump_dir);
908
909        result
910    }
911}