use parking_lot::RwLock;
use std::{collections::HashMap, sync::Arc};
use crate::{
error::ConsensusError,
protos::consensus::v1::Proposal,
scope::ConsensusScope,
scope_config::ScopeConfig,
session::{ConsensusConfig, ConsensusSession},
};
pub trait ConsensusStorage<Scope>: Clone + Send + Sync + 'static
where
Scope: ConsensusScope,
{
fn save_session(&self, scope: &Scope, session: ConsensusSession) -> Result<(), ConsensusError>;
fn get_session(
&self,
scope: &Scope,
proposal_id: u32,
) -> Result<Option<ConsensusSession>, ConsensusError>;
fn remove_session(
&self,
scope: &Scope,
proposal_id: u32,
) -> Result<Option<ConsensusSession>, ConsensusError>;
fn list_scope_sessions(
&self,
scope: &Scope,
) -> Result<Option<Vec<ConsensusSession>>, ConsensusError>;
fn stream_scope_sessions(
&self,
scope: &Scope,
) -> impl Iterator<Item = Result<ConsensusSession, ConsensusError>>;
fn replace_scope_sessions(
&self,
scope: &Scope,
sessions: Vec<ConsensusSession>,
) -> Result<(), ConsensusError>;
fn list_scopes(&self) -> Result<Option<Vec<Scope>>, ConsensusError>;
fn update_session<R, F>(
&self,
scope: &Scope,
proposal_id: u32,
mutator: F,
) -> Result<R, ConsensusError>
where
F: FnOnce(&mut ConsensusSession) -> Result<R, ConsensusError>;
fn update_scope_sessions<F>(&self, scope: &Scope, mutator: F) -> Result<(), ConsensusError>
where
F: FnOnce(&mut Vec<ConsensusSession>) -> Result<(), ConsensusError>;
fn get_scope_config(&self, scope: &Scope) -> Result<Option<ScopeConfig>, ConsensusError>;
fn set_scope_config(&self, scope: &Scope, config: ScopeConfig) -> Result<(), ConsensusError>;
fn delete_scope(&self, scope: &Scope) -> Result<(), ConsensusError>;
fn update_scope_config<F>(&self, scope: &Scope, updater: F) -> Result<(), ConsensusError>
where
F: FnOnce(&mut ScopeConfig) -> Result<(), ConsensusError>;
fn get_consensus_result(
&self,
scope: &Scope,
proposal_id: u32,
) -> Result<bool, ConsensusError> {
use crate::session::ConsensusState;
let session = self
.get_session(scope, proposal_id)?
.ok_or(ConsensusError::SessionNotFound)?;
match session.state {
ConsensusState::ConsensusReached(result) => Ok(result),
ConsensusState::Failed => Err(ConsensusError::ConsensusFailed),
ConsensusState::Active => Err(ConsensusError::ConsensusNotReached),
}
}
fn get_proposal(&self, scope: &Scope, proposal_id: u32) -> Result<Proposal, ConsensusError> {
let session = self
.get_session(scope, proposal_id)?
.ok_or(ConsensusError::SessionNotFound)?;
Ok(session.proposal)
}
fn get_proposal_config(
&self,
scope: &Scope,
proposal_id: u32,
) -> Result<ConsensusConfig, ConsensusError> {
let session = self
.get_session(scope, proposal_id)?
.ok_or(ConsensusError::SessionNotFound)?;
Ok(session.config)
}
fn get_active_proposals(&self, scope: &Scope) -> Result<Vec<Proposal>, ConsensusError> {
let sessions = self.list_scope_sessions(scope)?.unwrap_or_default();
Ok(sessions
.into_iter()
.filter(|s| s.is_active())
.map(|s| s.proposal)
.collect())
}
fn get_reached_proposals(&self, scope: &Scope) -> Result<HashMap<u32, bool>, ConsensusError> {
let sessions = self.list_scope_sessions(scope)?.unwrap_or_default();
Ok(sessions
.into_iter()
.filter_map(|s| {
s.get_consensus_result()
.ok()
.map(|result| (s.proposal.proposal_id, result))
})
.collect())
}
}
#[derive(Clone)]
pub struct InMemoryConsensusStorage<Scope>
where
Scope: ConsensusScope,
{
sessions: Arc<RwLock<HashMap<Scope, HashMap<u32, ConsensusSession>>>>,
scope_configs: Arc<RwLock<HashMap<Scope, ScopeConfig>>>,
}
impl<Scope> Default for InMemoryConsensusStorage<Scope>
where
Scope: ConsensusScope,
{
fn default() -> Self {
Self {
sessions: Arc::new(RwLock::new(HashMap::new())),
scope_configs: Arc::new(RwLock::new(HashMap::new())),
}
}
}
impl<Scope> InMemoryConsensusStorage<Scope>
where
Scope: ConsensusScope,
{
pub fn new() -> Self {
Self::default()
}
}
impl<Scope> ConsensusStorage<Scope> for InMemoryConsensusStorage<Scope>
where
Scope: ConsensusScope + Clone,
{
fn save_session(&self, scope: &Scope, session: ConsensusSession) -> Result<(), ConsensusError> {
let mut sessions = self.sessions.write();
let entry = sessions.entry(scope.clone()).or_default();
entry.insert(session.proposal.proposal_id, session);
Ok(())
}
fn get_session(
&self,
scope: &Scope,
proposal_id: u32,
) -> Result<Option<ConsensusSession>, ConsensusError> {
let sessions = self.sessions.read();
Ok(sessions
.get(scope)
.and_then(|scope| scope.get(&proposal_id))
.cloned())
}
fn remove_session(
&self,
scope: &Scope,
proposal_id: u32,
) -> Result<Option<ConsensusSession>, ConsensusError> {
let mut sessions = self.sessions.write();
Ok(sessions
.get_mut(scope)
.and_then(|scope| scope.remove(&proposal_id)))
}
fn list_scope_sessions(
&self,
scope: &Scope,
) -> Result<Option<Vec<ConsensusSession>>, ConsensusError> {
let sessions = self.sessions.read();
let result = sessions
.get(scope)
.map(|scope| scope.values().cloned().collect::<Vec<ConsensusSession>>());
Ok(result)
}
fn stream_scope_sessions(
&self,
scope: &Scope,
) -> impl Iterator<Item = Result<ConsensusSession, ConsensusError>> {
let guard = self.sessions.read();
let sessions = guard
.get(scope)
.map(|inner_map| inner_map.values().cloned().collect::<Vec<_>>())
.unwrap_or_default();
sessions.into_iter().map(Ok)
}
fn replace_scope_sessions(
&self,
scope: &Scope,
sessions_list: Vec<ConsensusSession>,
) -> Result<(), ConsensusError> {
let mut sessions = self.sessions.write();
let new_map = sessions_list
.into_iter()
.map(|session| (session.proposal.proposal_id, session))
.collect();
sessions.insert(scope.clone(), new_map);
Ok(())
}
fn list_scopes(&self) -> Result<Option<Vec<Scope>>, ConsensusError> {
let sessions = self.sessions.read();
let result = sessions.keys().cloned().collect::<Vec<Scope>>();
if result.is_empty() {
return Ok(None);
}
Ok(Some(result))
}
fn update_session<R, F>(
&self,
scope: &Scope,
proposal_id: u32,
mutator: F,
) -> Result<R, ConsensusError>
where
F: FnOnce(&mut ConsensusSession) -> Result<R, ConsensusError>,
{
let mut sessions = self.sessions.write();
let session = sessions
.get_mut(scope)
.and_then(|scope_sessions| scope_sessions.get_mut(&proposal_id))
.ok_or(ConsensusError::SessionNotFound)?;
let result = mutator(session)?;
Ok(result)
}
fn update_scope_sessions<F>(&self, scope: &Scope, mutator: F) -> Result<(), ConsensusError>
where
F: FnOnce(&mut Vec<ConsensusSession>) -> Result<(), ConsensusError>,
{
let mut sessions = self.sessions.write();
let scope_sessions = sessions.entry(scope.clone()).or_default();
let mut sessions_vec: Vec<ConsensusSession> = scope_sessions.values().cloned().collect();
mutator(&mut sessions_vec)?;
if sessions_vec.is_empty() {
sessions.remove(scope);
return Ok(());
}
let new_map: HashMap<u32, ConsensusSession> = sessions_vec
.into_iter()
.map(|session| (session.proposal.proposal_id, session))
.collect();
*scope_sessions = new_map;
Ok(())
}
fn get_scope_config(&self, scope: &Scope) -> Result<Option<ScopeConfig>, ConsensusError> {
let configs = self.scope_configs.read();
Ok(configs.get(scope).cloned())
}
fn set_scope_config(&self, scope: &Scope, config: ScopeConfig) -> Result<(), ConsensusError> {
config.validate()?;
let mut configs = self.scope_configs.write();
configs.insert(scope.clone(), config);
Ok(())
}
fn delete_scope(&self, scope: &Scope) -> Result<(), ConsensusError> {
let mut sessions = self.sessions.write();
sessions.remove(scope);
drop(sessions);
let mut configs = self.scope_configs.write();
configs.remove(scope);
Ok(())
}
fn update_scope_config<F>(&self, scope: &Scope, updater: F) -> Result<(), ConsensusError>
where
F: FnOnce(&mut ScopeConfig) -> Result<(), ConsensusError>,
{
let mut configs = self.scope_configs.write();
let config = configs.entry(scope.clone()).or_default();
updater(config)?;
config.validate()?;
Ok(())
}
}