protocol_sync.rs 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119
  1. use async_executor::Executor;
  2. use async_trait::async_trait;
  3. use darkfi::{
  4. consensus::{
  5. block::{BlockInfo, BlockOrder, BlockResponse},
  6. state::ValidatorStatePtr,
  7. },
  8. net::{
  9. ChannelPtr, MessageSubscription, P2pPtr, ProtocolBase, ProtocolBasePtr,
  10. ProtocolJobsManager, ProtocolJobsManagerPtr,
  11. },
  12. Result,
  13. };
  14. use log::debug;
  15. use std::sync::Arc;
  16. // Constant defining how many blocks we send during syncing.
  17. const BATCH: u64 = 10;
  18. pub struct ProtocolSync {
  19. channel: ChannelPtr,
  20. request_sub: MessageSubscription<BlockOrder>,
  21. block_sub: MessageSubscription<BlockInfo>,
  22. jobsman: ProtocolJobsManagerPtr,
  23. state: ValidatorStatePtr,
  24. p2p: P2pPtr,
  25. consensus_mode: bool,
  26. }
  27. impl ProtocolSync {
  28. pub async fn init(
  29. channel: ChannelPtr,
  30. state: ValidatorStatePtr,
  31. p2p: P2pPtr,
  32. consensus_mode: bool,
  33. ) -> ProtocolBasePtr {
  34. let message_subsytem = channel.get_message_subsystem();
  35. message_subsytem.add_dispatch::<BlockOrder>().await;
  36. message_subsytem.add_dispatch::<BlockInfo>().await;
  37. let request_sub =
  38. channel.subscribe_msg::<BlockOrder>().await.expect("Missing BlockOrder dispatcher!");
  39. let block_sub =
  40. channel.subscribe_msg::<BlockInfo>().await.expect("Missing BlockInfo dispatcher!");
  41. Arc::new(Self {
  42. channel: channel.clone(),
  43. request_sub,
  44. block_sub,
  45. jobsman: ProtocolJobsManager::new("SyncProtocol", channel),
  46. state,
  47. p2p,
  48. consensus_mode,
  49. })
  50. }
  51. async fn handle_receive_request(self: Arc<Self>) -> Result<()> {
  52. debug!(target: "ircd", "ProtocolSync::handle_receive_request() [START]");
  53. loop {
  54. let order = self.request_sub.receive().await?;
  55. debug!(
  56. target: "ircd",
  57. "ProtocolSync::handle_receive_request() received {:?}",
  58. order
  59. );
  60. // Extra validations can be added here.
  61. let key = order.sl;
  62. let blocks = self.state.read().unwrap().blockchain.get_with_info(key, BATCH)?;
  63. let response = BlockResponse { blocks };
  64. self.channel.send(response).await?;
  65. }
  66. }
  67. async fn handle_receive_block(self: Arc<Self>) -> Result<()> {
  68. debug!(target: "ircd", "ProtocolSync::handle_receive_block() [START]");
  69. loop {
  70. let info = self.block_sub.receive().await?;
  71. debug!(
  72. target: "ircd",
  73. "ProtocolSync::handle_receive_block() received {:?}",
  74. info
  75. );
  76. // Node stores finalized block, if it doesn't exists (checking by slot),
  77. // and removes its transactions from the unconfirmed_txs vector.
  78. // Consensus mode enabled nodes have already performed this steps,
  79. // during proposal finalization.
  80. // Extra validations can be added here.
  81. if !self.consensus_mode {
  82. let info_copy = (*info).clone();
  83. if !self.state.read().unwrap().blockchain.has_block(&info_copy)? {
  84. self.state.write().unwrap().blockchain.add_by_info(info_copy.clone())?;
  85. self.state.write().unwrap().remove_txs(info_copy.txs.clone())?;
  86. self.p2p.broadcast(info_copy).await?;
  87. }
  88. }
  89. }
  90. }
  91. }
  92. #[async_trait]
  93. impl ProtocolBase for ProtocolSync {
  94. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  95. debug!(target: "ircd", "ProtocolSync::start() [START]");
  96. self.jobsman.clone().start(executor.clone());
  97. self.jobsman.clone().spawn(self.clone().handle_receive_request(), executor.clone()).await;
  98. self.jobsman.clone().spawn(self.clone().handle_receive_block(), executor.clone()).await;
  99. debug!(target: "ircd", "ProtocolSync::start() [END]");
  100. Ok(())
  101. }
  102. fn name(&self) -> &'static str {
  103. "ProtocolSync"
  104. }
  105. }