1use tracing::{debug, info};
7
8use super::{
9 error::SyncError,
10 peer_types::{Address, ConnectionState, PeerInfo, PeerStatus},
11};
12use crate::{
13 Error, Result, Transaction, auth::crypto::PublicKey, crdt::doc::path, entry::ID,
14 store::DocStore, sync::PeerId,
15};
16
17pub(super) const PEERS_SUBTREE: &str = "peers"; pub(super) const TREES_SUBTREE: &str = "trees"; pub(super) struct PeerManager<'a> {
26 txn: &'a Transaction,
27}
28
29impl<'a> PeerManager<'a> {
30 pub(super) fn new(txn: &'a Transaction) -> Self {
32 Self { txn }
33 }
34
35 pub(super) async fn register_peer(
44 &self,
45 pubkey: &PublicKey,
46 display_name: Option<&str>,
47 ) -> Result<()> {
48 let pk_str = pubkey.to_string();
49 let now = self.txn.now_rfc3339()?;
50 let peer_info = PeerInfo::new_at(pubkey, display_name, now);
51 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
52
53 if peers.contains_path(path!(&pk_str)).await {
55 debug!(peer = %pk_str, "Peer already registered, skipping");
56 return Err(Error::Sync(Box::new(SyncError::PeerAlreadyExists(pk_str))));
57 }
58
59 debug!(peer = %pk_str, display_name = ?display_name, "Registering new peer");
60
61 peers
63 .set_path(path!(&pk_str, "pubkey"), peer_info.id.to_string())
64 .await?;
65 if let Some(name) = &peer_info.display_name {
66 peers
67 .set_path(path!(&pk_str, "display_name"), name.clone())
68 .await?;
69 }
70 peers
71 .set_path(path!(&pk_str, "first_seen"), peer_info.first_seen.clone())
72 .await?;
73 peers
74 .set_path(path!(&pk_str, "last_seen"), peer_info.last_seen.clone())
75 .await?;
76 peers
77 .set_path(
78 path!(&pk_str, "status"),
79 match peer_info.status {
80 PeerStatus::Active => "active".to_string(),
81 PeerStatus::Inactive => "inactive".to_string(),
82 PeerStatus::Blocked => "blocked".to_string(),
83 },
84 )
85 .await?;
86
87 if !peer_info.addresses.is_empty() {
89 let addresses_json = serde_json::to_string(&peer_info.addresses).unwrap_or_default();
90 peers
91 .set_path(path!(&pk_str, "addresses"), addresses_json)
92 .await?;
93 }
94
95 info!(peer = %pk_str, display_name = ?display_name, "Successfully registered new peer");
96 Ok(())
97 }
98
99 pub(super) async fn update_peer_info(
108 &self,
109 pubkey: &PublicKey,
110 peer_info: PeerInfo,
111 ) -> Result<()> {
112 let pk_str = pubkey.to_string();
113 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
114
115 if !peers.contains_path_str(&pk_str).await {
117 return Err(Error::Sync(Box::new(SyncError::PeerNotFound(pk_str))));
118 }
119
120 peers
122 .set_path(path!(&pk_str, "pubkey"), peer_info.id.to_string())
123 .await?;
124
125 if let Some(name) = &peer_info.display_name {
126 peers
127 .set_path(path!(&pk_str, "display_name"), name.clone())
128 .await?;
129 }
130
131 peers
132 .set_path(path!(&pk_str, "first_seen"), peer_info.first_seen.clone())
133 .await?;
134
135 peers
136 .set_path(path!(&pk_str, "last_seen"), peer_info.last_seen.clone())
137 .await?;
138
139 let status_str = match peer_info.status {
141 PeerStatus::Active => "active",
142 PeerStatus::Inactive => "inactive",
143 PeerStatus::Blocked => "blocked",
144 };
145 peers
146 .set_path(path!(&pk_str, "status"), status_str.to_string())
147 .await?;
148
149 let connection_state_str = match &peer_info.connection_state {
151 ConnectionState::Disconnected => "disconnected",
152 ConnectionState::Connecting => "connecting",
153 ConnectionState::Connected => "connected",
154 ConnectionState::Failed(msg) => &format!("failed:{msg}"),
155 };
156 peers
157 .set_path(
158 path!(&pk_str, "connection_state"),
159 connection_state_str.to_string(),
160 )
161 .await?;
162
163 peers
165 .set_path(
166 path!(&pk_str, "connection_attempts"),
167 peer_info.connection_attempts as i64,
168 )
169 .await?;
170
171 if let Some(error) = &peer_info.last_error {
172 peers
173 .set_path(path!(&pk_str, "last_error"), error.clone())
174 .await?;
175 }
176
177 if !peer_info.addresses.is_empty() {
179 let addresses_json = serde_json::to_string(&peer_info.addresses).unwrap_or_default();
180 peers
181 .set_path(path!(&pk_str, "addresses"), addresses_json)
182 .await?;
183 }
184
185 debug!(peer = %pk_str, "Successfully updated peer information");
186 Ok(())
187 }
188
189 pub(super) async fn update_peer_status(
198 &self,
199 pubkey: &PublicKey,
200 status: PeerStatus,
201 ) -> Result<()> {
202 let pk_str = pubkey.to_string();
203 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
204
205 if !peers.contains_path_str(&pk_str).await {
207 return Err(Error::Sync(Box::new(SyncError::PeerNotFound(pk_str))));
208 }
209
210 let status_str = match status {
212 PeerStatus::Active => "active",
213 PeerStatus::Inactive => "inactive",
214 PeerStatus::Blocked => "blocked",
215 };
216 peers
217 .set_path(path!(&pk_str, "status"), status_str.to_string())
218 .await?;
219
220 let now = self.txn.now_rfc3339()?;
222 peers.set_path(path!(&pk_str, "last_seen"), now).await?;
223
224 Ok(())
225 }
226
227 pub(super) async fn get_peer_info(&self, pubkey: &PublicKey) -> Result<Option<PeerInfo>> {
235 self.get_peer_info_str(&pubkey.to_string()).await
236 }
237
238 async fn get_peer_info_str(&self, pk_str: &str) -> Result<Option<PeerInfo>> {
240 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
241
242 if !peers.contains_path_str(pk_str).await {
244 return Ok(None);
245 }
246
247 let peer_pubkey = peers
249 .get_path_as::<String>(path!(pk_str, "pubkey"))
250 .await
251 .map_err(|_| {
252 Error::Sync(Box::new(SyncError::SerializationError(
253 "Missing pubkey field".to_string(),
254 )))
255 })?;
256
257 let display_name = peers
258 .get_path_as::<String>(path!(pk_str, "display_name"))
259 .await
260 .ok();
261
262 let first_seen = peers
263 .get_path_as::<String>(path!(pk_str, "first_seen"))
264 .await
265 .map_err(|_| {
266 Error::Sync(Box::new(SyncError::SerializationError(
267 "Missing first_seen field".to_string(),
268 )))
269 })?;
270
271 let last_seen = peers
272 .get_path_as::<String>(path!(pk_str, "last_seen"))
273 .await
274 .map_err(|_| {
275 Error::Sync(Box::new(SyncError::SerializationError(
276 "Missing last_seen field".to_string(),
277 )))
278 })?;
279
280 let status_str = peers
281 .get_path_as::<String>(path!(pk_str, "status"))
282 .await
283 .unwrap_or_else(|_| "active".to_string());
284 let status = match status_str.as_str() {
285 "active" => PeerStatus::Active,
286 "inactive" => PeerStatus::Inactive,
287 "blocked" => PeerStatus::Blocked,
288 _ => PeerStatus::Active, };
290
291 let connection_state_str = peers
293 .get_path_as::<String>(path!(pk_str, "connection_state"))
294 .await
295 .unwrap_or_else(|_| "disconnected".to_string());
296 let connection_state = match connection_state_str.as_str() {
297 "disconnected" => ConnectionState::Disconnected,
298 "connecting" => ConnectionState::Connecting,
299 "connected" => ConnectionState::Connected,
300 s if s.starts_with("failed:") => {
301 ConnectionState::Failed(s.strip_prefix("failed:").unwrap_or("").to_string())
302 }
303 _ => ConnectionState::Disconnected,
304 };
305
306 let connection_attempts = peers
307 .get_path_as::<i64>(path!(pk_str, "connection_attempts"))
308 .await
309 .map(|v| v as u32)
310 .unwrap_or(0);
311
312 let last_error = peers
313 .get_path_as::<String>(path!(pk_str, "last_error"))
314 .await
315 .ok();
316
317 let mut peer_info = PeerInfo {
318 id: PeerId::new(PublicKey::from_prefixed_string(&peer_pubkey)?),
319 display_name,
320 first_seen,
321 last_seen,
322 status,
323 addresses: Vec::new(),
324 connection_state,
325 connection_attempts,
326 last_error,
327 };
328
329 if let Ok(addresses_json) = peers
331 .get_path_as::<String>(path!(pk_str, "addresses"))
332 .await
333 && let Ok(addresses) = serde_json::from_str(&addresses_json)
334 {
335 peer_info.addresses = addresses;
336 }
337
338 if peer_info.status != PeerStatus::Blocked {
340 Ok(Some(peer_info))
341 } else {
342 Ok(None)
343 }
344 }
345
346 pub(super) async fn list_peers(&self) -> Result<Vec<PeerInfo>> {
351 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
352 let all_peers = peers.get_all().await?;
353 let mut peer_list = Vec::new();
354
355 for pubkey_str in all_peers.keys() {
357 if let Some(peer_info) = self.get_peer_info_str(pubkey_str).await? {
359 peer_list.push(peer_info);
360 }
361 }
362
363 Ok(peer_list)
364 }
365
366 pub(super) async fn remove_peer(&self, pubkey: &PublicKey) -> Result<()> {
376 let pk_str = pubkey.to_string();
377 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
378
379 if peers.contains_path_str(&pk_str).await {
381 peers
382 .set_path(path!(&pk_str, "status"), "blocked".to_string())
383 .await?;
384 }
385
386 let trees = self.txn.get_store::<DocStore>(TREES_SUBTREE).await?;
388 let all_keys = trees.get_all().await?.keys().cloned().collect::<Vec<_>>();
389 for tree_id in all_keys {
390 let peer_list_path = path!(&tree_id, "peer_pubkeys");
391 if let Ok(peer_list_json) = trees.get_path_as::<String>(&peer_list_path).await
392 && let Ok(mut peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json)
393 {
394 let initial_len = peer_pubkeys.len();
395 peer_pubkeys.retain(|p| p.as_str() != pk_str);
396
397 if peer_pubkeys.len() != initial_len {
398 if peer_pubkeys.is_empty() {
400 trees.delete(&tree_id).await?;
401 } else {
402 let updated_json = serde_json::to_string(&peer_pubkeys).unwrap_or_default();
403 trees.set_path(&peer_list_path, updated_json).await?;
404 }
405 }
406 }
407 }
408
409 Ok(())
410 }
411
412 pub(super) async fn add_tree_sync(
423 &self,
424 peer_pubkey: &PublicKey,
425 tree_root_id: &ID,
426 ) -> Result<()> {
427 let pk_str = peer_pubkey.to_string();
428 let tree_root_str = tree_root_id.to_string();
429
430 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
432 if !peers.contains_path_str(&pk_str).await {
433 return Err(Error::Sync(Box::new(SyncError::PeerNotFound(pk_str))));
434 }
435
436 let trees = self.txn.get_store::<DocStore>(TREES_SUBTREE).await?;
437
438 let peer_list_path = path!(&tree_root_str, "peer_pubkeys");
440 let peer_list_result = trees.get_path_as::<String>(&peer_list_path).await;
441 let mut peer_pubkeys: Vec<String> = peer_list_result
442 .ok()
443 .and_then(|json| serde_json::from_str(&json).ok())
444 .unwrap_or_else(Vec::new);
445
446 if !peer_pubkeys.contains(&pk_str) {
448 peer_pubkeys.push(pk_str.clone());
449
450 let peer_list_json = serde_json::to_string(&peer_pubkeys).unwrap_or_default();
452 trees.set_path(&peer_list_path, peer_list_json).await?;
453
454 trees
456 .set_path(path!(&tree_root_str, "tree_id"), tree_root_str.clone())
457 .await?;
458 } else {
459 debug!(peer = %pk_str, tree = %tree_root_str, "Peer already syncing with tree");
460 }
461
462 Ok(())
463 }
464
465 pub(super) async fn remove_tree_sync(
474 &self,
475 peer_pubkey: &PublicKey,
476 tree_root_id: &ID,
477 ) -> Result<()> {
478 let pk_str = peer_pubkey.to_string();
479 let tree_root_str = tree_root_id.to_string();
480 info!(peer = %pk_str, tree = %tree_root_str, "Removing tree sync relationship");
481 let trees = self.txn.get_store::<DocStore>(TREES_SUBTREE).await?;
482
483 let peer_list_path = path!(&tree_root_str, "peer_pubkeys");
485 if let Ok(peer_list_json) = trees.get_path_as::<String>(&peer_list_path).await
486 && let Ok(mut peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json)
487 {
488 let initial_len = peer_pubkeys.len();
490 peer_pubkeys.retain(|p| p.as_str() != pk_str);
491
492 if peer_pubkeys.len() != initial_len {
493 if peer_pubkeys.is_empty() {
495 trees.delete(&tree_root_str).await?;
497 } else {
498 let updated_json = serde_json::to_string(&peer_pubkeys).unwrap_or_default();
500 trees.set_path(&peer_list_path, updated_json).await?;
501 }
502 }
503 }
504
505 Ok(())
506 }
507
508 pub(super) async fn get_peer_trees(&self, peer_pubkey: &PublicKey) -> Result<Vec<ID>> {
516 let pk_str = peer_pubkey.to_string();
517 let trees = self.txn.get_store::<DocStore>(TREES_SUBTREE).await?;
518 let all_trees = trees.get_all().await?;
519 let mut synced_trees = Vec::new();
520
521 for tree_id_str in all_trees.keys() {
522 let peer_list_path = path!(tree_id_str, "peer_pubkeys");
523 if let Ok(peer_list_json) = trees.get_path_as::<String>(&peer_list_path).await
524 && let Ok(peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json)
525 && peer_pubkeys.contains(&pk_str)
526 && let Ok(id) = ID::parse(tree_id_str)
527 {
528 synced_trees.push(id);
529 }
530 }
531
532 Ok(synced_trees)
533 }
534
535 pub(super) async fn get_tree_peers(&self, tree_root_id: &ID) -> Result<Vec<PeerId>> {
543 let tree_root_str = tree_root_id.to_string();
544 let trees = self.txn.get_store::<DocStore>(TREES_SUBTREE).await?;
545 let peer_list_path = path!(&tree_root_str, "peer_pubkeys");
546 let peer_list_result = trees.get_path_as::<String>(&peer_list_path).await;
547 let string_vec: Vec<String> = peer_list_result
548 .ok()
549 .and_then(|json| serde_json::from_str(&json).ok())
550 .unwrap_or_else(Vec::new);
551 Ok(string_vec
552 .into_iter()
553 .filter_map(|s| PublicKey::from_prefixed_string(&s).ok())
554 .map(PeerId::new)
555 .collect())
556 }
557
558 pub(super) async fn is_tree_synced_with_peer(
567 &self,
568 peer_pubkey: &PublicKey,
569 tree_root_id: &ID,
570 ) -> Result<bool> {
571 let pk_str = peer_pubkey.to_string();
572 let tree_root_str = tree_root_id.to_string();
573 let trees = self.txn.get_store::<DocStore>(TREES_SUBTREE).await?;
574 let peer_list_path = path!(&tree_root_str, "peer_pubkeys");
575 match trees.get_path_as::<String>(&peer_list_path).await {
576 Ok(peer_list_json) => {
577 if let Ok(peer_pubkeys) = serde_json::from_str::<Vec<String>>(&peer_list_json) {
578 Ok(peer_pubkeys.contains(&pk_str))
579 } else {
580 Ok(false)
581 }
582 }
583 Err(_) => Ok(false),
584 }
585 }
586
587 pub(super) async fn add_address(
598 &self,
599 peer_pubkey: &PublicKey,
600 address: Address,
601 ) -> Result<()> {
602 let pk_str = peer_pubkey.to_string();
603 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
604
605 if !peers.contains_path_str(&pk_str).await {
607 return Err(Error::Sync(Box::new(SyncError::PeerNotFound(pk_str))));
608 }
609
610 let addresses_result = peers
612 .get_path_as::<String>(path!(&pk_str, "addresses"))
613 .await;
614 let mut all_addresses: Vec<Address> = addresses_result
615 .ok()
616 .and_then(|json| serde_json::from_str(&json).ok())
617 .unwrap_or_else(Vec::new);
618
619 if !all_addresses.contains(&address) {
621 all_addresses.push(address);
622
623 let addresses_json = serde_json::to_string(&all_addresses).unwrap_or_default();
625 peers
626 .set_path(path!(&pk_str, "addresses"), addresses_json)
627 .await?;
628
629 let now = self.txn.now_rfc3339()?;
631 peers.set_path(path!(&pk_str, "last_seen"), now).await?;
632 }
633
634 Ok(())
635 }
636
637 pub(super) async fn remove_address(
646 &self,
647 peer_pubkey: &PublicKey,
648 address: &Address,
649 ) -> Result<bool> {
650 let pk_str = peer_pubkey.to_string();
651 let peers = self.txn.get_store::<DocStore>(PEERS_SUBTREE).await?;
652
653 if !peers.contains_path_str(&pk_str).await {
655 return Err(Error::Sync(Box::new(SyncError::PeerNotFound(pk_str))));
656 }
657
658 let addresses_result = peers
660 .get_path_as::<String>(path!(&pk_str, "addresses"))
661 .await;
662 let mut all_addresses: Vec<Address> = addresses_result
663 .ok()
664 .and_then(|json| serde_json::from_str(&json).ok())
665 .unwrap_or_else(Vec::new);
666
667 let initial_len = all_addresses.len();
669 all_addresses.retain(|a| a != address);
670
671 if all_addresses.len() != initial_len {
672 let addresses_json = serde_json::to_string(&all_addresses).unwrap_or_default();
674 peers
675 .set_path(path!(&pk_str, "addresses"), addresses_json)
676 .await?;
677
678 let now = self.txn.now_rfc3339()?;
680 peers.set_path(path!(&pk_str, "last_seen"), now).await?;
681
682 Ok(true)
683 } else {
684 Ok(false)
685 }
686 }
687
688 pub(super) async fn get_addresses(
697 &self,
698 peer_pubkey: &PublicKey,
699 transport_type: Option<&str>,
700 ) -> Result<Vec<Address>> {
701 if let Some(peer_info) = self.get_peer_info(peer_pubkey).await? {
702 match transport_type {
703 Some(transport) => Ok(peer_info
704 .get_addresses(transport)
705 .into_iter()
706 .cloned()
707 .collect()),
708 None => Ok(peer_info.addresses),
709 }
710 } else {
711 Ok(Vec::new())
712 }
713 }
714}