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
// use std::{thread, time::Duration};
// use tokio::spawn;
// pub enum Message {
// Start { label: String },
// Stop,
// GetTotalMessageCount(oneshot::Sender<usize>),
// }
// #[derive(Clone)]
// pub struct MessageSender(tokio::sync::mpsc::Sender<Message>);
// impl MessageSender {
// pub async fn start<T: Into<String>>(&self, label: T) {
// self.0
// .send(Message::Start {
// label: label.into(),
// })
// .await
// .unwrap();
// }
// pub async fn get_total_message_count(&self) -> usize {
// let (tx, rx) = oneshot::channel();
// self.0
// .send(Message::GetTotalMessageCount(tx))
// .await
// .unwrap();
// rx.await.unwrap()
// }
// pub async fn stop(&self) {
// self.0.send(Message::Stop).await.unwrap();
// }
// }
// trait MessageProtocol {
// fn start(label: String);
// fn stop();
// fn get_total_message_count() -> usize;
// }
// #[tokio::main]
// async fn main() {
// let (tx, mut rx) = tokio::sync::mpsc::channel::<Message>(100);
// let context = MessageSender(tx);
// let other_context = context.clone();
// spawn(async move {
// loop {
// dbg!(thread::current().name());
// context.start("with label").await;
// tokio::time::sleep(Duration::from_secs(1)).await;
// context.stop().await;
// tokio::time::sleep(Duration::from_secs(1)).await;
// println!("{}", context.get_total_message_count().await);
// }
// });
// spawn(async move {
// loop {
// dbg!(thread::current().name());
// other_context.start("with label 2").await;
// tokio::time::sleep(Duration::from_secs(1)).await;
// other_context.stop().await;
// tokio::time::sleep(Duration::from_secs(1)).await;
// println!("{}", other_context.get_total_message_count().await);
// }
// });
// let mut i = 0;
// while let Some(message) = rx.recv().await {
// match message {
// Message::Start { label } => println!("start: {label}"),
// Message::Stop => println!("stop"),
// Message::GetTotalMessageCount(sender) => {
// sender.send(i).unwrap();
// }
// }
// i += 1;
// }
// }