Skip to content
This repository was archived by the owner on Feb 1, 2024. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
32 changes: 24 additions & 8 deletions validator/sawtooth_validator/config/validator.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,8 @@ def load_default_validator_config():
peering='static',
scheduler='serial',
minimum_peer_connectivity=3,
maximum_peer_connectivity=10)
maximum_peer_connectivity=10,
state_pruning_block_depth=100)


def load_toml_validator_config(filename):
Expand Down Expand Up @@ -63,7 +64,7 @@ def load_toml_validator_config(filename):
'network_private_key', 'scheduler', 'permissions', 'roles',
'opentsdb_url', 'opentsdb_db', 'opentsdb_username',
'opentsdb_password', 'minimum_peer_connectivity',
'maximum_peer_connectivity'])
'maximum_peer_connectivity', 'state_pruning_block_depth'])
if invalid_keys:
raise LocalConfigurationError(
"Invalid keys in validator config: "
Expand Down Expand Up @@ -104,7 +105,9 @@ def load_toml_validator_config(filename):
minimum_peer_connectivity=toml_config.get(
"minimum_peer_connectivity", None),
maximum_peer_connectivity=toml_config.get(
"maximum_peer_connectivity", None)
"maximum_peer_connectivity", None),
state_pruning_block_depth=toml_config.get(
"state_pruning_block_depth", None)
)

return config
Expand Down Expand Up @@ -133,6 +136,7 @@ def merge_validator_config(configs):
opentsdb_password = None
minimum_peer_connectivity = None
maximum_peer_connectivity = None
state_pruning_block_depth = None

for config in reversed(configs):
if config.bind_network is not None:
Expand Down Expand Up @@ -169,6 +173,8 @@ def merge_validator_config(configs):
minimum_peer_connectivity = config.minimum_peer_connectivity
if config.maximum_peer_connectivity is not None:
maximum_peer_connectivity = config.maximum_peer_connectivity
if config.state_pruning_block_depth is not None:
state_pruning_block_depth = config.state_pruning_block_depth

return ValidatorConfig(
bind_network=bind_network,
Expand All @@ -187,7 +193,8 @@ def merge_validator_config(configs):
opentsdb_username=opentsdb_username,
opentsdb_password=opentsdb_password,
minimum_peer_connectivity=minimum_peer_connectivity,
maximum_peer_connectivity=maximum_peer_connectivity)
maximum_peer_connectivity=maximum_peer_connectivity,
state_pruning_block_depth=state_pruning_block_depth)


def parse_permissions(permissions):
Expand Down Expand Up @@ -235,7 +242,8 @@ def __init__(self, bind_network=None, bind_component=None,
roles=None, opentsdb_url=None, opentsdb_db=None,
opentsdb_username=None, opentsdb_password=None,
minimum_peer_connectivity=None,
maximum_peer_connectivity=None):
maximum_peer_connectivity=None,
state_pruning_block_depth=None):

self._bind_network = bind_network
self._bind_component = bind_component
Expand All @@ -254,6 +262,7 @@ def __init__(self, bind_network=None, bind_component=None,
self._opentsdb_password = opentsdb_password
self._minimum_peer_connectivity = minimum_peer_connectivity
self._maximum_peer_connectivity = maximum_peer_connectivity
self._state_pruning_block_depth = state_pruning_block_depth

@property
def bind_network(self):
Expand Down Expand Up @@ -323,6 +332,10 @@ def minimum_peer_connectivity(self):
def maximum_peer_connectivity(self):
return self._maximum_peer_connectivity

@property
def state_pruning_block_depth(self):
return self._state_pruning_block_depth

def __repr__(self):
# not including password for opentsdb
return (
Expand All @@ -331,7 +344,8 @@ def __repr__(self):
"network_public_key={}, network_private_key={}, "
"scheduler={}, permissions={}, roles={} "
"opentsdb_url={}, opentsdb_db={}, opentsdb_username={}, "
"minimum_peer_connectivity={}, maximum_peer_connectivity={})"
"minimum_peer_connectivity={}, maximum_peer_connectivity={}, "
"state_pruning_block_depth={})"
).format(
self.__class__.__name__,
repr(self._bind_network),
Expand All @@ -349,7 +363,8 @@ def __repr__(self):
repr(self._opentsdb_db),
repr(self._opentsdb_username),
repr(self._minimum_peer_connectivity),
repr(self._maximum_peer_connectivity))
repr(self._maximum_peer_connectivity),
repr(self._state_pruning_block_depth))

