1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
use crate::protocol::FloodsubMessage;
use crate::subscription::Subscription;
use crate::Topic;
use crate::{FloodsubConfig, FloodsubError};
use futures::channel::{mpsc, oneshot};
use futures::SinkExt;
use libp2prs_core::PeerId;
pub(crate) enum ControlCommand {
Publish(FloodsubMessage, oneshot::Sender<()>),
Subscribe(Topic, oneshot::Sender<Subscription>),
Ls(oneshot::Sender<Vec<Topic>>),
GetPeers(Topic, oneshot::Sender<Vec<PeerId>>),
}
#[derive(Clone)]
pub struct Control {
config: FloodsubConfig,
control_sender: mpsc::UnboundedSender<ControlCommand>,
}
type Result<T> = std::result::Result<T, FloodsubError>;
impl Control {
pub(crate) fn new(control_sender: mpsc::UnboundedSender<ControlCommand>, config: FloodsubConfig) -> Self {
Control { control_sender, config }
}
pub fn close(&mut self) {
self.control_sender.close_channel();
}
pub async fn publish(&mut self, topic: Topic, data: impl Into<Vec<u8>>) -> Result<()> {
let msg = FloodsubMessage {
source: self.config.local_peer_id,
data: data.into(),
sequence_number: rand::random::<[u8; 20]>().to_vec(),
topics: vec![topic.clone()],
};
let (tx, rx) = oneshot::channel();
self.control_sender.send(ControlCommand::Publish(msg, tx)).await?;
Ok(rx.await?)
}
pub async fn subscribe(&mut self, topic: Topic) -> Result<Subscription> {
let (tx, rx) = oneshot::channel();
self.control_sender.send(ControlCommand::Subscribe(topic, tx)).await?;
Ok(rx.await?)
}
pub async fn ls(&mut self) -> Result<Vec<Topic>> {
let (tx, rx) = oneshot::channel();
self.control_sender.send(ControlCommand::Ls(tx)).await?;
Ok(rx.await?)
}
pub async fn get_peers(&mut self, topic: Topic) -> Result<Vec<PeerId>> {
let (tx, rx) = oneshot::channel();
self.control_sender.send(ControlCommand::GetPeers(topic, tx)).await?;
Ok(rx.await?)
}
}