protocol_vote.rs 2.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566
  1. use async_executor::Executor;
  2. use async_trait::async_trait;
  3. use darkfi::{
  4. consensus::{state::StatePtr, vote::Vote},
  5. net::{
  6. ChannelPtr, MessageSubscription, ProtocolBase, ProtocolBasePtr, ProtocolJobsManager,
  7. ProtocolJobsManagerPtr,
  8. },
  9. Result,
  10. };
  11. use log::debug;
  12. use std::sync::Arc;
  13. pub struct ProtocolVote {
  14. vote_sub: MessageSubscription<Vote>,
  15. jobsman: ProtocolJobsManagerPtr,
  16. state: StatePtr,
  17. }
  18. impl ProtocolVote {
  19. pub async fn init(channel: ChannelPtr, state: StatePtr) -> ProtocolBasePtr {
  20. let message_subsytem = channel.get_message_subsystem();
  21. message_subsytem.add_dispatch::<Vote>().await;
  22. let vote_sub = channel.subscribe_msg::<Vote>().await.expect("Missing Vote dispatcher!");
  23. Arc::new(Self {
  24. vote_sub,
  25. jobsman: ProtocolJobsManager::new("VoteProtocol", channel),
  26. state,
  27. })
  28. }
  29. // TODO: 1. Nodes count not hardcoded.
  30. async fn handle_receive_vote(self: Arc<Self>) -> Result<()> {
  31. debug!(target: "ircd", "ProtocolVote::handle_receive_vote() [START]");
  32. loop {
  33. let vote = self.vote_sub.receive().await?;
  34. debug!(
  35. target: "ircd",
  36. "ProtocolVote::handle_receive_vote() received {:?}",
  37. vote
  38. );
  39. let vote_copy = (*vote).clone();
  40. let nodes_count = 4;
  41. self.state.write().unwrap().receive_vote(&vote_copy, nodes_count);
  42. }
  43. }
  44. }
  45. #[async_trait]
  46. impl ProtocolBase for ProtocolVote {
  47. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  48. debug!(target: "ircd", "ProtocolVote::start() [START]");
  49. self.jobsman.clone().start(executor.clone());
  50. self.jobsman.clone().spawn(self.clone().handle_receive_vote(), executor.clone()).await;
  51. debug!(target: "ircd", "ProtocolVote::start() [END]");
  52. Ok(())
  53. }
  54. fn name(&self) -> &'static str {
  55. "ProtocolVote"
  56. }
  57. }