def to_dict(self):
return collections.OrderedDict([
Expand All @@ -369,7 +384,8 @@ def to_dict(self):
('opentsdb_username', self._opentsdb_username),
('opentsdb_password', self._opentsdb_password),
('minimum_peer_connectivity', self._minimum_peer_connectivity),
('maximum_peer_connectivity', self._maximum_peer_connectivity)
('maximum_peer_connectivity', self._maximum_peer_connectivity),
('state_pruning_block_depth', self._state_pruning_block_depth)
])

def to_toml_string(self):
Expand Down
4 changes: 4 additions & 0 deletions validator/sawtooth_validator/journal/chain.py
Original file line number Diff line number Diff line change
Expand Up @@ -41,8 +41,10 @@ def __init__(
block_store,
block_cache,
block_validator,
state_database,
chain_head_lock,
on_chain_updated,
state_pruning_block_depth=1000,
data_dir=None,
observers=None
):
Expand All @@ -59,9 +61,11 @@ def __init__(
ctypes.py_object(block_store),
ctypes.py_object(block_cache),
ctypes.py_object(block_validator),
state_database.pointer,
ctypes.py_object(chain_head_lock),
ctypes.py_object(on_chain_updated),
ctypes.py_object(observers),
ctypes.c_long(state_pruning_block_depth),
ctypes.c_char_p(data_dir.encode()),
ctypes.byref(self.pointer))

