protocol_sync_consensus.rs 2.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172
  1. use async_executor::Executor;
  2. use async_trait::async_trait;
  3. use darkfi::{
  4. consensus::state::{ConsensusRequest, ConsensusResponse, ValidatorStatePtr},
  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 ProtocolSyncConsensus {
  14. channel: ChannelPtr,
  15. request_sub: MessageSubscription<ConsensusRequest>,
  16. jobsman: ProtocolJobsManagerPtr,
  17. state: ValidatorStatePtr,
  18. }
  19. impl ProtocolSyncConsensus {
  20. pub async fn init(channel: ChannelPtr, state: ValidatorStatePtr) -> ProtocolBasePtr {
  21. let message_subsytem = channel.get_message_subsystem();
  22. message_subsytem.add_dispatch::<ConsensusRequest>().await;
  23. let request_sub = channel
  24. .subscribe_msg::<ConsensusRequest>()
  25. .await
  26. .expect("Missing ConsensusRequest dispatcher!");
  27. Arc::new(Self {
  28. channel: channel.clone(),
  29. request_sub,
  30. jobsman: ProtocolJobsManager::new("SyncConsensusProtocol", channel),
  31. state,
  32. })
  33. }
  34. async fn handle_receive_request(self: Arc<Self>) -> Result<()> {
  35. debug!(target: "ircd", "ProtocolSyncConsensus::handle_receive_request() [START]");
  36. loop {
  37. let order = self.request_sub.receive().await?;
  38. debug!(
  39. target: "ircd",
  40. "ProtocolSyncConsensus::handle_receive_request() received {:?}",
  41. order
  42. );
  43. // Extra validations can be added here.
  44. let consensus = self.state.read().unwrap().consensus.clone();
  45. let response = ConsensusResponse { consensus };
  46. self.channel.send(response).await?;
  47. }
  48. }
  49. }
  50. #[async_trait]
  51. impl ProtocolBase for ProtocolSyncConsensus {
  52. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  53. debug!(target: "ircd", "ProtocolSyncConsensus::start() [START]");
  54. self.jobsman.clone().start(executor.clone());
  55. self.jobsman.clone().spawn(self.clone().handle_receive_request(), executor.clone()).await;
  56. debug!(target: "ircd", "ProtocolSyncConsensus::start() [END]");
  57. Ok(())
  58. }
  59. fn name(&self) -> &'static str {
  60. "ProtocolSyncConsensus"
  61. }
  62. }