protocol_participant.rs 2.3 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677
  1. use async_executor::Executor;
  2. use async_trait::async_trait;
  3. use darkfi::{
  4. consensus::{participant::Participant, 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 ProtocolParticipant {
  14. participant_sub: MessageSubscription<Participant>,
  15. jobsman: ProtocolJobsManagerPtr,
  16. state: ValidatorStatePtr,
  17. p2p: P2pPtr,
  18. }
  19. impl ProtocolParticipant {
  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::<Participant>().await;
  27. let participant_sub =
  28. channel.subscribe_msg::<Participant>().await.expect("Missing Participant dispatcher!");
  29. Arc::new(Self {
  30. participant_sub,
  31. jobsman: ProtocolJobsManager::new("ParticipantProtocol", channel),
  32. state,
  33. p2p,
  34. })
  35. }
  36. async fn handle_receive_participant(self: Arc<Self>) -> Result<()> {
  37. debug!(target: "ircd", "ProtocolParticipant::handle_receive_participant() [START]");
  38. loop {
  39. let participant = self.participant_sub.receive().await?;
  40. debug!(
  41. target: "ircd",
  42. "ProtocolParticipant::handle_receive_participant() received {:?}",
  43. participant
  44. );
  45. let participant_copy = (*participant).clone();
  46. if self.state.write().unwrap().append_participant(participant_copy.clone()) {
  47. self.p2p.broadcast(participant_copy).await?;
  48. }
  49. }
  50. }
  51. }
  52. #[async_trait]
  53. impl ProtocolBase for ProtocolParticipant {
  54. async fn start(self: Arc<Self>, executor: Arc<Executor<'_>>) -> Result<()> {
  55. debug!(target: "ircd", "ProtocolParticipant::start() [START]");
  56. self.jobsman.clone().start(executor.clone());
  57. self.jobsman
  58. .clone()
  59. .spawn(self.clone().handle_receive_participant(), executor.clone())
  60. .await;
  61. debug!(target: "ircd", "ProtocolParticipant::start() [END]");
  62. Ok(())
  63. }
  64. fn name(&self) -> &'static str {
  65. "ProtocolParticipant"
  66. }
  67. }