Expand Down
1 change: 1 addition & 0 deletions validator/sawtooth_validator/server/cli.py
Original file line number Diff line number Diff line change
Expand Up @@ -258,6 +258,7 @@ def main(args):
validator_config.permissions,
validator_config.minimum_peer_connectivity,
validator_config.maximum_peer_connectivity,
validator_config.state_pruning_block_depth,
validator_config.network_public_key,
validator_config.network_private_key,
roles=validator_config.roles)
Expand Down
3 changes: 3 additions & 0 deletions validator/sawtooth_validator/server/core.py
Original file line number Diff line number Diff line change
Expand Up @@ -79,6 +79,7 @@ def __init__(self,
permissions,
minimum_peer_connectivity,
maximum_peer_connectivity,
state_pruning_block_depth,
network_public_key=None,
network_private_key=None,
roles=None):
Expand Down Expand Up @@ -297,8 +298,10 @@ def __init__(self,
block_store=block_store,
block_cache=block_cache,
block_validator=block_validator,
state_database=global_state_db,
chain_head_lock=block_publisher.chain_head_lock,
on_chain_updated=block_publisher.on_chain_updated,
state_pruning_block_depth=state_pruning_block_depth,
data_dir=data_dir,
observers=[
event_broadcaster,
Expand Down
7 changes: 7 additions & 0 deletions validator/src/database/lmdb.rs
Original file line number Diff line number Diff line change
Expand Up @@ -213,6 +213,13 @@ impl<'a> LmdbDatabaseReaderCursor<'a> {
.map(|(key, value): (&[u8], &[u8])| (Vec::from(key), Vec::from(value)))
}

pub fn next(&mut self) -> Option<(Vec<u8>, Vec<u8>)> {
self.cursor
.next(&self.access)
.ok()
.map(|(key, value): (&[u8], &[u8])| (Vec::from(key), Vec::from(value)))
}

pub fn last(&mut self) -> Option<(Vec<u8>, Vec<u8>)> {
self.cursor
.last(&self.access)
Expand Down
4 changes: 4 additions & 0 deletions validator/src/journal/block_wrapper.rs
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,10 @@ impl BlockWrapper {
pub fn batches(&self) -> &[Batch] {
&self.block.batches
}

pub fn state_root_hash(&self) -> &str {
&self.block.state_root_hash
}
}

impl fmt::Display for BlockWrapper {
Expand Down
41 changes: 41 additions & 0 deletions validator/src/journal/chain.rs
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,7 @@ use journal;
use journal::block_validator::{BlockValidationResult, BlockValidator, ValidationError};
use journal::block_wrapper::BlockWrapper;
use metrics;
use state::state_pruning_manager::StatePruningManager;

use proto::transaction_receipt::TransactionReceipt;
use scheduler::TxnExecutionResult;
Expand Down Expand Up @@ -124,6 +125,10 @@ pub enum ChainReadError {
pub trait ChainReader: Send + Sync {
fn chain_head(&self) -> Result<Option<BlockWrapper>, ChainReadError>;
fn count_committed_transactions(&self) -> Result<usize, ChainReadError>;
fn get_block_by_block_num(
&self,
block_num: u64,
) -> Result<Option<BlockWrapper>, ChainReadError>;
}

pub trait ChainWriter: Send + Sync {
Expand All @@ -144,6 +149,8 @@ struct ChainControllerState<BC: BlockCache, BV: BlockValidator, CW: ChainWriter>
chain_id_manager: ChainIdManager,
chain_head_update_observer: Box<ChainHeadUpdateObserver>,
observers: Vec<Box<ChainObserver>>,

state_pruning_manager: StatePruningManager,
}

#[derive(Clone)]
Expand All @@ -152,6 +159,7 @@ pub struct ChainController<BC: BlockCache, BV: BlockValidator, CW: ChainWriter>
stop_handle: Arc<Mutex<Option<ChainThreadStopHandle>>>,
block_queue_sender: Option<Sender<BlockWrapper>>,
validation_result_sender: Option<Sender<(bool, BlockValidationResult)>>,
state_pruning_block_depth: u32,
}

impl<BC: BlockCache + 'static, BV: BlockValidator + 'static, CW: ChainWriter + 'static>
Expand All @@ -165,7 +173,9 @@ impl<BC: BlockCache + 'static, BV: BlockValidator + 'static, CW: ChainWriter + '
chain_head_lock: Box<ExternalLock + 'static>,
data_dir: String,
chain_head_update_observer: Box<ChainHeadUpdateObserver>,
state_pruning_block_depth: u32,
observers: Vec<Box<ChainObserver>>,
state_pruning_manager: StatePruningManager,
) -> Self {
let mut chain_controller = ChainController {
state: Arc::new(RwLock::new(ChainControllerState {
Expand All @@ -178,10 +188,12 @@ impl<BC: BlockCache + 'static, BV: BlockValidator + 'static, CW: ChainWriter + '
chain_head_update_observer,
observers,
chain_head: None,
state_pruning_manager,
})),
stop_handle: Arc::new(Mutex::new(None)),
block_queue_sender: None,
validation_result_sender: None,
state_pruning_block_depth,
};

chain_controller.initialize_chain_head();
Expand Down Expand Up @@ -262,6 +274,19 @@ impl<BC: BlockCache + 'static, BV: BlockValidator + 'static, CW: ChainWriter + '
let chain_head_block = new_block.clone();
state.chain_head = Some(new_block);

state.state_pruning_manager.update_queue(
&result
.new_chain
.iter()
.map(|block| block.state_root_hash())
.collect::<Vec<_>>(),
&result
.current_chain
.iter()
.map(|block| (block.block_num(), block.state_root_hash()))
.collect::<Vec<_>>(),
);

if let Err(err) = state
.chain_writer
.update_chain(&result.new_chain, &result.current_chain)
Expand Down Expand Up @@ -331,6 +356,21 @@ impl<BC: BlockCache + 'static, BV: BlockValidator + 'static, CW: ChainWriter + '
let mut committed_transactions_gauge =
COLLECTOR.gauge("ChainController.committed_transactions_gauge", None, None);
committed_transactions_gauge.set_value(total_committed_txns);

let chain_head_block_num = state.chain_head.as_ref().unwrap().block_num();
if chain_head_block_num + 1 > self.state_pruning_block_depth as u64 {
let prune_at = chain_head_block_num - (self.state_pruning_block_depth as u64);
match state.chain_reader.get_block_by_block_num(prune_at) {
Ok(Some(block)) => state
.state_pruning_manager
.add_to_queue(block.block_num(), block.state_root_hash()),
Ok(None) => warn!("No block at block height {}; ignoring...", prune_at),
Err(err) => error!("Unable to fetch block at height {}: {:?}", prune_at, err),
}

// Execute pruning:
state.state_pruning_manager.execute(prune_at)
}
} else {
info!("Rejected new chain head: {}", new_block);
}
Expand All @@ -346,6 +386,7 @@ impl<BC: BlockCache + 'static, BV: BlockValidator + 'static, CW: ChainWriter + '
stop_handle: Arc::new(Mutex::new(None)),
block_queue_sender: self.block_queue_sender.clone(),
validation_result_sender: self.validation_result_sender.clone(),
state_pruning_block_depth: self.state_pruning_block_depth,
}
}

Expand Down
34 changes: 34 additions & 0 deletions validator/src/journal/chain_ffi.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,11 +16,13 @@
*/
use cpython;
use cpython::{FromPyObject, ObjectProtocol, PyList, PyObject, Python, PythonObject, ToPyObject};
use database::lmdb::LmdbDatabase;
use journal::block_validator::{BlockValidationResult, BlockValidator, ValidationError};
use journal::block_wrapper::BlockWrapper;
use journal::chain::*;
use py_ffi;
use pylogger;
use state::state_pruning_manager::StatePruningManager;
use std::ffi::CStr;
use std::os::raw::{c_char, c_void};
use std::sync::mpsc::Sender;
Expand Down Expand Up @@ -55,16 +57,19 @@ pub extern "C" fn chain_controller_new(
block_store: *mut py_ffi::PyObject,
block_cache: *mut py_ffi::PyObject,
block_validator: *mut py_ffi::PyObject,
state_database: *const c_void,
chain_head_lock: *mut py_ffi::PyObject,
on_chain_updated: *mut py_ffi::PyObject,
observers: *mut py_ffi::PyObject,
state_pruning_block_depth: u32,
data_directory: *const c_char,
chain_controller_ptr: *mut *const c_void,
) -> ErrorCode {
check_null!(
block_store,
block_cache,
block_validator,
state_database,
chain_head_lock,
on_chain_updated,
observers,
Expand Down Expand Up @@ -100,6 +105,10 @@ pub extern "C" fn chain_controller_new(
return ErrorCode::InvalidPythonObject;
};

let state_database = unsafe { (*(state_database as *const LmdbDatabase)).clone() };

let state_pruning_manager = StatePruningManager::new(state_database);

let chain_controller = ChainController::new(
PyBlockCache::new(py_block_cache),
PyBlockValidator::new(py_block_validator),
Expand All @@ -108,7 +117,9 @@ pub extern "C" fn chain_controller_new(
Box::new(chain_head_lock),
data_dir.into(),
Box::new(PyChainHeadUpdateObserver::new(py_on_chain_updated)),
state_pruning_block_depth,
observer_wrappers,
state_pruning_manager,
);

unsafe {
Expand Down Expand Up @@ -575,6 +586,29 @@ impl ChainReader for PyBlockStore {
})
}

fn get_block_by_block_num(
&self,
block_num: u64,
) -> Result<Option<BlockWrapper>, ChainReadError> {
let gil_guard = Python::acquire_gil();
let py = gil_guard.python();

self.py_block_store
.call_method(py, "get_block_by_number", (block_num,), None)
.and_then(|result| result.extract(py))
.or_else(|py_err| {
if py_err.get_type(py).name(py) == "KeyError" {
Ok(None)
} else {
Err(py_err)
}
})
.map_err(|py_err| {
pylogger::exception(py, "Unable to call block_store.chain_head", py_err);
ChainReadError::GeneralReadError("Unable to read from python block store".into())
})
}

fn count_committed_transactions(&self) -> Result<usize, ChainReadError> {
let gil_guard = Python::acquire_gil();
let py = gil_guard.python();
Expand Down
1 change: 0 additions & 1 deletion validator/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,6 @@
*/

extern crate cbor;
#[macro_use]
extern crate cpython;
extern crate crypto;
extern crate hex;
Expand Down
Loading