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
use super::connection::Connection;
use crate::command::Command;
use futures::StreamExt;
impl Connection {
/// Run the connection task.
pub(crate) async fn run(mut self) {
// Get the next command.
while let Some(cmd) = self.command_receiver.next().await {
match cmd {
Command::SendMessage(msg) => self.send_message(msg),
Command::SendMessageOneshot(msg, response) => {
self.send_message_oneshot(msg, response)
}
Command::SendMessageMpcs(msg, response_reply_serial, response) => {
self.send_message_mpsc(msg, response_reply_serial, response)
}
Command::AddPath(path, object) => {
// Add the handler.
self.path_handler.insert(path, object);
}
Command::DeletePath(path) => {
// Remove the handler.
self.path_handler.remove(&path);
}
Command::DeleteSender(sender_other) => {
// Remove the handler by `Sender<Message>` object.
self.path_handler
.retain(|_path, sender| !sender_other.same_receiver(sender));
}
Command::DeleteReceiver(_receiver_other) => {
// TODO: Wait until the is_connect PR is merged:
// https://github.com/rust-lang/futures-rs/pull/2179
}
Command::ListPath(path, sender) => self.list_path(&path, sender),
Command::AddInterface(interface, sender) => {
// Add an interface handler
self.interface_handler.insert(interface, sender);
}
Command::AddSignalHandler(path, sender) => {
// Add a signal handler.
if let Some(vec) = self.signals.get_mut(&path) {
vec.push(sender);
} else {
self.signals.insert(path, vec![sender]);
}
}
Command::DeleteSignalHandler(sender_other) => {
// Remove the signal handler by `Sender<Message>` object.
for vec_sender_message in self.signals.values_mut() {
vec_sender_message.retain(|sender| !sender_other.same_receiver(sender));
}
}
Command::ReceiveMessage(msg) => self.receive_message(msg),
Command::Close => {
// Stop the server.
return;
}
}
}
}
}