protocol_proposal.rs 2.9 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788
  1. use async_executor::Executor;
  2. use async_trait::async_trait;
  3. use darkfi::{
  4. consensus::{block::BlockProposal, state::ValidatorStatePtr},
  5. net::{
  6. ChannelPtr, MessageSubscription, P2pPtr, ProtocolBase, ProtocolBasePtr,
  7. ProtocolJobsManager, ProtocolJobsManagerPtr,
  8. },
  9. Result,
  10. };
  11. use log::debug;
  12. use std::sync::Arc;
  13. pub struct ProtocolProposal {
  14. proposal_sub: MessageSubscription<BlockProposal>,
  15. jobsman: ProtocolJobsManagerPtr,
  16. state: ValidatorStatePtr,
  17. p2p: P2pPtr,
  18. }
  19. impl ProtocolProposal {
  20. pub async fn init(
  21. channel: ChannelPtr,
  22. state: ValidatorStatePtr,
  23. p2p: P2pPtr,
  24. ) -> ProtocolBasePtr {
  25. let message_subsytem = channel.get_message_subsystem();
  26. message_subsytem.add_dispatch::<BlockProposal>().await;
  27. let proposal_sub =
  28. channel.subscribe_msg::<BlockProposal>().await.expect("Missing Proposal dispatcher!");
  29. Arc::new(Self {
  30. proposal_sub,
  31. jobsman: ProtocolJobsManager::new("ProposalProtocol", channel),
  32. state,
  33. p2p,
  34. })
  35. }
  36. async fn handle_receive_proposal(self: Arc<Self>) -> Result<()> {
  37. debug!(target: "ircd", "ProtocolBlock::handle_receive_proposal() [START]");
  38. loop {
  39. let proposal = self.proposal_sub.receive().await?;
  40. debug!(
  41. target: "ircd",
  42. "ProtocolProposal::handle_receive_proposal() received {:?}",
  43. proposal
  44. );
  45. let proposal_copy = (*proposal).clone();
  46. let vote = self.state.write().unwrap().receive_proposal(&proposal_copy);
  47. match vote {
  48. Ok(x) => {
  49. if x.is_none() {
  50. debug!("Node did not vote for the proposed block.");
  51. } else {
  52. let vote = x.unwrap();
  53. self.state.write().unwrap().receive_vote(&vote)?;
  54. // Broadcasting block to rest nodes
  55. self.p2p.broadcast(proposal_copy).await?;
  56. // Broadcasting vote
  57. self.p2p.broadcast(vote).await?;
  58. }
  59. }
  60. Err(e) => {
  61. debug!(target: "ircd", "ProtocolBlock::handle_receive_proposal() error prosessing proposal: {:?}", e)
  62. }
  63. }
  64. }
  65. }
  66. }
  67. #[async_trait]
  68. impl ProtocolBase for ProtocolProposal {
  69. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  70. debug!(target: "ircd", "ProtocolProposal::start() [START]");
  71. self.jobsman.clone().start(executor.clone());
  72. self.jobsman.clone().spawn(self.clone().handle_receive_proposal(), executor.clone()).await;
  73. debug!(target: "ircd", "ProtocolProposal::start() [END]");
  74. Ok(())
  75. }
  76. fn name(&self) -> &'static str {
  77. "ProtocolProposal"
  78. }
  79. }