use tracing::{debug, info};
use super::{
error::SyncError,
peer_types::{Address, ConnectionState, PeerInfo, PeerStatus},
};
use crate::{Error, Result, Transaction, crdt::doc::path, store::DocStore};
pub(super) const PEERS_SUBTREE: &str = "peers"; pub(super) const TREES_SUBTREE: &str = "trees";
pub(super) struct PeerManager<'a> {
op: &'a Transaction,
}
impl<'a> PeerManager<'a> {
pub(super) fn new(op: &'a Transaction) -> Self {
Self { op }
}
pub(super) fn register_peer(
&self,
pubkey: impl Into<String>,
display_name: Option<&str>,
) -> Result<()> {
let pubkey = pubkey.into();
let peer_info = PeerInfo::new(&pubkey, display_name);
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
debug!(peer = %pubkey, display_name = ?display_name, "Registering new peer");
if peers.contains_path(path!(&pubkey as &str)) {
debug!(peer = %pubkey, "Peer already registered, skipping");
return Err(Error::Sync(SyncError::PeerAlreadyExists(pubkey.clone())));
}
peers.set_path(path!(&pubkey as &str, "pubkey"), peer_info.pubkey.clone())?;
if let Some(name) = &peer_info.display_name {
peers.set_path(path!(&pubkey as &str, "display_name"), name.clone())?;
}
peers.set_path(
path!(&pubkey as &str, "first_seen"),
peer_info.first_seen.clone(),
)?;
peers.set_path(
path!(&pubkey as &str, "last_seen"),
peer_info.last_seen.clone(),
)?;
peers.set_path(
path!(&pubkey as &str, "status"),
match peer_info.status {
PeerStatus::Active => "active".to_string(),
PeerStatus::Inactive => "inactive".to_string(),
PeerStatus::Blocked => "blocked".to_string(),
},
)?;
if !peer_info.addresses.is_empty() {
let addresses_json = serde_json::to_string(&peer_info.addresses).unwrap_or_default();
peers.set_path(path!(&pubkey as &str, "addresses"), addresses_json)?;
}
info!(peer = %pubkey, display_name = ?display_name, "Successfully registered new peer");
Ok(())
}
pub(super) fn update_peer_info(
&self,
pubkey: impl AsRef<str>,
peer_info: PeerInfo,
) -> Result<()> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
if !peers.contains_path_str(pubkey.as_ref()) {
return Err(Error::Sync(SyncError::PeerNotFound(
pubkey.as_ref().to_string(),
)));
}
peers.set_path(path!(pubkey.as_ref(), "pubkey"), peer_info.pubkey.clone())?;
if let Some(name) = &peer_info.display_name {
peers.set_path(path!(pubkey.as_ref(), "display_name"), name.clone())?;
}
peers.set_path(
path!(pubkey.as_ref(), "first_seen"),
peer_info.first_seen.clone(),
)?;
peers.set_path(
path!(pubkey.as_ref(), "last_seen"),
peer_info.last_seen.clone(),
)?;
let status_str = match peer_info.status {
PeerStatus::Active => "active",
PeerStatus::Inactive => "inactive",
PeerStatus::Blocked => "blocked",
};
peers.set_path(path!(pubkey.as_ref(), "status"), status_str.to_string())?;
let connection_state_str = match &peer_info.connection_state {
ConnectionState::Disconnected => "disconnected",
ConnectionState::Connecting => "connecting",
ConnectionState::Connected => "connected",
ConnectionState::Failed(msg) => &format!("failed:{msg}"),
};
peers.set_path(
path!(pubkey.as_ref(), "connection_state"),
connection_state_str.to_string(),
)?;
if let Some(last_sync) = &peer_info.last_successful_sync {
peers.set_path(
path!(pubkey.as_ref(), "last_successful_sync"),
last_sync.clone(),
)?;
}
peers.set_path(
path!(pubkey.as_ref(), "connection_attempts"),
peer_info.connection_attempts as i64,
)?;
if let Some(error) = &peer_info.last_error {
peers.set_path(path!(pubkey.as_ref(), "last_error"), error.clone())?;
}
if !peer_info.addresses.is_empty() {
let addresses_json = serde_json::to_string(&peer_info.addresses).unwrap_or_default();
peers.set_path(path!(pubkey.as_ref(), "addresses"), addresses_json)?;
}
debug!(peer = %pubkey.as_ref(), "Successfully updated peer information");
Ok(())
}
pub(super) fn update_peer_status(
&self,
pubkey: impl AsRef<str>,
status: PeerStatus,
) -> Result<()> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
if !peers.contains_path_str(pubkey.as_ref()) {
return Err(Error::Sync(SyncError::PeerNotFound(
pubkey.as_ref().to_string(),
)));
}
let status_str = match status {
PeerStatus::Active => "active",
PeerStatus::Inactive => "inactive",
PeerStatus::Blocked => "blocked",
};
peers.set_path(path!(pubkey.as_ref(), "status"), status_str.to_string())?;
let now = chrono::Utc::now().to_rfc3339();
peers.set_path(path!(pubkey.as_ref(), "last_seen"), now)?;
Ok(())
}
pub(super) fn get_peer_info(&self, pubkey: impl AsRef<str>) -> Result<Option<PeerInfo>> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
if !peers.contains_path_str(pubkey.as_ref()) {
return Ok(None);
}
let peer_pubkey = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "pubkey"))
.map_err(|_| {
Error::Sync(SyncError::SerializationError(
"Missing pubkey field".to_string(),
))
})?;
let display_name = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "display_name"))
.ok();
let first_seen = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "first_seen"))
.map_err(|_| {
Error::Sync(SyncError::SerializationError(
"Missing first_seen field".to_string(),
))
})?;
let last_seen = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "last_seen"))
.map_err(|_| {
Error::Sync(SyncError::SerializationError(
"Missing last_seen field".to_string(),
))
})?;
let status_str = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "status"))
.unwrap_or_else(|_| "active".to_string());
let status = match status_str.as_str() {
"active" => PeerStatus::Active,
"inactive" => PeerStatus::Inactive,
"blocked" => PeerStatus::Blocked,
_ => PeerStatus::Active, };
let connection_state_str = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "connection_state"))
.unwrap_or_else(|_| "disconnected".to_string());
let connection_state = match connection_state_str.as_str() {
"disconnected" => ConnectionState::Disconnected,
"connecting" => ConnectionState::Connecting,
"connected" => ConnectionState::Connected,
s if s.starts_with("failed:") => {
ConnectionState::Failed(s.strip_prefix("failed:").unwrap_or("").to_string())
}
_ => ConnectionState::Disconnected,
};
let last_successful_sync = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "last_successful_sync"))
.ok();
let connection_attempts = peers
.get_path_as::<i64>(path!(pubkey.as_ref(), "connection_attempts"))
.map(|v| v as u32)
.unwrap_or(0);
let last_error = peers
.get_path_as::<String>(path!(pubkey.as_ref(), "last_error"))
.ok();
let mut peer_info = PeerInfo {
pubkey: peer_pubkey,
display_name,
first_seen,
last_seen,
status,
addresses: Vec::new(),
connection_state,
last_successful_sync,
connection_attempts,
last_error,
};
if let Ok(addresses_json) = peers.get_path_as::<String>(path!(pubkey.as_ref(), "addresses"))
&& let Ok(addresses) = serde_json::from_str(&addresses_json)
{
peer_info.addresses = addresses;
}
if peer_info.status != PeerStatus::Blocked {
Ok(Some(peer_info))
} else {
Ok(None)
}
}
pub(super) fn list_peers(&self) -> Result<Vec<PeerInfo>> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
let all_peers = peers.get_all()?;
let mut peer_list = Vec::new();
for pubkey in all_peers.keys() {
if let Some(peer_info) = self.get_peer_info(pubkey)? {
peer_list.push(peer_info);
}
}
Ok(peer_list)
}
pub(super) fn remove_peer(&self, pubkey: impl AsRef<str>) -> Result<()> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
if peers.contains_path_str(pubkey.as_ref()) {
peers.set_path(path!(pubkey.as_ref(), "status"), "blocked".to_string())?;
}
let trees = self.op.get_store::<DocStore>(TREES_SUBTREE)?;
let all_keys = trees.get_all()?.keys().cloned().collect::<Vec<_>>();
for tree_id in all_keys {
let peer_list_path = path!(&tree_id, "peer_pubkeys");
if let Ok(peer_list_json) = trees.get_path_as::<String>(&peer_list_path)
&& let Ok(mut peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json)
{
let initial_len = peer_pubkeys.len();
peer_pubkeys.retain(|p| p.as_str() != pubkey.as_ref());
if peer_pubkeys.len() != initial_len {
if peer_pubkeys.is_empty() {
trees.delete(&tree_id)?;
} else {
let updated_json = serde_json::to_string(&peer_pubkeys).unwrap_or_default();
trees.set_path(&peer_list_path, updated_json)?;
}
}
}
}
Ok(())
}
pub(super) fn add_tree_sync(
&self,
peer_pubkey: impl AsRef<str>,
tree_root_id: impl AsRef<str>,
) -> Result<()> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
if !peers.contains_path_str(peer_pubkey.as_ref()) {
return Err(Error::Sync(SyncError::PeerNotFound(
peer_pubkey.as_ref().to_string(),
)));
}
let trees = self.op.get_store::<DocStore>(TREES_SUBTREE)?;
let peer_list_path = path!(tree_root_id.as_ref(), "peer_pubkeys");
let mut peer_pubkeys: Vec<String> = trees
.get_path_as::<String>(&peer_list_path)
.ok()
.and_then(|json| serde_json::from_str(&json).ok())
.unwrap_or_else(Vec::new);
if !peer_pubkeys.contains(&peer_pubkey.as_ref().to_string()) {
peer_pubkeys.push(peer_pubkey.as_ref().to_string());
let peer_list_json = serde_json::to_string(&peer_pubkeys).unwrap_or_default();
trees.set_path(&peer_list_path, peer_list_json)?;
trees.set_path(
path!(tree_root_id.as_ref(), "tree_id"),
tree_root_id.as_ref().to_string(),
)?;
} else {
debug!(peer = %peer_pubkey.as_ref(), tree = %tree_root_id.as_ref(), "Peer already syncing with tree");
}
Ok(())
}
pub(super) fn remove_tree_sync(
&self,
peer_pubkey: impl AsRef<str>,
tree_root_id: impl AsRef<str>,
) -> Result<()> {
info!(peer = %peer_pubkey.as_ref(), tree = %tree_root_id.as_ref(), "Removing tree sync relationship");
let trees = self.op.get_store::<DocStore>(TREES_SUBTREE)?;
let peer_list_path = path!(tree_root_id.as_ref(), "peer_pubkeys");
if let Ok(peer_list_json) = trees.get_path_as::<String>(&peer_list_path)
&& let Ok(mut peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json)
{
let initial_len = peer_pubkeys.len();
peer_pubkeys.retain(|p| p.as_str() != peer_pubkey.as_ref());
if peer_pubkeys.len() != initial_len {
if peer_pubkeys.is_empty() {
trees.delete(tree_root_id.as_ref())?;
} else {
let updated_json = serde_json::to_string(&peer_pubkeys).unwrap_or_default();
trees.set_path(&peer_list_path, updated_json)?;
}
}
}
Ok(())
}
pub(super) fn get_peer_trees(&self, peer_pubkey: impl AsRef<str>) -> Result<Vec<String>> {
let trees = self.op.get_store::<DocStore>(TREES_SUBTREE)?;
let all_trees = trees.get_all()?;
let mut synced_trees = Vec::new();
for tree_id in all_trees.keys() {
let peer_list_path = path!(tree_id, "peer_pubkeys");
if let Ok(peer_list_json) = trees.get_path_as::<String>(&peer_list_path)
&& let Ok(peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json)
&& peer_pubkeys.contains(&peer_pubkey.as_ref().to_string())
{
synced_trees.push(tree_id.clone());
}
}
Ok(synced_trees)
}
pub(super) fn get_tree_peers(&self, tree_root_id: impl AsRef<str>) -> Result<Vec<String>> {
let trees = self.op.get_store::<DocStore>(TREES_SUBTREE)?;
let peer_list_path = path!(tree_root_id.as_ref(), "peer_pubkeys");
Ok(trees
.get_path_as::<String>(&peer_list_path)
.ok()
.and_then(|json| serde_json::from_str(&json).ok())
.unwrap_or_else(Vec::new))
}
pub(super) fn is_tree_synced_with_peer(
&self,
peer_pubkey: impl AsRef<str>,
tree_root_id: impl AsRef<str>,
) -> Result<bool> {
let trees = self.op.get_store::<DocStore>(TREES_SUBTREE)?;
let peer_list_path = path!(tree_root_id.as_ref(), "peer_pubkeys");
match trees.get_path_as::<String>(&peer_list_path) {
Ok(peer_list_json) => {
if let Ok(peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json) {
Ok(peer_pubkeys.contains(&peer_pubkey.as_ref().to_string()))
} else {
Ok(false)
}
}
Err(_) => Ok(false),
}
}
pub(super) fn add_address(&self, peer_pubkey: impl AsRef<str>, address: Address) -> Result<()> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
if !peers.contains_path_str(peer_pubkey.as_ref()) {
return Err(Error::Sync(SyncError::PeerNotFound(
peer_pubkey.as_ref().to_string(),
)));
}
let mut all_addresses: Vec<Address> = peers
.get_path_as::<String>(path!(peer_pubkey.as_ref(), "addresses"))
.ok()
.and_then(|json| serde_json::from_str(&json).ok())
.unwrap_or_else(Vec::new);
if !all_addresses.contains(&address) {
all_addresses.push(address);
let addresses_json = serde_json::to_string(&all_addresses).unwrap_or_default();
peers.set_path(path!(peer_pubkey.as_ref(), "addresses"), addresses_json)?;
let now = chrono::Utc::now().to_rfc3339();
peers.set_path(path!(peer_pubkey.as_ref(), "last_seen"), now)?;
}
Ok(())
}
pub(super) fn remove_address(
&self,
peer_pubkey: impl AsRef<str>,
address: &Address,
) -> Result<bool> {
let peers = self.op.get_store::<DocStore>(PEERS_SUBTREE)?;
if !peers.contains_path_str(peer_pubkey.as_ref()) {
return Err(Error::Sync(SyncError::PeerNotFound(
peer_pubkey.as_ref().to_string(),
)));
}
let mut all_addresses: Vec<Address> = peers
.get_path_as::<String>(path!(peer_pubkey.as_ref(), "addresses"))
.ok()
.and_then(|json| serde_json::from_str(&json).ok())
.unwrap_or_else(Vec::new);
let initial_len = all_addresses.len();
all_addresses.retain(|a| a != address);
if all_addresses.len() != initial_len {
let addresses_json = serde_json::to_string(&all_addresses).unwrap_or_default();
peers.set_path(path!(peer_pubkey.as_ref(), "addresses"), addresses_json)?;
let now = chrono::Utc::now().to_rfc3339();
peers.set_path(path!(peer_pubkey.as_ref(), "last_seen"), now)?;
Ok(true)
} else {
Ok(false)
}
}
pub(super) fn get_addresses(
&self,
peer_pubkey: impl AsRef<str>,
transport_type: Option<&str>,
) -> Result<Vec<Address>> {
if let Some(peer_info) = self.get_peer_info(peer_pubkey)? {
match transport_type {
Some(transport) => Ok(peer_info
.get_addresses(transport)
.into_iter()
.cloned()
.collect()),
None => Ok(peer_info.addresses),
}
} else {
Ok(Vec::new())
}
}
}