pub mod api;
mod errors;
pub(crate) mod protocol_execution;
pub(crate) mod protocol_transport;
#[cfg(test)]
mod tests;
use std::{
collections::HashMap,
sync::{Arc, Mutex},
};
use entropy_protocol::{Listener, SessionId};
pub use self::{errors::*, protocol_execution::ProtocolMessage};
#[derive(Default, Debug, Clone)]
pub struct ListenerState {
pub listeners: Arc<Mutex<HashMap<SessionId, Listener>>>,
}
impl ListenerState {
pub fn contains_listener(&self, session_id: &SessionId) -> Result<bool, SubscribeErr> {
Ok(self
.listeners
.lock()
.map_err(|e| SubscribeErr::LockError(e.to_string()))?
.contains_key(session_id))
}
pub fn unsubscribed_peers(
&self,
session_id: &SessionId,
) -> Result<Vec<subxt::utils::AccountId32>, SubscribeErr> {
let listeners =
self.listeners.lock().map_err(|e| SubscribeErr::LockError(e.to_string()))?;
let listener =
listeners.get(session_id).ok_or(SubscribeErr::NoSessionId(session_id.clone()))?;
let unsubscribed_peers =
listener.validators.keys().map(|id| subxt::utils::AccountId32(*id)).collect();
Ok(unsubscribed_peers)
}
}