Refactor SyncState (#3297)

* Refactor SyncState

Method sync_error() retrun type was simplified.
update_txhashset_download() was  made type safe, which eliminates a runtime enum variant's  check, added an atomic status update
This commit is contained in:
hashmap
2020-04-20 12:30:04 +02:00
committed by GitHub
parent e64e90623b
commit d8c6eef485
7 changed files with 130 additions and 106 deletions
+2 -1
View File
@@ -46,5 +46,6 @@ pub use crate::chain::{Chain, MAX_ORPHAN_SIZE};
pub use crate::error::{Error, ErrorKind};
pub use crate::store::ChainStore;
pub use crate::types::{
BlockStatus, ChainAdapter, Options, SyncState, SyncStatus, Tip, TxHashsetWriteStatus,
BlockStatus, ChainAdapter, Options, SyncState, SyncStatus, Tip, TxHashsetDownloadStats,
TxHashsetWriteStatus,
};
+78 -31
View File
@@ -15,14 +15,13 @@
//! Base types that the block chain pipeline requires.
use chrono::prelude::{DateTime, Utc};
use std::sync::Arc;
use crate::core::core::hash::{Hash, Hashed, ZERO_HASH};
use crate::core::core::{Block, BlockHeader, HeaderVersion};
use crate::core::pow::Difficulty;
use crate::core::ser::{self, PMMRIndexHashable, Readable, Reader, Writeable, Writer};
use crate::error::{Error, ErrorKind};
use crate::util::RwLock;
use crate::util::{RwLock, RwLockWriteGuard};
bitflags! {
/// Options for block validation
@@ -40,7 +39,6 @@ bitflags! {
/// Various status sync can be in, whether it's fast sync or archival.
#[derive(Debug, Clone, Copy, Eq, PartialEq, Deserialize, Serialize)]
#[allow(missing_docs)]
pub enum SyncStatus {
/// Initial State (we do not yet know if we are/should be syncing)
Initial,
@@ -51,28 +49,27 @@ pub enum SyncStatus {
AwaitingPeers(bool),
/// Downloading block headers
HeaderSync {
/// current node height
current_height: u64,
/// height of the most advanced peer
highest_height: u64,
},
/// Downloading the various txhashsets
TxHashsetDownload {
start_time: DateTime<Utc>,
prev_update_time: DateTime<Utc>,
update_time: DateTime<Utc>,
prev_downloaded_size: u64,
downloaded_size: u64,
total_size: u64,
},
TxHashsetDownload(TxHashsetDownloadStats),
/// Setting up before validation
TxHashsetSetup,
/// Validating the kernels
TxHashsetKernelsValidation {
/// kernels validated
kernels: u64,
/// kernels in total
kernels_total: u64,
},
/// Validating the range proofs
TxHashsetRangeProofsValidation {
/// range proofs validated
rproofs: u64,
/// range proofs in total
rproofs_total: u64,
},
/// Finalizing the new state
@@ -81,16 +78,49 @@ pub enum SyncStatus {
TxHashsetDone,
/// Downloading blocks
BodySync {
/// current node height
current_height: u64,
/// height of the most advanced peer
highest_height: u64,
},
/// Shutdown
Shutdown,
}
/// Stats for TxHashsetDownload stage
#[derive(Debug, Clone, Copy, Eq, PartialEq, Deserialize, Serialize)]
pub struct TxHashsetDownloadStats {
/// when download started
pub start_time: DateTime<Utc>,
/// time of the previous update
pub prev_update_time: DateTime<Utc>,
/// time of the latest update
pub update_time: DateTime<Utc>,
/// size of the previous chunk
pub prev_downloaded_size: u64,
/// size of the the latest chunk
pub downloaded_size: u64,
/// downloaded since the start
pub total_size: u64,
}
impl Default for TxHashsetDownloadStats {
fn default() -> Self {
TxHashsetDownloadStats {
start_time: Utc::now(),
update_time: Utc::now(),
prev_update_time: Utc::now(),
prev_downloaded_size: 0,
downloaded_size: 0,
total_size: 0,
}
}
}
/// Current sync state. Encapsulates the current SyncStatus.
pub struct SyncState {
current: RwLock<SyncStatus>,
sync_error: Arc<RwLock<Option<Error>>>,
sync_error: RwLock<Option<Error>>,
}
impl SyncState {
@@ -98,7 +128,7 @@ impl SyncState {
pub fn new() -> SyncState {
SyncState {
current: RwLock::new(SyncStatus::Initial),
sync_error: Arc::new(RwLock::new(None)),
sync_error: RwLock::new(None),
}
}
@@ -114,37 +144,54 @@ impl SyncState {
}
/// Update the syncing status
pub fn update(&self, new_status: SyncStatus) {
if self.status() == new_status {
return;
}
let mut status = self.current.write();
debug!("sync_state: sync_status: {:?} -> {:?}", *status, new_status,);
*status = new_status;
pub fn update(&self, new_status: SyncStatus) -> bool {
let status = self.current.write();
self.update_with_guard(new_status, status)
}
/// Update txhashset downloading progress
pub fn update_txhashset_download(&self, new_status: SyncStatus) -> bool {
if let SyncStatus::TxHashsetDownload { .. } = new_status {
let mut status = self.current.write();
*status = new_status;
true
fn update_with_guard(
&self,
new_status: SyncStatus,
mut status: RwLockWriteGuard<SyncStatus>,
) -> bool {
if *status == new_status {
return false;
}
debug!("sync_state: sync_status: {:?} -> {:?}", *status, new_status,);
*status = new_status;
true
}
/// Update the syncing status if predicate f is satisfied
pub fn update_if<F>(&self, new_status: SyncStatus, f: F) -> bool
where
F: Fn(SyncStatus) -> bool,
{
let status = self.current.write();
if f(*status) {
self.update_with_guard(new_status, status)
} else {
false
}
}
/// Update txhashset downloading progress
pub fn update_txhashset_download(&self, stats: TxHashsetDownloadStats) {
*self.current.write() = SyncStatus::TxHashsetDownload(stats);
}
/// Communicate sync error
pub fn set_sync_error(&self, error: Error) {
*self.sync_error.write() = Some(error);
}
/// Get sync error
pub fn sync_error(&self) -> Arc<RwLock<Option<Error>>> {
Arc::clone(&self.sync_error)
pub fn sync_error(&self) -> Option<String> {
self.sync_error
.read()
.as_ref()
.and_then(|e| Some(e.to_string()))
}
/// Clear sync error