From a04644130a6528f065774509d95756ec95e982fa Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 10:47:11 +0200 Subject: [PATCH 01/14] feat(sdk): detect stale proofs --- packages/rs-drive-proof-verifier/src/error.rs | 7 ++ packages/rs-sdk/src/sdk.rs | 65 +++++++++++++++++-- 2 files changed, 67 insertions(+), 5 deletions(-) diff --git a/packages/rs-drive-proof-verifier/src/error.rs b/packages/rs-drive-proof-verifier/src/error.rs index 3203eb73174..069e4e9f28b 100644 --- a/packages/rs-drive-proof-verifier/src/error.rs +++ b/packages/rs-drive-proof-verifier/src/error.rs @@ -82,6 +82,13 @@ pub enum Error { /// Context provider error #[error("context provider error: {0}")] ContextProviderError(#[from] ContextProviderError), + + /// Proof is stale; try another server + #[error("proof is stale; try another server")] + StaleProof { + expected_height: u64, + actual_height: u64, + }, } /// Errors returned by the context provider diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index f7f938703d7..2fc5d617eff 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -36,7 +36,8 @@ use std::fmt::Debug; use std::num::NonZeroUsize; #[cfg(feature = "mocks")] use std::path::{Path, PathBuf}; -use std::sync::Arc; +use std::sync::atomic::Ordering; +use std::sync::{atomic, Arc}; use std::time::{SystemTime, UNIX_EPOCH}; #[cfg(feature = "mocks")] use tokio::sync::{Mutex, MutexGuard}; @@ -50,6 +51,9 @@ pub const DEFAULT_QUORUM_PUBLIC_KEYS_CACHE_SIZE: usize = 100; /// The default identity nonce stale time in seconds pub const DEFAULT_IDENTITY_NONCE_STALE_TIME_S: u64 = 1200; //20 mins +/// How many blocks difference is allowed between the last proof height and the current proof height. +const PROOF_HEIGHT_TOLERANCE: u64 = 1; + /// a type to represent staleness in seconds pub type StalenessInSeconds = u64; @@ -100,6 +104,11 @@ pub struct Sdk { /// Note that setting this to None can panic. context_provider: ArcSwapOption>, + /// Last proof height; used to determine if the proof is stale. + /// + /// This is clone-able and can be shared between threads. + previous_proof_height: Arc, + /// Cancellation token; once cancelled, all pending requests should be aborted. pub(crate) cancel_token: CancellationToken, @@ -115,6 +124,7 @@ impl Clone for Sdk { internal_cache: Arc::clone(&self.internal_cache), context_provider: ArcSwapOption::new(self.context_provider.load_full()), cancel_token: self.cancel_token.clone(), + previous_proof_height: self.previous_proof_height.clone(), #[cfg(feature = "mocks")] dump_dir: self.dump_dir.clone(), } @@ -221,7 +231,7 @@ impl Sdk { .context_provider() .ok_or(drive_proof_verifier::Error::ContextProviderNotSet)?; - match self.inner { + let (object, mtd) = match self.inner { SdkInstance::Dapi { .. } => O::maybe_from_proof_with_metadata( request, response, @@ -237,9 +247,53 @@ impl Sdk { .parse_proof_with_metadata(request, response) .map(|(a, b, _)| (a, b)) } - } + }?; + + self.verify_metadata(&mtd)?; + + Ok((object, mtd)) } + /// Verify metadata contained in the response. + /// + /// This method is used to verify metadata contained in the response. + /// It also updates the last proof height. + fn verify_metadata( + &self, + metadata: &ResponseMetadata, + ) -> Result<(), drive_proof_verifier::Error> { + let mut prev = self.previous_proof_height.load(Ordering::Relaxed); + let received = metadata.height; + + // Same height, no need to update. + if received == prev { + return Ok(()); + } + + // If received proof height is behind previous proof height by more than PROOF_HEIGHT_TOLERANCE, the proof is stale. + if received < prev - PROOF_HEIGHT_TOLERANCE { + return Err(drive_proof_verifier::Error::StaleProof { + expected_height: prev, + actual_height: metadata.height, + }); + } + + // New proof is ahead of the previous proof, so we update the previous proof height. + while let Err(stored) = self.previous_proof_height.compare_exchange( + prev, + received, + Ordering::SeqCst, + Ordering::Relaxed, + ) { + // The value was changed to a higher value by another thread, so we need to retry. + if stored >= metadata.height { + break; + } + prev = stored; + } + + Ok(()) + } /// Retrieve object `O` from proof contained in `request` (of type `R`) and `response`. /// /// This method is used to retrieve objects from proofs returned by Dash Platform. @@ -792,9 +846,10 @@ impl SdkBuilder { proofs:self.proofs, context_provider: ArcSwapOption::new( self.context_provider.map(Arc::new)), cancel_token: self.cancel_token, + internal_cache: Default::default(), + previous_proof_height: Arc::new(atomic::AtomicU64::new(0)), #[cfg(feature = "mocks")] dump_dir: self.dump_dir, - internal_cache: Default::default(), }; // if context provider is not set correctly (is None), it means we need to fallback to core wallet if sdk.context_provider.load().is_none() { @@ -850,13 +905,13 @@ impl SdkBuilder { mock:mock_sdk.clone(), dapi, version:self.version, - }, dump_dir: self.dump_dir.clone(), proofs:self.proofs, internal_cache: Default::default(), context_provider:ArcSwapAny::new( Some(Arc::new(context_provider))), cancel_token: self.cancel_token, + previous_proof_height: Arc::new(atomic::AtomicU64::new(0)), }; let mut guard = mock_sdk.try_lock().expect("mock sdk is in use by another thread and connot be reconfigured"); guard.set_sdk(sdk.clone()); From 4bbf930d421b4fb020fcca7614860b20edc0b80a Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Thu, 15 Aug 2024 11:17:02 +0200 Subject: [PATCH 02/14] refactor(dapi-client): replace CanRetry.is_node_failure with !CanRetry.can_retry() --- packages/rs-dapi-client/src/dapi_client.rs | 14 +++++++------- packages/rs-dapi-client/src/lib.rs | 10 +++++++++- packages/rs-dapi-client/src/transport/grpc.rs | 4 ++-- 3 files changed, 18 insertions(+), 10 deletions(-) diff --git a/packages/rs-dapi-client/src/dapi_client.rs b/packages/rs-dapi-client/src/dapi_client.rs index 372b28bc3fa..611aa96316e 100644 --- a/packages/rs-dapi-client/src/dapi_client.rs +++ b/packages/rs-dapi-client/src/dapi_client.rs @@ -39,12 +39,12 @@ pub enum DapiClientError { } impl CanRetry for DapiClientError { - fn is_node_failure(&self) -> bool { + fn can_retry(&self) -> bool { use DapiClientError::*; match self { - NoAvailableAddresses => false, - Transport(transport_error, _) => transport_error.is_node_failure(), - AddressList(_) => false, + NoAvailableAddresses => true, + Transport(transport_error, _) => transport_error.can_retry(), + AddressList(_) => true, #[cfg(feature = "mocks")] Mock(_) => false, } @@ -233,7 +233,7 @@ impl DapiRequestExecutor for DapiClient { tracing::trace!(?response, "received {} response", response_name); } Err(error) => { - if error.is_node_failure() { + if !error.can_retry() { if applied_settings.ban_failed_address { let mut address_list = self .address_list @@ -264,12 +264,12 @@ impl DapiRequestExecutor for DapiClient { duration.as_secs_f32() ) }) - .when(|e| e.is_node_failure()) + .when(|e| !e.can_retry()) .instrument(tracing::info_span!("request routine")) .await; if let Err(error) = &result { - if error.is_node_failure() { + if !error.can_retry() { tracing::error!(?error, "request failed"); } } diff --git a/packages/rs-dapi-client/src/lib.rs b/packages/rs-dapi-client/src/lib.rs index 976537097eb..7a352019594 100644 --- a/packages/rs-dapi-client/src/lib.rs +++ b/packages/rs-dapi-client/src/lib.rs @@ -74,6 +74,14 @@ impl DapiRequest for T { /// Allows to flag the transport error variant how tolerant we are of it and whether we can /// try to do a request again. pub trait CanRetry { + /// Returns true if the operation can be retried safely, false means it's unspecified + fn can_retry(&self) -> bool; + /// Get boolean flag that indicates if the error is retryable. - fn is_node_failure(&self) -> bool; + /// + /// Depreacted in favor of [CanRetry::can_retry]. + #[deprecated = "Use !can_retry() instead"] + fn is_node_failure(&self) -> bool { + !self.can_retry() + } } diff --git a/packages/rs-dapi-client/src/transport/grpc.rs b/packages/rs-dapi-client/src/transport/grpc.rs index d5180099d0a..dd141d3db7a 100644 --- a/packages/rs-dapi-client/src/transport/grpc.rs +++ b/packages/rs-dapi-client/src/transport/grpc.rs @@ -117,11 +117,11 @@ impl TransportClient for CoreGrpcClient { } impl CanRetry for dapi_grpc::tonic::Status { - fn is_node_failure(&self) -> bool { + fn can_retry(&self) -> bool { let code = self.code(); use dapi_grpc::tonic::Code::*; - matches!( + !matches!( code, Ok | DataLoss | Cancelled From cbabee83898d8351f1829f788c2660911ae5c71b Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Thu, 15 Aug 2024 11:17:17 +0200 Subject: [PATCH 03/14] feat(sdk): Error implements CanRetry --- packages/rs-sdk/src/error.rs | 12 +++++++++++- 1 file changed, 11 insertions(+), 1 deletion(-) diff --git a/packages/rs-sdk/src/error.rs b/packages/rs-sdk/src/error.rs index ce8b3f309ad..500d105771b 100644 --- a/packages/rs-sdk/src/error.rs +++ b/packages/rs-sdk/src/error.rs @@ -5,7 +5,7 @@ use std::time::Duration; use dapi_grpc::mock::Mockable; use dpp::version::PlatformVersionError; use dpp::ProtocolError; -use rs_dapi_client::DapiClientError; +use rs_dapi_client::{CanRetry, DapiClientError}; pub use drive_proof_verifier::error::ContextProviderError; @@ -80,3 +80,13 @@ impl From for Error { Self::Protocol(value.into()) } } +impl CanRetry for Error { + /// Returns true if the operation can be retried, false means it's unspecified + /// False means + fn can_retry(&self) -> bool { + matches!( + self, + Error::CoreLockedHeightNotYetAvailable(_, _) | Error::QuorumNotFound { .. } + ) + } +} From ad7c235279edb97931dee9169b3d035203be163d Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 12:27:19 +0200 Subject: [PATCH 04/14] refactor: CanRetry minor refactoring --- packages/rs-dapi-client/src/dapi_client.rs | 14 ++++++------- packages/rs-dapi-client/src/lib.rs | 12 ++++++----- packages/rs-dapi-client/src/transport/grpc.rs | 13 +++++++++--- packages/rs-sdk/src/error.rs | 20 +++++++++++++------ 4 files changed, 38 insertions(+), 21 deletions(-) diff --git a/packages/rs-dapi-client/src/dapi_client.rs b/packages/rs-dapi-client/src/dapi_client.rs index 611aa96316e..4a11059d7de 100644 --- a/packages/rs-dapi-client/src/dapi_client.rs +++ b/packages/rs-dapi-client/src/dapi_client.rs @@ -39,14 +39,14 @@ pub enum DapiClientError { } impl CanRetry for DapiClientError { - fn can_retry(&self) -> bool { + fn can_retry(&self) -> Option { use DapiClientError::*; match self { - NoAvailableAddresses => true, + NoAvailableAddresses => Some(true), Transport(transport_error, _) => transport_error.can_retry(), - AddressList(_) => true, + AddressList(_) => Some(true), #[cfg(feature = "mocks")] - Mock(_) => false, + Mock(_) => None, } } } @@ -233,7 +233,7 @@ impl DapiRequestExecutor for DapiClient { tracing::trace!(?response, "received {} response", response_name); } Err(error) => { - if !error.can_retry() { + if !error.can_retry().unwrap_or(false) { if applied_settings.ban_failed_address { let mut address_list = self .address_list @@ -264,12 +264,12 @@ impl DapiRequestExecutor for DapiClient { duration.as_secs_f32() ) }) - .when(|e| !e.can_retry()) + .when(|e| !e.can_retry().unwrap_or(false)) .instrument(tracing::info_span!("request routine")) .await; if let Err(error) = &result { - if !error.can_retry() { + if !error.can_retry().unwrap_or(false) { tracing::error!(?error, "request failed"); } } diff --git a/packages/rs-dapi-client/src/lib.rs b/packages/rs-dapi-client/src/lib.rs index 7a352019594..8b127f52279 100644 --- a/packages/rs-dapi-client/src/lib.rs +++ b/packages/rs-dapi-client/src/lib.rs @@ -71,17 +71,19 @@ impl DapiRequest for T { } } -/// Allows to flag the transport error variant how tolerant we are of it and whether we can -/// try to do a request again. +/// Returns true if the operation can be retried. +/// `None` means unspecified - in this case, you should either inspect the error +/// in more details, or assume `false`. pub trait CanRetry { - /// Returns true if the operation can be retried safely, false means it's unspecified - fn can_retry(&self) -> bool; + /// Returns true if the operation can be retried safely. + /// None value means it's unspecified + fn can_retry(&self) -> Option; /// Get boolean flag that indicates if the error is retryable. /// /// Depreacted in favor of [CanRetry::can_retry]. #[deprecated = "Use !can_retry() instead"] fn is_node_failure(&self) -> bool { - !self.can_retry() + !self.can_retry().unwrap_or(false) } } diff --git a/packages/rs-dapi-client/src/transport/grpc.rs b/packages/rs-dapi-client/src/transport/grpc.rs index dd141d3db7a..735df4f0b66 100644 --- a/packages/rs-dapi-client/src/transport/grpc.rs +++ b/packages/rs-dapi-client/src/transport/grpc.rs @@ -117,11 +117,12 @@ impl TransportClient for CoreGrpcClient { } impl CanRetry for dapi_grpc::tonic::Status { - fn can_retry(&self) -> bool { + fn can_retry(&self) -> Option { let code = self.code(); use dapi_grpc::tonic::Code::*; - !matches!( + + let retry = !matches!( code, Ok | DataLoss | Cancelled @@ -131,7 +132,13 @@ impl CanRetry for dapi_grpc::tonic::Status { | Aborted | Internal | Unavailable - ) + ); + + if retry { + Some(true) + } else { + None + } } } diff --git a/packages/rs-sdk/src/error.rs b/packages/rs-sdk/src/error.rs index 500d105771b..a084d195131 100644 --- a/packages/rs-sdk/src/error.rs +++ b/packages/rs-sdk/src/error.rs @@ -80,13 +80,21 @@ impl From for Error { Self::Protocol(value.into()) } } + impl CanRetry for Error { - /// Returns true if the operation can be retried, false means it's unspecified - /// False means - fn can_retry(&self) -> bool { - matches!( + fn can_retry(&self) -> Option { + let retry = matches!( self, - Error::CoreLockedHeightNotYetAvailable(_, _) | Error::QuorumNotFound { .. } - ) + Error::Proof(drive_proof_verifier::Error::StaleProof { .. }) + | Error::DapiClientError(_) + | Error::CoreClientError(_) + | Error::TimeoutReached(_, _) + ); + + if retry { + Some(true) + } else { + None + } } } From afd750df6edc47e866c6d692310a8e6a2ce395cd Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 13:44:22 +0200 Subject: [PATCH 05/14] chore: proof tolerance configurable --- packages/rs-drive-proof-verifier/src/error.rs | 26 ++- packages/rs-sdk/src/error.rs | 2 +- packages/rs-sdk/src/sdk.rs | 211 ++++++++++++++---- 3 files changed, 198 insertions(+), 41 deletions(-) diff --git a/packages/rs-drive-proof-verifier/src/error.rs b/packages/rs-drive-proof-verifier/src/error.rs index 069e4e9f28b..f14636aa51e 100644 --- a/packages/rs-drive-proof-verifier/src/error.rs +++ b/packages/rs-drive-proof-verifier/src/error.rs @@ -85,9 +85,33 @@ pub enum Error { /// Proof is stale; try another server #[error("proof is stale; try another server")] - StaleProof { + StaleProof(#[from] StaleProofError), +} + +/// Received proof is stale; try another server +#[derive(Debug, thiserror::Error)] +pub enum StaleProofError { + /// Stale proof height + #[error("stale proof height: expected height {expected_height}, received {actual_height}, tolerance {tolerance}, try another server")] + StaleProofHeight { + /// Expected height - last block height seen by the Sdk expected_height: u64, + /// Actual height - block height received from the server in the proof actual_height: u64, + /// Tolerance - how many blocks can be behind the expected height + tolerance: u64, + }, + /// Proof time is stale + #[error( + "received outdated proof time: expected {expected_ms} ms, received {actual_ms} ms, tolerance {tolerance_ms} ms, try another server" + )] + Time { + /// Expected time in milliseconds - is local time when the proof was received + expected_ms: u64, + /// Actual time in milliseconds - time received from the server in the proof + actual_ms: u64, + /// Tolerance in milliseconds + tolerance_ms: u64, }, } diff --git a/packages/rs-sdk/src/error.rs b/packages/rs-sdk/src/error.rs index a084d195131..4d468557663 100644 --- a/packages/rs-sdk/src/error.rs +++ b/packages/rs-sdk/src/error.rs @@ -85,7 +85,7 @@ impl CanRetry for Error { fn can_retry(&self) -> Option { let retry = matches!( self, - Error::Proof(drive_proof_verifier::Error::StaleProof { .. }) + Error::Proof(drive_proof_verifier::Error::StaleProof(..)) | Error::DapiClientError(_) | Error::CoreClientError(_) | Error::TimeoutReached(_, _) diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index 2fc5d617eff..c37383b5af3 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -51,9 +51,6 @@ pub const DEFAULT_QUORUM_PUBLIC_KEYS_CACHE_SIZE: usize = 100; /// The default identity nonce stale time in seconds pub const DEFAULT_IDENTITY_NONCE_STALE_TIME_S: u64 = 1200; //20 mins -/// How many blocks difference is allowed between the last proof height and the current proof height. -const PROOF_HEIGHT_TOLERANCE: u64 = 1; - /// a type to represent staleness in seconds pub type StalenessInSeconds = u64; @@ -109,6 +106,20 @@ pub struct Sdk { /// This is clone-able and can be shared between threads. previous_proof_height: Arc, + /// How many blocks difference is allowed between the last proof height and the current proof height. + /// If current proof height is behind previous proof height by more than this value, the proof is stale. + /// If None, proof height is not checked. + /// + /// This is set to `1` by default. + proof_height_tolerance: Option, + + /// How many milliseconds difference is allowed between the proof time and current local time. + /// If current proof time differs from local time by more than this value, the proof is stale. + /// If None, proof time is not checked. + /// + /// This is set to `None` by fefault. + proof_time_tolerance_ms: Option, + /// Cancellation token; once cancelled, all pending requests should be aborted. pub(crate) cancel_token: CancellationToken, @@ -125,6 +136,8 @@ impl Clone for Sdk { context_provider: ArcSwapOption::new(self.context_provider.load_full()), cancel_token: self.cancel_token.clone(), previous_proof_height: self.previous_proof_height.clone(), + proof_height_tolerance: self.proof_height_tolerance, + proof_time_tolerance_ms: self.proof_time_tolerance_ms, #[cfg(feature = "mocks")] dump_dir: self.dump_dir.clone(), } @@ -249,51 +262,29 @@ impl Sdk { } }?; - self.verify_metadata(&mtd)?; - + self.verify_proof_metadata(&mtd)?; Ok((object, mtd)) } - /// Verify metadata contained in the response. - /// - /// This method is used to verify metadata contained in the response. - /// It also updates the last proof height. - fn verify_metadata( + /// Verify proof metadata against the current state of the SDK. + fn verify_proof_metadata( &self, metadata: &ResponseMetadata, ) -> Result<(), drive_proof_verifier::Error> { - let mut prev = self.previous_proof_height.load(Ordering::Relaxed); - let received = metadata.height; - - // Same height, no need to update. - if received == prev { - return Ok(()); - } - - // If received proof height is behind previous proof height by more than PROOF_HEIGHT_TOLERANCE, the proof is stale. - if received < prev - PROOF_HEIGHT_TOLERANCE { - return Err(drive_proof_verifier::Error::StaleProof { - expected_height: prev, - actual_height: metadata.height, - }); - } - - // New proof is ahead of the previous proof, so we update the previous proof height. - while let Err(stored) = self.previous_proof_height.compare_exchange( - prev, - received, - Ordering::SeqCst, - Ordering::Relaxed, - ) { - // The value was changed to a higher value by another thread, so we need to retry. - if stored >= metadata.height { - break; - } - prev = stored; - } + if let Some(height_tolerance) = self.proof_height_tolerance { + verify_proof_height( + metadata, + height_tolerance, + Arc::clone(&(self.previous_proof_height)), + )?; + }; + if let Some(time_tolerance) = self.proof_time_tolerance_ms { + verify_proof_time(metadata, time_tolerance)?; + }; Ok(()) } + /// Retrieve object `O` from proof contained in `request` (of type `R`) and `response`. /// /// This method is used to retrieve objects from proofs returned by Dash Platform. @@ -598,6 +589,103 @@ impl Sdk { } } +/// If current proof time differs from local time by more than `tolerance`, the proof is considered stale. +/// This is used to verify the freshness of the proof. +fn verify_proof_time( + metadata: &ResponseMetadata, + tolerance: u64, +) -> Result<(), drive_proof_verifier::Error> { + let now = chrono::Utc::now().timestamp_millis() as u64; + + let proof_time = metadata.time_ms; + + // proof_time - tolerance <= now <= proof_time + tolerance + if now.abs_diff(proof_time) > tolerance { + tracing::warn!( + expected_time = now, + actual_time = proof_time, + tolerance, + "received proof with stale time; you should retry with another server" + ); + return Err(drive_proof_verifier::error::StaleProofError::Time { + expected_ms: now, + actual_ms: proof_time, + tolerance_ms: tolerance, + } + .into()); + } + + tracing::trace!( + expected_time = now, + actual_time = proof_time, + tolerance, + "received proof with valid time" + ); + Ok(()) +} + +/// If current proof height is behind previous proof height by more than `tolerance`, the proof is considered stale. +/// This is used to verify the freshness of the proof. +fn verify_proof_height( + metadata: &ResponseMetadata, + tolerance: u64, + previous_proof_height: Arc, +) -> Result<(), drive_proof_verifier::Error> { + let mut prev = previous_proof_height.load(Ordering::Relaxed); + let received = metadata.height; + + // Same height, no need to update. + if received == prev { + tracing::trace!( + expected_height = prev, + actual_height = received, + tolerance, + "received proof with the same height as previous" + ); + return Ok(()); + } + + // If received proof height is behind previous proof height by more than PROOF_HEIGHT_TOLERANCE, the proof is stale. + // If prev is less than tolerance, then Sdk just started, so we just trust the proof, assuming we connected to a + // trusted node. + // FIXME: in future, we need to implement t + if prev > tolerance && received < prev - tolerance { + tracing::warn!( + expected_height = prev, + actual_height = received, + tolerance, + "received proof with stale height; you should retry with another server" + ); + return Err( + drive_proof_verifier::error::StaleProofError::StaleProofHeight { + expected_height: prev, + actual_height: received, + tolerance, + } + .into(), + ); + } + + // New proof is ahead of the previous proof, so we update the previous proof height. + tracing::trace!( + expected_height = prev, + actual_height = received, + tolerance, + "received proof with new height" + ); + while let Err(stored) = + previous_proof_height.compare_exchange(prev, received, Ordering::SeqCst, Ordering::Relaxed) + { + // The value was changed to a higher value by another thread, so we need to retry. + if stored >= metadata.height { + break; + } + prev = stored; + } + + Ok(()) +} + #[async_trait::async_trait] impl DapiRequestExecutor for Sdk { async fn execute( @@ -659,6 +747,20 @@ pub struct SdkBuilder { /// Context provider used by the SDK. context_provider: Option>, + /// How many blocks difference is allowed between the last proof height and the current proof height. + /// If current proof height is behind previous proof height by more than this value, the proof is stale. + /// If None, proof height is not checked. + /// + /// This is set to `1` by default. + proof_height_tolerance: Option, + + /// How many milliseconds difference is allowed between the proof time and current local time. + /// If current proof time differs from local time by more than this value, the proof is stale. + /// If None, proof time is not checked. + /// + /// This is set to `None` by fefault. + proof_time_tolerance_ms: Option, + /// directory where dump files will be stored #[cfg(feature = "mocks")] dump_dir: Option, @@ -680,6 +782,8 @@ impl Default for SdkBuilder { core_user: "".to_string(), proofs: true, + proof_height_tolerance: Some(1), + proof_time_tolerance_ms: None, #[cfg(feature = "mocks")] data_contract_cache_size: NonZeroUsize::new(DEFAULT_CONTRACT_CACHE_SIZE) @@ -804,6 +908,30 @@ impl SdkBuilder { self } + /// Change number of blocks difference allowed between the last proof height and the current proof height. + /// + /// If current proof height is behind previous proof height by more than this value, the proof is stale. + /// If None, proof height is not checked. + /// + /// This is set to `1` by default. + pub fn with_proof_height_tolerance(mut self, tolerance: Option) -> Self { + self.proof_height_tolerance = tolerance; + self + } + + /// How many milliseconds difference is allowed between the proof time and current local time. + /// If current proof time differs from local time by more than this value, the proof is stale. + /// If None, proof time is not checked. + /// + /// Note that enabling this check can cause issues if the local time is not synchronized with the network time, + /// when the network is stalled or time between blocks increases significantly. + /// + /// This is set to `None` by fefault. + pub fn with_proof_time_tolerance(mut self, tolerance_ms: Option) -> Self { + self.proof_time_tolerance_ms = tolerance_ms; + self + } + /// Configure directory where dumps of all requests and responses will be saved. /// Useful for debugging. /// @@ -848,6 +976,8 @@ impl SdkBuilder { cancel_token: self.cancel_token, internal_cache: Default::default(), previous_proof_height: Arc::new(atomic::AtomicU64::new(0)), + proof_height_tolerance: self.proof_height_tolerance, + proof_time_tolerance_ms: self.proof_time_tolerance_ms, #[cfg(feature = "mocks")] dump_dir: self.dump_dir, }; @@ -912,6 +1042,8 @@ impl SdkBuilder { context_provider:ArcSwapAny::new( Some(Arc::new(context_provider))), cancel_token: self.cancel_token, previous_proof_height: Arc::new(atomic::AtomicU64::new(0)), + proof_height_tolerance: self.proof_height_tolerance, + proof_time_tolerance_ms: self.proof_time_tolerance_ms, }; let mut guard = mock_sdk.try_lock().expect("mock sdk is in use by another thread and connot be reconfigured"); guard.set_sdk(sdk.clone()); @@ -957,3 +1089,4 @@ pub fn prettify_proof(proof: &Proof) -> String { proof.quorum_type, ) } + From d390820b051271669c8857c119c1c34845b3a166 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 13:55:19 +0200 Subject: [PATCH 06/14] test: sdk --- packages/rs-sdk/Cargo.toml | 2 +- packages/rs-sdk/src/sdk.rs | 155 ++++++++++++++++++++++++++++++++++--- 2 files changed, 145 insertions(+), 12 deletions(-) diff --git a/packages/rs-sdk/Cargo.toml b/packages/rs-sdk/Cargo.toml index ebf783f6523..98f5dd4c023 100644 --- a/packages/rs-sdk/Cargo.toml +++ b/packages/rs-sdk/Cargo.toml @@ -6,6 +6,7 @@ edition = "2021" [dependencies] arc-swap = { version = "1.7.1" } +chrono = { version = "0.4.38" } dpp = { path = "../rs-dpp", default-features = false, features = [ "dash-sdk-features", ] } @@ -52,7 +53,6 @@ data-contracts = { path = "../data-contracts" } tokio-test = { version = "0.4.4" } clap = { version = "4.5.4", features = ["derive"] } sanitize-filename = { version = "0.5.0" } -chrono = { version = "0.4.38" } test-case = { version = "3.3.1" } [features] diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index c37383b5af3..e8dd8e70b20 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -279,7 +279,8 @@ impl Sdk { )?; }; if let Some(time_tolerance) = self.proof_time_tolerance_ms { - verify_proof_time(metadata, time_tolerance)?; + let now = chrono::Utc::now().timestamp_millis() as u64; + verify_proof_time(metadata, now, time_tolerance)?; }; Ok(()) @@ -591,34 +592,39 @@ impl Sdk { /// If current proof time differs from local time by more than `tolerance`, the proof is considered stale. /// This is used to verify the freshness of the proof. +/// +/// ## Parameters +/// +/// - `metadata`: Metadata of the proof +/// - `now_ms`: Current local time in milliseconds +/// - `tolerance_ms`: Tolerance in milliseconds fn verify_proof_time( metadata: &ResponseMetadata, - tolerance: u64, + now_ms: u64, + tolerance_ms: u64, ) -> Result<(), drive_proof_verifier::Error> { - let now = chrono::Utc::now().timestamp_millis() as u64; - let proof_time = metadata.time_ms; // proof_time - tolerance <= now <= proof_time + tolerance - if now.abs_diff(proof_time) > tolerance { + if now_ms.abs_diff(proof_time) > tolerance_ms { tracing::warn!( - expected_time = now, + expected_time = now_ms, actual_time = proof_time, - tolerance, + tolerance_ms, "received proof with stale time; you should retry with another server" ); return Err(drive_proof_verifier::error::StaleProofError::Time { - expected_ms: now, + expected_ms: now_ms, actual_ms: proof_time, - tolerance_ms: tolerance, + tolerance_ms, } .into()); } tracing::trace!( - expected_time = now, + expected_time = now_ms, actual_time = proof_time, - tolerance, + tolerance_ms, "received proof with valid time" ); Ok(()) @@ -1090,3 +1096,130 @@ pub fn prettify_proof(proof: &Proof) -> String { ) } +#[cfg(test)] +mod test { + use std::sync::Arc; + + use dapi_grpc::platform::v0::ResponseMetadata; + use test_case::test_matrix; + + use crate::SdkBuilder; + + #[test_matrix(97..102, 100, 2, false; "valid height")] + #[test_case(103, 100, 2, true; "invalid height")] + fn test_verify_proof_height(expected: u64, received: u64, tolerance: u64, expect_err: bool) { + let metadata = ResponseMetadata { + height: received, + ..Default::default() + }; + + let previous_proof_height = + std::sync::Arc::new(std::sync::atomic::AtomicU64::new(expected)); + + let result = + super::verify_proof_height(&metadata, tolerance, Arc::clone(&previous_proof_height)); + + assert_eq!(result.is_err(), expect_err); + if result.is_ok() { + assert_eq!( + previous_proof_height.load(std::sync::atomic::Ordering::Relaxed), + received, + "previous proof height should be updated" + ); + } + } + + #[test] + fn cloned_sdk_verify_proof_height() { + let sdk1 = SdkBuilder::new_mock() + .build() + .expect("mock Sdk should be created"); + + // First message verified, height 1. + let metadata = ResponseMetadata { + height: 1, + ..Default::default() + }; + + sdk1.verify_proof_metadata(&metadata) + .expect("proof should be valid"); + + assert_eq!( + sdk1.previous_proof_height + .load(std::sync::atomic::Ordering::Relaxed), + metadata.height, + "initial height" + ); + + // now, we clone sdk and do two requests. + let sdk2 = sdk1.clone(); + let sdk3 = sdk1.clone(); + + // Second message verified, height 2. + let metadata = ResponseMetadata { + height: 2, + ..Default::default() + }; + sdk2.verify_proof_metadata(&metadata) + .expect("proof should be valid"); + + assert_eq!( + sdk1.previous_proof_height + .load(std::sync::atomic::Ordering::Relaxed), + metadata.height, + "first sdk should see height from second sdk" + ); + assert_eq!( + sdk3.previous_proof_height + .load(std::sync::atomic::Ordering::Relaxed), + metadata.height, + "third sdk should see height from second sdk" + ); + + // Third message verified, height 3. + let metadata = ResponseMetadata { + height: 3, + ..Default::default() + }; + sdk3.verify_proof_metadata(&metadata) + .expect("proof should be valid"); + + assert_eq!( + sdk1.previous_proof_height + .load(std::sync::atomic::Ordering::Relaxed), + metadata.height, + "first sdk should see height from third sdk" + ); + + assert_eq!( + sdk2.previous_proof_height + .load(std::sync::atomic::Ordering::Relaxed), + metadata.height, + "second sdk should see height from third sdk" + ); + + // Now, using sdk1 for height 1 again should fail, as we are already at 3, with default tolerance 1. + let metadata = ResponseMetadata { + height: 1, + ..Default::default() + }; + + sdk1.verify_proof_metadata(&metadata) + .expect_err("proof should be invalid"); + } + + #[test_matrix([90,91,100,109,110], 100, 10, false; "valid time")] + #[test_matrix([0,89,111], 100, 10, true; "invalid time")] + #[test_matrix([0,100], [0,100], 100, false; "zero time")] + #[test_matrix([99,101], 100, 0, true; "zero tolerance")] + fn test_verify_proof_time(received: u64, now: u64, tolerance: u64, expect_err: bool) { + let metadata = ResponseMetadata { + time_ms: received, + ..Default::default() + }; + + let result = super::verify_proof_time(&metadata, now, tolerance); + + assert_eq!(result.is_err(), expect_err); + } +} From 34fee90173b3d78d8926814f2bc650aa750790dc Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 14:03:47 +0200 Subject: [PATCH 07/14] chore: simplify proof verification in sdk --- packages/rs-sdk/src/sdk.rs | 32 ++++++++------------------------ 1 file changed, 8 insertions(+), 24 deletions(-) diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index e8dd8e70b20..d6d99ae0a19 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -240,29 +240,10 @@ impl Sdk { where O::Request: Mockable, { - let provider = self - .context_provider() - .ok_or(drive_proof_verifier::Error::ContextProviderNotSet)?; - - let (object, mtd) = match self.inner { - SdkInstance::Dapi { .. } => O::maybe_from_proof_with_metadata( - request, - response, - self.network, - self.version(), - &provider, - ) - .map(|(a, b, _)| (a, b)), - #[cfg(feature = "mocks")] - SdkInstance::Mock { ref mock, .. } => { - let guard = mock.lock().await; - guard - .parse_proof_with_metadata(request, response) - .map(|(a, b, _)| (a, b)) - } - }?; + let (object, mtd, _proof) = self + .parse_proof_with_metadata_and_proof(request, response) + .await?; - self.verify_proof_metadata(&mtd)?; Ok((object, mtd)) } @@ -306,7 +287,7 @@ impl Sdk { .context_provider() .ok_or(drive_proof_verifier::Error::ContextProviderNotSet)?; - match self.inner { + let (object, mtd, proof) = match self.inner { SdkInstance::Dapi { .. } => O::maybe_from_proof_with_metadata( request, response, @@ -319,7 +300,10 @@ impl Sdk { let guard = mock.lock().await; guard.parse_proof_with_metadata(request, response) } - } + }?; + + self.verify_proof_metadata(&mtd)?; + Ok((object, mtd, proof)) } /// Return [ContextProvider] used by the SDK. From f7abb603f5e08908d27ed8cee3fce9c7d065fb2d Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 15:48:34 +0200 Subject: [PATCH 08/14] chore: self-review --- packages/rs-dapi-client/src/dapi_client.rs | 9 ++++-- packages/rs-drive-proof-verifier/src/error.rs | 16 +++++----- packages/rs-sdk/src/sdk.rs | 31 ++++++++++++------- 3 files changed, 34 insertions(+), 22 deletions(-) diff --git a/packages/rs-dapi-client/src/dapi_client.rs b/packages/rs-dapi-client/src/dapi_client.rs index 4a11059d7de..854c042d983 100644 --- a/packages/rs-dapi-client/src/dapi_client.rs +++ b/packages/rs-dapi-client/src/dapi_client.rs @@ -269,8 +269,13 @@ impl DapiRequestExecutor for DapiClient { .await; if let Err(error) = &result { - if !error.can_retry().unwrap_or(false) { - tracing::error!(?error, "request failed"); + match error.can_retry() { + Some(false) | None => { + tracing::error!(?error, "request failed"); + } + Some(true) => { + tracing::warn!(?error, "request failed, retrying"); + } } } diff --git a/packages/rs-drive-proof-verifier/src/error.rs b/packages/rs-drive-proof-verifier/src/error.rs index f14636aa51e..cf2196661fe 100644 --- a/packages/rs-drive-proof-verifier/src/error.rs +++ b/packages/rs-drive-proof-verifier/src/error.rs @@ -84,7 +84,7 @@ pub enum Error { ContextProviderError(#[from] ContextProviderError), /// Proof is stale; try another server - #[error("proof is stale; try another server")] + #[error(transparent)] StaleProof(#[from] StaleProofError), } @@ -92,24 +92,24 @@ pub enum Error { #[derive(Debug, thiserror::Error)] pub enum StaleProofError { /// Stale proof height - #[error("stale proof height: expected height {expected_height}, received {actual_height}, tolerance {tolerance}, try another server")] + #[error("stale proof height: expected height {expected_height}, received {received_height}, tolerance {tolerance_blocks}; try another server")] StaleProofHeight { /// Expected height - last block height seen by the Sdk expected_height: u64, /// Actual height - block height received from the server in the proof - actual_height: u64, + received_height: u64, /// Tolerance - how many blocks can be behind the expected height - tolerance: u64, + tolerance_blocks: u64, }, /// Proof time is stale #[error( - "received outdated proof time: expected {expected_ms} ms, received {actual_ms} ms, tolerance {tolerance_ms} ms, try another server" + "stale proof time: expected {expected_timestamp_ms}ms, received {received_timestamp_ms} ms, tolerance {tolerance_ms} ms; try another server" )] Time { /// Expected time in milliseconds - is local time when the proof was received - expected_ms: u64, - /// Actual time in milliseconds - time received from the server in the proof - actual_ms: u64, + expected_timestamp_ms: u64, + /// Time received from the server in the proof, in milliseconds + received_timestamp_ms: u64, /// Tolerance in milliseconds tolerance_ms: u64, }, diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index d6d99ae0a19..1d86d9ec2df 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -135,7 +135,7 @@ impl Clone for Sdk { internal_cache: Arc::clone(&self.internal_cache), context_provider: ArcSwapOption::new(self.context_provider.load_full()), cancel_token: self.cancel_token.clone(), - previous_proof_height: self.previous_proof_height.clone(), + previous_proof_height: Arc::clone(&self.previous_proof_height), proof_height_tolerance: self.proof_height_tolerance, proof_time_tolerance_ms: self.proof_time_tolerance_ms, #[cfg(feature = "mocks")] @@ -593,13 +593,13 @@ fn verify_proof_time( if now_ms.abs_diff(proof_time) > tolerance_ms { tracing::warn!( expected_time = now_ms, - actual_time = proof_time, + received_time = proof_time, tolerance_ms, "received proof with stale time; you should retry with another server" ); return Err(drive_proof_verifier::error::StaleProofError::Time { - expected_ms: now_ms, - actual_ms: proof_time, + expected_timestamp_ms: now_ms, + received_timestamp_ms: proof_time, tolerance_ms, } .into()); @@ -607,7 +607,7 @@ fn verify_proof_time( tracing::trace!( expected_time = now_ms, - actual_time = proof_time, + received_time = proof_time, tolerance_ms, "received proof with valid time" ); @@ -628,7 +628,7 @@ fn verify_proof_height( if received == prev { tracing::trace!( expected_height = prev, - actual_height = received, + received_height = received, tolerance, "received proof with the same height as previous" ); @@ -636,21 +636,23 @@ fn verify_proof_height( } // If received proof height is behind previous proof height by more than PROOF_HEIGHT_TOLERANCE, the proof is stale. + // // If prev is less than tolerance, then Sdk just started, so we just trust the proof, assuming we connected to a // trusted node. - // FIXME: in future, we need to implement t + // FIXME: in future, we need to securely trust some node (like a seed) from which we will fetch list of nodes and + // initial height. if prev > tolerance && received < prev - tolerance { tracing::warn!( expected_height = prev, - actual_height = received, + received_height = received, tolerance, "received proof with stale height; you should retry with another server" ); return Err( drive_proof_verifier::error::StaleProofError::StaleProofHeight { expected_height: prev, - actual_height: received, - tolerance, + received_height: received, + tolerance_blocks: tolerance, } .into(), ); @@ -659,7 +661,7 @@ fn verify_proof_height( // New proof is ahead of the previous proof, so we update the previous proof height. tracing::trace!( expected_height = prev, - actual_height = received, + received_height = received, tolerance, "received proof with new height" ); @@ -913,10 +915,15 @@ impl SdkBuilder { /// If current proof time differs from local time by more than this value, the proof is stale. /// If None, proof time is not checked. /// + /// This is set to `None` by default. + /// /// Note that enabling this check can cause issues if the local time is not synchronized with the network time, /// when the network is stalled or time between blocks increases significantly. /// - /// This is set to `None` by fefault. + /// Selecting a safe value for this parameter depends on maximum time between blocks mined on the network. + /// For example, if the network is configured to mine a block every maximum 3 minutes, setting this value + /// to a bit more than 6 minutes (to account for misbehaving proposers, network delays and local time + /// synchronization issues) should be safe. pub fn with_proof_time_tolerance(mut self, tolerance_ms: Option) -> Self { self.proof_time_tolerance_ms = tolerance_ms; self From f4b4e8b25a7119cf4dd1d9cc90fc93210885d869 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 15:57:48 +0200 Subject: [PATCH 09/14] refactor: err name --- packages/rs-drive-proof-verifier/src/error.rs | 4 ++-- packages/rs-sdk/src/sdk.rs | 14 ++++++-------- 2 files changed, 8 insertions(+), 10 deletions(-) diff --git a/packages/rs-drive-proof-verifier/src/error.rs b/packages/rs-drive-proof-verifier/src/error.rs index cf2196661fe..4a03830db10 100644 --- a/packages/rs-drive-proof-verifier/src/error.rs +++ b/packages/rs-drive-proof-verifier/src/error.rs @@ -93,10 +93,10 @@ pub enum Error { pub enum StaleProofError { /// Stale proof height #[error("stale proof height: expected height {expected_height}, received {received_height}, tolerance {tolerance_blocks}; try another server")] - StaleProofHeight { + Height { /// Expected height - last block height seen by the Sdk expected_height: u64, - /// Actual height - block height received from the server in the proof + /// Block height received from the server in the proof received_height: u64, /// Tolerance - how many blocks can be behind the expected height tolerance_blocks: u64, diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index 1d86d9ec2df..5750fe75ccc 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -648,14 +648,12 @@ fn verify_proof_height( tolerance, "received proof with stale height; you should retry with another server" ); - return Err( - drive_proof_verifier::error::StaleProofError::StaleProofHeight { - expected_height: prev, - received_height: received, - tolerance_blocks: tolerance, - } - .into(), - ); + return Err(drive_proof_verifier::error::StaleProofError::Height { + expected_height: prev, + received_height: received, + tolerance_blocks: tolerance, + } + .into()); } // New proof is ahead of the previous proof, so we update the previous proof height. From 8d7deba54f2e11b854a199fa35e04951007c4eae Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 17:02:48 +0200 Subject: [PATCH 10/14] refactor: some renames --- packages/rs-drive-proof-verifier/src/error.rs | 22 +- packages/rs-sdk/src/error.rs | 2 +- packages/rs-sdk/src/sdk.rs | 257 +++++++++--------- 3 files changed, 144 insertions(+), 137 deletions(-) diff --git a/packages/rs-drive-proof-verifier/src/error.rs b/packages/rs-drive-proof-verifier/src/error.rs index 4a03830db10..aa61b3281c3 100644 --- a/packages/rs-drive-proof-verifier/src/error.rs +++ b/packages/rs-drive-proof-verifier/src/error.rs @@ -83,32 +83,32 @@ pub enum Error { #[error("context provider error: {0}")] ContextProviderError(#[from] ContextProviderError), - /// Proof is stale; try another server + /// Remote node is stale; try another server #[error(transparent)] - StaleProof(#[from] StaleProofError), + StaleNode(#[from] StaleNodeError), } -/// Received proof is stale; try another server +/// Server returned stale metadata #[derive(Debug, thiserror::Error)] -pub enum StaleProofError { - /// Stale proof height - #[error("stale proof height: expected height {expected_height}, received {received_height}, tolerance {tolerance_blocks}; try another server")] +pub enum StaleNodeError { + /// Server returned metadata with outdated height + #[error("received height is outdated: expected {expected_height}, received {received_height}, tolerance {tolerance_blocks}; try another server")] Height { /// Expected height - last block height seen by the Sdk expected_height: u64, - /// Block height received from the server in the proof + /// Block height received from the server received_height: u64, /// Tolerance - how many blocks can be behind the expected height tolerance_blocks: u64, }, - /// Proof time is stale + /// Server returned metadata with time outside of the tolerance #[error( - "stale proof time: expected {expected_timestamp_ms}ms, received {received_timestamp_ms} ms, tolerance {tolerance_ms} ms; try another server" + "received invalid time: expected {expected_timestamp_ms}ms, received {received_timestamp_ms} ms, tolerance {tolerance_ms} ms; try another server" )] Time { - /// Expected time in milliseconds - is local time when the proof was received + /// Expected time in milliseconds - is local time when the message was received expected_timestamp_ms: u64, - /// Time received from the server in the proof, in milliseconds + /// Time received from the server in the message, in milliseconds received_timestamp_ms: u64, /// Tolerance in milliseconds tolerance_ms: u64, diff --git a/packages/rs-sdk/src/error.rs b/packages/rs-sdk/src/error.rs index 4d468557663..6fb572b4f36 100644 --- a/packages/rs-sdk/src/error.rs +++ b/packages/rs-sdk/src/error.rs @@ -85,7 +85,7 @@ impl CanRetry for Error { fn can_retry(&self) -> Option { let retry = matches!( self, - Error::Proof(drive_proof_verifier::Error::StaleProof(..)) + Error::Proof(drive_proof_verifier::Error::StaleNode(..)) | Error::DapiClientError(_) | Error::CoreClientError(_) | Error::TimeoutReached(_, _) diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index 5750fe75ccc..68251b6d255 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -101,24 +101,20 @@ pub struct Sdk { /// Note that setting this to None can panic. context_provider: ArcSwapOption>, - /// Last proof height; used to determine if the proof is stale. + /// Last seen height; used to determine if the remote node is stale. /// /// This is clone-able and can be shared between threads. - previous_proof_height: Arc, + metadata_last_seen_height: Arc, - /// How many blocks difference is allowed between the last proof height and the current proof height. - /// If current proof height is behind previous proof height by more than this value, the proof is stale. - /// If None, proof height is not checked. + /// How many blocks difference is allowed between the last height and the current height received in metadata. /// - /// This is set to `1` by default. - proof_height_tolerance: Option, + /// See [SdkBuilder::with_height_tolerance] for more information. + metadata_height_tolerance: Option, - /// How many milliseconds difference is allowed between the proof time and current local time. - /// If current proof time differs from local time by more than this value, the proof is stale. - /// If None, proof time is not checked. + /// How many milliseconds difference is allowed between the time received in response and current local time. /// - /// This is set to `None` by fefault. - proof_time_tolerance_ms: Option, + /// See [SdkBuilder::with_time_tolerance] for more information. + metadata_time_tolerance_ms: Option, /// Cancellation token; once cancelled, all pending requests should be aborted. pub(crate) cancel_token: CancellationToken, @@ -135,9 +131,9 @@ impl Clone for Sdk { internal_cache: Arc::clone(&self.internal_cache), context_provider: ArcSwapOption::new(self.context_provider.load_full()), cancel_token: self.cancel_token.clone(), - previous_proof_height: Arc::clone(&self.previous_proof_height), - proof_height_tolerance: self.proof_height_tolerance, - proof_time_tolerance_ms: self.proof_time_tolerance_ms, + metadata_last_seen_height: Arc::clone(&self.metadata_last_seen_height), + metadata_height_tolerance: self.metadata_height_tolerance, + metadata_time_tolerance_ms: self.metadata_time_tolerance_ms, #[cfg(feature = "mocks")] dump_dir: self.dump_dir.clone(), } @@ -240,28 +236,28 @@ impl Sdk { where O::Request: Mockable, { - let (object, mtd, _proof) = self + let (object, metadata, _proof) = self .parse_proof_with_metadata_and_proof(request, response) .await?; - Ok((object, mtd)) + Ok((object, metadata)) } - /// Verify proof metadata against the current state of the SDK. - fn verify_proof_metadata( + /// Verify response metadata against the current state of the SDK. + fn verify_response_metadata( &self, metadata: &ResponseMetadata, ) -> Result<(), drive_proof_verifier::Error> { - if let Some(height_tolerance) = self.proof_height_tolerance { - verify_proof_height( + if let Some(height_tolerance) = self.metadata_height_tolerance { + verify_metadata_height( metadata, height_tolerance, - Arc::clone(&(self.previous_proof_height)), + Arc::clone(&(self.metadata_last_seen_height)), )?; }; - if let Some(time_tolerance) = self.proof_time_tolerance_ms { + if let Some(time_tolerance) = self.metadata_time_tolerance_ms { let now = chrono::Utc::now().timestamp_millis() as u64; - verify_proof_time(metadata, now, time_tolerance)?; + verify_metadata_time(metadata, now, time_tolerance)?; }; Ok(()) @@ -287,7 +283,7 @@ impl Sdk { .context_provider() .ok_or(drive_proof_verifier::Error::ContextProviderNotSet)?; - let (object, mtd, proof) = match self.inner { + let (object, metadata, proof) = match self.inner { SdkInstance::Dapi { .. } => O::maybe_from_proof_with_metadata( request, response, @@ -302,8 +298,8 @@ impl Sdk { } }?; - self.verify_proof_metadata(&mtd)?; - Ok((object, mtd, proof)) + self.verify_response_metadata(&metadata)?; + Ok((object, metadata, proof)) } /// Return [ContextProvider] used by the SDK. @@ -574,32 +570,31 @@ impl Sdk { } } -/// If current proof time differs from local time by more than `tolerance`, the proof is considered stale. -/// This is used to verify the freshness of the proof. +/// If received metadata time differs from local time by more than `tolerance`, the remote node is considered stale. /// /// ## Parameters /// -/// - `metadata`: Metadata of the proof +/// - `metadata`: Metadata of the received response /// - `now_ms`: Current local time in milliseconds /// - `tolerance_ms`: Tolerance in milliseconds -fn verify_proof_time( +fn verify_metadata_time( metadata: &ResponseMetadata, now_ms: u64, tolerance_ms: u64, ) -> Result<(), drive_proof_verifier::Error> { - let proof_time = metadata.time_ms; + let metadata_time = metadata.time_ms; - // proof_time - tolerance <= now <= proof_time + tolerance - if now_ms.abs_diff(proof_time) > tolerance_ms { + // metadata_time - tolerance_ms <= now_ms <= metadata_time + tolerance_ms + if now_ms.abs_diff(metadata_time) > tolerance_ms { tracing::warn!( expected_time = now_ms, - received_time = proof_time, + received_time = metadata_time, tolerance_ms, - "received proof with stale time; you should retry with another server" + "received response with stale time; you should retry with another server" ); - return Err(drive_proof_verifier::error::StaleProofError::Time { + return Err(drive_proof_verifier::error::StaleNodeError::Time { expected_timestamp_ms: now_ms, - received_timestamp_ms: proof_time, + received_timestamp_ms: metadata_time, tolerance_ms, } .into()); @@ -607,70 +602,68 @@ fn verify_proof_time( tracing::trace!( expected_time = now_ms, - received_time = proof_time, + received_time = metadata_time, tolerance_ms, - "received proof with valid time" + "received response with valid time" ); Ok(()) } -/// If current proof height is behind previous proof height by more than `tolerance`, the proof is considered stale. -/// This is used to verify the freshness of the proof. -fn verify_proof_height( +/// If current metadata height is behind previously seen height by more than `tolerance`, the remote node +/// is considered stale. +fn verify_metadata_height( metadata: &ResponseMetadata, tolerance: u64, - previous_proof_height: Arc, + last_seen_height: Arc, ) -> Result<(), drive_proof_verifier::Error> { - let mut prev = previous_proof_height.load(Ordering::Relaxed); - let received = metadata.height; + let mut expected_height = last_seen_height.load(Ordering::Relaxed); + let received_height = metadata.height; // Same height, no need to update. - if received == prev { + if received_height == expected_height { tracing::trace!( - expected_height = prev, - received_height = received, + expected_height, + received_height, tolerance, - "received proof with the same height as previous" + "received message has the same height as previously seen" ); return Ok(()); } - // If received proof height is behind previous proof height by more than PROOF_HEIGHT_TOLERANCE, the proof is stale. - // - // If prev is less than tolerance, then Sdk just started, so we just trust the proof, assuming we connected to a - // trusted node. - // FIXME: in future, we need to securely trust some node (like a seed) from which we will fetch list of nodes and - // initial height. - if prev > tolerance && received < prev - tolerance { + // If expected_height <= tolerance, then Sdk just started, so we just assume what we got is correct. + if expected_height > tolerance && received_height < expected_height - tolerance { tracing::warn!( - expected_height = prev, - received_height = received, + expected_height, + received_height, tolerance, - "received proof with stale height; you should retry with another server" + "received message with stale height; you should retry with another server" ); - return Err(drive_proof_verifier::error::StaleProofError::Height { - expected_height: prev, - received_height: received, + return Err(drive_proof_verifier::error::StaleNodeError::Height { + expected_height, + received_height, tolerance_blocks: tolerance, } .into()); } - // New proof is ahead of the previous proof, so we update the previous proof height. + // New height is ahead of the last seen height, so we update the last seen height. tracing::trace!( - expected_height = prev, - received_height = received, + expected_height = expected_height, + received_height = received_height, tolerance, - "received proof with new height" + "received message with new height" ); - while let Err(stored) = - previous_proof_height.compare_exchange(prev, received, Ordering::SeqCst, Ordering::Relaxed) - { + while let Err(stored_height) = last_seen_height.compare_exchange( + expected_height, + received_height, + Ordering::SeqCst, + Ordering::Relaxed, + ) { // The value was changed to a higher value by another thread, so we need to retry. - if stored >= metadata.height { + if stored_height >= metadata.height { break; } - prev = stored; + expected_height = stored_height; } Ok(()) @@ -737,19 +730,16 @@ pub struct SdkBuilder { /// Context provider used by the SDK. context_provider: Option>, - /// How many blocks difference is allowed between the last proof height and the current proof height. - /// If current proof height is behind previous proof height by more than this value, the proof is stale. - /// If None, proof height is not checked. + /// How many blocks difference is allowed between the last seen metadata height and the height received in response + /// metadata. /// - /// This is set to `1` by default. - proof_height_tolerance: Option, + /// See [SdkBuilder::with_height_tolerance] for more information. + metadata_height_tolerance: Option, - /// How many milliseconds difference is allowed between the proof time and current local time. - /// If current proof time differs from local time by more than this value, the proof is stale. - /// If None, proof time is not checked. + /// How many milliseconds difference is allowed between the time received in response metadata and current local time. /// - /// This is set to `None` by fefault. - proof_time_tolerance_ms: Option, + /// See [SdkBuilder::with_time_tolerance] for more information. + metadata_time_tolerance_ms: Option, /// directory where dump files will be stored #[cfg(feature = "mocks")] @@ -772,8 +762,8 @@ impl Default for SdkBuilder { core_user: "".to_string(), proofs: true, - proof_height_tolerance: Some(1), - proof_time_tolerance_ms: None, + metadata_height_tolerance: Some(1), + metadata_time_tolerance_ms: None, #[cfg(feature = "mocks")] data_contract_cache_size: NonZeroUsize::new(DEFAULT_CONTRACT_CACHE_SIZE) @@ -898,20 +888,26 @@ impl SdkBuilder { self } - /// Change number of blocks difference allowed between the last proof height and the current proof height. + /// Change number of blocks difference allowed between the last height and the height received in current response. /// - /// If current proof height is behind previous proof height by more than this value, the proof is stale. - /// If None, proof height is not checked. + /// If height received in response metadata is behind previously seen height by more than this value, the node + /// is considered stale, and the request will fail. + /// + /// If None, the height is not checked. + /// + /// Note that this feature doesn't guarantee that you are getting latest data, but it significantly decreases + /// probability of getting old data. /// /// This is set to `1` by default. - pub fn with_proof_height_tolerance(mut self, tolerance: Option) -> Self { - self.proof_height_tolerance = tolerance; + pub fn with_height_tolerance(mut self, tolerance: Option) -> Self { + self.metadata_height_tolerance = tolerance; self } - /// How many milliseconds difference is allowed between the proof time and current local time. - /// If current proof time differs from local time by more than this value, the proof is stale. - /// If None, proof time is not checked. + /// How many milliseconds difference is allowed between the time received in response and current local time. + /// If the received time differs from local time by more than this value, the remote node is stale. + /// + /// If None, the time is not checked. /// /// This is set to `None` by default. /// @@ -922,8 +918,8 @@ impl SdkBuilder { /// For example, if the network is configured to mine a block every maximum 3 minutes, setting this value /// to a bit more than 6 minutes (to account for misbehaving proposers, network delays and local time /// synchronization issues) should be safe. - pub fn with_proof_time_tolerance(mut self, tolerance_ms: Option) -> Self { - self.proof_time_tolerance_ms = tolerance_ms; + pub fn with_time_tolerance(mut self, tolerance_ms: Option) -> Self { + self.metadata_time_tolerance_ms = tolerance_ms; self } @@ -970,9 +966,10 @@ impl SdkBuilder { context_provider: ArcSwapOption::new( self.context_provider.map(Arc::new)), cancel_token: self.cancel_token, internal_cache: Default::default(), - previous_proof_height: Arc::new(atomic::AtomicU64::new(0)), - proof_height_tolerance: self.proof_height_tolerance, - proof_time_tolerance_ms: self.proof_time_tolerance_ms, + // Note: in future, we need to securely initialize initial height during Sdk bootstrap or first request. + metadata_last_seen_height: Arc::new(atomic::AtomicU64::new(0)), + metadata_height_tolerance: self.metadata_height_tolerance, + metadata_time_tolerance_ms: self.metadata_time_tolerance_ms, #[cfg(feature = "mocks")] dump_dir: self.dump_dir, }; @@ -1036,9 +1033,9 @@ impl SdkBuilder { internal_cache: Default::default(), context_provider:ArcSwapAny::new( Some(Arc::new(context_provider))), cancel_token: self.cancel_token, - previous_proof_height: Arc::new(atomic::AtomicU64::new(0)), - proof_height_tolerance: self.proof_height_tolerance, - proof_time_tolerance_ms: self.proof_time_tolerance_ms, + metadata_last_seen_height: Arc::new(atomic::AtomicU64::new(0)), + metadata_height_tolerance: self.metadata_height_tolerance, + metadata_time_tolerance_ms: self.metadata_time_tolerance_ms, }; let mut guard = mock_sdk.try_lock().expect("mock sdk is in use by another thread and connot be reconfigured"); guard.set_sdk(sdk.clone()); @@ -1096,30 +1093,35 @@ mod test { #[test_matrix(97..102, 100, 2, false; "valid height")] #[test_case(103, 100, 2, true; "invalid height")] - fn test_verify_proof_height(expected: u64, received: u64, tolerance: u64, expect_err: bool) { + fn test_verify_metadata_height( + expected_height: u64, + received_height: u64, + tolerance: u64, + expect_err: bool, + ) { let metadata = ResponseMetadata { - height: received, + height: received_height, ..Default::default() }; - let previous_proof_height = - std::sync::Arc::new(std::sync::atomic::AtomicU64::new(expected)); + let last_seen_height = + std::sync::Arc::new(std::sync::atomic::AtomicU64::new(expected_height)); let result = - super::verify_proof_height(&metadata, tolerance, Arc::clone(&previous_proof_height)); + super::verify_metadata_height(&metadata, tolerance, Arc::clone(&last_seen_height)); assert_eq!(result.is_err(), expect_err); if result.is_ok() { assert_eq!( - previous_proof_height.load(std::sync::atomic::Ordering::Relaxed), - received, - "previous proof height should be updated" + last_seen_height.load(std::sync::atomic::Ordering::Relaxed), + received_height, + "previous height should be updated" ); } } #[test] - fn cloned_sdk_verify_proof_height() { + fn cloned_sdk_verify_metadata_height() { let sdk1 = SdkBuilder::new_mock() .build() .expect("mock Sdk should be created"); @@ -1130,11 +1132,11 @@ mod test { ..Default::default() }; - sdk1.verify_proof_metadata(&metadata) - .expect("proof should be valid"); + sdk1.verify_response_metadata(&metadata) + .expect("metadata should be valid"); assert_eq!( - sdk1.previous_proof_height + sdk1.metadata_last_seen_height .load(std::sync::atomic::Ordering::Relaxed), metadata.height, "initial height" @@ -1149,17 +1151,17 @@ mod test { height: 2, ..Default::default() }; - sdk2.verify_proof_metadata(&metadata) - .expect("proof should be valid"); + sdk2.verify_response_metadata(&metadata) + .expect("metadata should be valid"); assert_eq!( - sdk1.previous_proof_height + sdk1.metadata_last_seen_height .load(std::sync::atomic::Ordering::Relaxed), metadata.height, "first sdk should see height from second sdk" ); assert_eq!( - sdk3.previous_proof_height + sdk3.metadata_last_seen_height .load(std::sync::atomic::Ordering::Relaxed), metadata.height, "third sdk should see height from second sdk" @@ -1170,18 +1172,18 @@ mod test { height: 3, ..Default::default() }; - sdk3.verify_proof_metadata(&metadata) - .expect("proof should be valid"); + sdk3.verify_response_metadata(&metadata) + .expect("metadata should be valid"); assert_eq!( - sdk1.previous_proof_height + sdk1.metadata_last_seen_height .load(std::sync::atomic::Ordering::Relaxed), metadata.height, "first sdk should see height from third sdk" ); assert_eq!( - sdk2.previous_proof_height + sdk2.metadata_last_seen_height .load(std::sync::atomic::Ordering::Relaxed), metadata.height, "second sdk should see height from third sdk" @@ -1193,21 +1195,26 @@ mod test { ..Default::default() }; - sdk1.verify_proof_metadata(&metadata) - .expect_err("proof should be invalid"); + sdk1.verify_response_metadata(&metadata) + .expect_err("metadata should be invalid"); } #[test_matrix([90,91,100,109,110], 100, 10, false; "valid time")] #[test_matrix([0,89,111], 100, 10, true; "invalid time")] #[test_matrix([0,100], [0,100], 100, false; "zero time")] #[test_matrix([99,101], 100, 0, true; "zero tolerance")] - fn test_verify_proof_time(received: u64, now: u64, tolerance: u64, expect_err: bool) { + fn test_verify_metadata_time( + received_time: u64, + now_time: u64, + tolerance: u64, + expect_err: bool, + ) { let metadata = ResponseMetadata { - time_ms: received, + time_ms: received_time, ..Default::default() }; - let result = super::verify_proof_time(&metadata, now, tolerance); + let result = super::verify_metadata_time(&metadata, now_time, tolerance); assert_eq!(result.is_err(), expect_err); } From 58bf6b415ba9817506c2aca30d9ed2e577099dec Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 17:16:51 +0200 Subject: [PATCH 11/14] refactor: CanRetry returns bool, not Option --- packages/rs-dapi-client/src/dapi_client.rs | 21 +++++++------------ packages/rs-dapi-client/src/lib.rs | 7 ++----- packages/rs-dapi-client/src/transport/grpc.rs | 12 +++-------- packages/rs-sdk/src/error.rs | 12 +++-------- 4 files changed, 16 insertions(+), 36 deletions(-) diff --git a/packages/rs-dapi-client/src/dapi_client.rs b/packages/rs-dapi-client/src/dapi_client.rs index 854c042d983..611aa96316e 100644 --- a/packages/rs-dapi-client/src/dapi_client.rs +++ b/packages/rs-dapi-client/src/dapi_client.rs @@ -39,14 +39,14 @@ pub enum DapiClientError { } impl CanRetry for DapiClientError { - fn can_retry(&self) -> Option { + fn can_retry(&self) -> bool { use DapiClientError::*; match self { - NoAvailableAddresses => Some(true), + NoAvailableAddresses => true, Transport(transport_error, _) => transport_error.can_retry(), - AddressList(_) => Some(true), + AddressList(_) => true, #[cfg(feature = "mocks")] - Mock(_) => None, + Mock(_) => false, } } } @@ -233,7 +233,7 @@ impl DapiRequestExecutor for DapiClient { tracing::trace!(?response, "received {} response", response_name); } Err(error) => { - if !error.can_retry().unwrap_or(false) { + if !error.can_retry() { if applied_settings.ban_failed_address { let mut address_list = self .address_list @@ -264,18 +264,13 @@ impl DapiRequestExecutor for DapiClient { duration.as_secs_f32() ) }) - .when(|e| !e.can_retry().unwrap_or(false)) + .when(|e| !e.can_retry()) .instrument(tracing::info_span!("request routine")) .await; if let Err(error) = &result { - match error.can_retry() { - Some(false) | None => { - tracing::error!(?error, "request failed"); - } - Some(true) => { - tracing::warn!(?error, "request failed, retrying"); - } + if !error.can_retry() { + tracing::error!(?error, "request failed"); } } diff --git a/packages/rs-dapi-client/src/lib.rs b/packages/rs-dapi-client/src/lib.rs index 8b127f52279..760d9ce2e78 100644 --- a/packages/rs-dapi-client/src/lib.rs +++ b/packages/rs-dapi-client/src/lib.rs @@ -72,18 +72,15 @@ impl DapiRequest for T { } /// Returns true if the operation can be retried. -/// `None` means unspecified - in this case, you should either inspect the error -/// in more details, or assume `false`. pub trait CanRetry { /// Returns true if the operation can be retried safely. - /// None value means it's unspecified - fn can_retry(&self) -> Option; + fn can_retry(&self) -> bool; /// Get boolean flag that indicates if the error is retryable. /// /// Depreacted in favor of [CanRetry::can_retry]. #[deprecated = "Use !can_retry() instead"] fn is_node_failure(&self) -> bool { - !self.can_retry().unwrap_or(false) + !self.can_retry() } } diff --git a/packages/rs-dapi-client/src/transport/grpc.rs b/packages/rs-dapi-client/src/transport/grpc.rs index 735df4f0b66..e32d6c70f02 100644 --- a/packages/rs-dapi-client/src/transport/grpc.rs +++ b/packages/rs-dapi-client/src/transport/grpc.rs @@ -117,12 +117,12 @@ impl TransportClient for CoreGrpcClient { } impl CanRetry for dapi_grpc::tonic::Status { - fn can_retry(&self) -> Option { + fn can_retry(&self) -> bool { let code = self.code(); use dapi_grpc::tonic::Code::*; - let retry = !matches!( + !matches!( code, Ok | DataLoss | Cancelled @@ -132,13 +132,7 @@ impl CanRetry for dapi_grpc::tonic::Status { | Aborted | Internal | Unavailable - ); - - if retry { - Some(true) - } else { - None - } + ) } } diff --git a/packages/rs-sdk/src/error.rs b/packages/rs-sdk/src/error.rs index 6fb572b4f36..f3d8cb9019d 100644 --- a/packages/rs-sdk/src/error.rs +++ b/packages/rs-sdk/src/error.rs @@ -82,19 +82,13 @@ impl From for Error { } impl CanRetry for Error { - fn can_retry(&self) -> Option { - let retry = matches!( + fn can_retry(&self) -> bool { + matches!( self, Error::Proof(drive_proof_verifier::Error::StaleNode(..)) | Error::DapiClientError(_) | Error::CoreClientError(_) | Error::TimeoutReached(_, _) - ); - - if retry { - Some(true) - } else { - None - } + ) } } From 46aeeee51e67a7c2acd3070efa6fc32386db6fa1 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 17:23:12 +0200 Subject: [PATCH 12/14] refactor: move StaleNodeError to Sdk --- packages/rs-drive-proof-verifier/src/error.rs | 31 ----------------- packages/rs-sdk/src/error.rs | 33 ++++++++++++++++++- packages/rs-sdk/src/sdk.rs | 21 +++++------- 3 files changed, 41 insertions(+), 44 deletions(-) diff --git a/packages/rs-drive-proof-verifier/src/error.rs b/packages/rs-drive-proof-verifier/src/error.rs index aa61b3281c3..3203eb73174 100644 --- a/packages/rs-drive-proof-verifier/src/error.rs +++ b/packages/rs-drive-proof-verifier/src/error.rs @@ -82,37 +82,6 @@ pub enum Error { /// Context provider error #[error("context provider error: {0}")] ContextProviderError(#[from] ContextProviderError), - - /// Remote node is stale; try another server - #[error(transparent)] - StaleNode(#[from] StaleNodeError), -} - -/// Server returned stale metadata -#[derive(Debug, thiserror::Error)] -pub enum StaleNodeError { - /// Server returned metadata with outdated height - #[error("received height is outdated: expected {expected_height}, received {received_height}, tolerance {tolerance_blocks}; try another server")] - Height { - /// Expected height - last block height seen by the Sdk - expected_height: u64, - /// Block height received from the server - received_height: u64, - /// Tolerance - how many blocks can be behind the expected height - tolerance_blocks: u64, - }, - /// Server returned metadata with time outside of the tolerance - #[error( - "received invalid time: expected {expected_timestamp_ms}ms, received {received_timestamp_ms} ms, tolerance {tolerance_ms} ms; try another server" - )] - Time { - /// Expected time in milliseconds - is local time when the message was received - expected_timestamp_ms: u64, - /// Time received from the server in the message, in milliseconds - received_timestamp_ms: u64, - /// Tolerance in milliseconds - tolerance_ms: u64, - }, } /// Errors returned by the context provider diff --git a/packages/rs-sdk/src/error.rs b/packages/rs-sdk/src/error.rs index f3d8cb9019d..d9c7e8dd2bb 100644 --- a/packages/rs-sdk/src/error.rs +++ b/packages/rs-sdk/src/error.rs @@ -67,6 +67,10 @@ pub enum Error { /// Operation cancelled - cancel token was triggered, timeout, etc. #[error("Operation cancelled: {0}")] Cancelled(String), + + /// Remote node is stale; try another server + #[error(transparent)] + StaleNode(#[from] StaleNodeError), } impl From> for Error { @@ -85,10 +89,37 @@ impl CanRetry for Error { fn can_retry(&self) -> bool { matches!( self, - Error::Proof(drive_proof_verifier::Error::StaleNode(..)) + Error::StaleNode(..) | Error::DapiClientError(_) | Error::CoreClientError(_) | Error::TimeoutReached(_, _) ) } } + +/// Server returned stale metadata +#[derive(Debug, thiserror::Error)] +pub enum StaleNodeError { + /// Server returned metadata with outdated height + #[error("received height is outdated: expected {expected_height}, received {received_height}, tolerance {tolerance_blocks}; try another server")] + Height { + /// Expected height - last block height seen by the Sdk + expected_height: u64, + /// Block height received from the server + received_height: u64, + /// Tolerance - how many blocks can be behind the expected height + tolerance_blocks: u64, + }, + /// Server returned metadata with time outside of the tolerance + #[error( + "received invalid time: expected {expected_timestamp_ms}ms, received {received_timestamp_ms} ms, tolerance {tolerance_ms} ms; try another server" + )] + Time { + /// Expected time in milliseconds - is local time when the message was received + expected_timestamp_ms: u64, + /// Time received from the server in the message, in milliseconds + received_timestamp_ms: u64, + /// Tolerance in milliseconds + tolerance_ms: u64, + }, +} diff --git a/packages/rs-sdk/src/sdk.rs b/packages/rs-sdk/src/sdk.rs index 68251b6d255..2b9268f2751 100644 --- a/packages/rs-sdk/src/sdk.rs +++ b/packages/rs-sdk/src/sdk.rs @@ -1,6 +1,6 @@ //! [Sdk] entrypoint to Dash Platform. -use crate::error::Error; +use crate::error::{Error, StaleNodeError}; use crate::internal_cache::InternalSdkCache; use crate::mock::MockResponse; #[cfg(feature = "mocks")] @@ -211,7 +211,7 @@ impl Sdk { &self, request: O::Request, response: O::Response, - ) -> Result, drive_proof_verifier::Error> + ) -> Result, Error> where O::Request: Mockable, { @@ -232,7 +232,7 @@ impl Sdk { &self, request: O::Request, response: O::Response, - ) -> Result<(Option, ResponseMetadata), drive_proof_verifier::Error> + ) -> Result<(Option, ResponseMetadata), Error> where O::Request: Mockable, { @@ -244,10 +244,7 @@ impl Sdk { } /// Verify response metadata against the current state of the SDK. - fn verify_response_metadata( - &self, - metadata: &ResponseMetadata, - ) -> Result<(), drive_proof_verifier::Error> { + fn verify_response_metadata(&self, metadata: &ResponseMetadata) -> Result<(), Error> { if let Some(height_tolerance) = self.metadata_height_tolerance { verify_metadata_height( metadata, @@ -275,7 +272,7 @@ impl Sdk { &self, request: O::Request, response: O::Response, - ) -> Result<(Option, ResponseMetadata, Proof), drive_proof_verifier::Error> + ) -> Result<(Option, ResponseMetadata, Proof), Error> where O::Request: Mockable, { @@ -581,7 +578,7 @@ fn verify_metadata_time( metadata: &ResponseMetadata, now_ms: u64, tolerance_ms: u64, -) -> Result<(), drive_proof_verifier::Error> { +) -> Result<(), Error> { let metadata_time = metadata.time_ms; // metadata_time - tolerance_ms <= now_ms <= metadata_time + tolerance_ms @@ -592,7 +589,7 @@ fn verify_metadata_time( tolerance_ms, "received response with stale time; you should retry with another server" ); - return Err(drive_proof_verifier::error::StaleNodeError::Time { + return Err(StaleNodeError::Time { expected_timestamp_ms: now_ms, received_timestamp_ms: metadata_time, tolerance_ms, @@ -615,7 +612,7 @@ fn verify_metadata_height( metadata: &ResponseMetadata, tolerance: u64, last_seen_height: Arc, -) -> Result<(), drive_proof_verifier::Error> { +) -> Result<(), Error> { let mut expected_height = last_seen_height.load(Ordering::Relaxed); let received_height = metadata.height; @@ -638,7 +635,7 @@ fn verify_metadata_height( tolerance, "received message with stale height; you should retry with another server" ); - return Err(drive_proof_verifier::error::StaleNodeError::Height { + return Err(StaleNodeError::Height { expected_height, received_height, tolerance_blocks: tolerance, From 178fed5c4eca497bf5557d13d4f457ee4e767bba Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 17:36:22 +0200 Subject: [PATCH 13/14] chore: fix CanRetry --- packages/rs-dapi-client/src/dapi_client.rs | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/packages/rs-dapi-client/src/dapi_client.rs b/packages/rs-dapi-client/src/dapi_client.rs index 611aa96316e..468ca399749 100644 --- a/packages/rs-dapi-client/src/dapi_client.rs +++ b/packages/rs-dapi-client/src/dapi_client.rs @@ -42,9 +42,9 @@ impl CanRetry for DapiClientError { fn can_retry(&self) -> bool { use DapiClientError::*; match self { - NoAvailableAddresses => true, + NoAvailableAddresses => false, Transport(transport_error, _) => transport_error.can_retry(), - AddressList(_) => true, + AddressList(_) => false, #[cfg(feature = "mocks")] Mock(_) => false, } @@ -264,7 +264,7 @@ impl DapiRequestExecutor for DapiClient { duration.as_secs_f32() ) }) - .when(|e| !e.can_retry()) + .when(|e| e.can_retry()) .instrument(tracing::info_span!("request routine")) .await; From 65a15f549c0c3764213d6cec9b3456b0b468d7a1 Mon Sep 17 00:00:00 2001 From: Lukasz Klimek <842586+lklimek@users.noreply.github.com> Date: Fri, 18 Oct 2024 17:43:47 +0200 Subject: [PATCH 14/14] chore: apply feedback --- packages/rs-sdk/src/error.rs | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/packages/rs-sdk/src/error.rs b/packages/rs-sdk/src/error.rs index d9c7e8dd2bb..e55bda4742e 100644 --- a/packages/rs-sdk/src/error.rs +++ b/packages/rs-sdk/src/error.rs @@ -87,13 +87,7 @@ impl From for Error { impl CanRetry for Error { fn can_retry(&self) -> bool { - matches!( - self, - Error::StaleNode(..) - | Error::DapiClientError(_) - | Error::CoreClientError(_) - | Error::TimeoutReached(_, _) - ) + matches!(self, Error::StaleNode(..) | Error::TimeoutReached(_, _)) } }