consensus_sync.rs 1.5 KB

1234567891011121314151617181920212223242526272829303132333435363738394041
  1. use log::{info, warn};
  2. use crate::{
  3. consensus::{
  4. state::{ConsensusRequest, ConsensusResponse},
  5. ValidatorStatePtr,
  6. },
  7. net::P2pPtr,
  8. Result,
  9. };
  10. /// async task used for consensus state syncing.
  11. pub async fn consensus_sync_task(p2p: P2pPtr, state: ValidatorStatePtr) -> Result<()> {
  12. info!("Starting consensus state sync...");
  13. // Using len here beacuse is_empty() uses unstable library feature
  14. // called 'exact_size_is_empty'.
  15. if p2p.channels().lock().await.values().len() != 0 {
  16. // Nodes ask for the consensus state of the last channel peer
  17. let channel = p2p.channels().lock().await.values().last().unwrap().clone();
  18. // Communication setup
  19. let msg_subsystem = channel.get_message_subsystem();
  20. msg_subsystem.add_dispatch::<ConsensusResponse>().await;
  21. let response_sub = channel.subscribe_msg::<ConsensusResponse>().await?;
  22. // Node creates a `ConsensusRequest` and sends it
  23. let request = ConsensusRequest { address: state.read().await.address };
  24. channel.send(request).await?;
  25. // Node stores response data. Extra validations can be added here.
  26. let response = response_sub.receive().await?;
  27. state.write().await.consensus = response.consensus.clone();
  28. } else {
  29. warn!("Node is not connected to other nodes, resetting consensus state.");
  30. state.write().await.reset_consensus_state()?;
  31. }
  32. info!("Consensus state synced!");
  33. Ok(())
  34. }