mocra_core/queue/
channel.rs1use crate::common::model::message::TaskEvent;
2use crate::common::model::message::{TaskErrorEvent, TaskParserEvent};
3use crate::common::model::{Request, Response};
4use crate::queue::QueuedItem;
5use crate::utils::logger::LogModel;
6use std::sync::Arc;
7use tokio::sync::Mutex;
8use tokio::sync::mpsc::{Receiver, Sender, channel};
9
10pub struct Channel {
17 pub task_sender: Sender<QueuedItem<TaskEvent>>,
21 pub task_receiver: Arc<Mutex<Receiver<QueuedItem<TaskEvent>>>>,
22
23 pub request_sender: Sender<QueuedItem<Request>>,
27 pub request_receiver: Arc<Mutex<Receiver<QueuedItem<Request>>>>,
28
29 pub download_request_sender: Sender<QueuedItem<Request>>,
31 pub download_request_receiver: Arc<Mutex<Receiver<QueuedItem<Request>>>>,
32
33 pub response_sender: Sender<QueuedItem<Response>>,
38 pub response_receiver: Arc<Mutex<Receiver<QueuedItem<Response>>>>,
39
40 pub error_sender: Sender<QueuedItem<TaskErrorEvent>>,
44 pub error_receiver: Arc<Mutex<Receiver<QueuedItem<TaskErrorEvent>>>>,
45
46 pub log_sender: Sender<QueuedItem<LogModel>>,
48 pub log_receiver: Arc<Mutex<Receiver<QueuedItem<LogModel>>>>,
49
50 pub parser_task_sender: Sender<QueuedItem<TaskParserEvent>>,
54 pub parser_task_receiver: Arc<Mutex<Receiver<QueuedItem<TaskParserEvent>>>>,
55
56 pub remote_response_sender: Sender<QueuedItem<Response>>,
58 pub remote_response_receiver: Arc<Mutex<Receiver<QueuedItem<Response>>>>,
59
60 pub remote_parser_task_sender: Sender<QueuedItem<TaskParserEvent>>,
61 pub remote_parser_task_receiver: Arc<Mutex<Receiver<QueuedItem<TaskParserEvent>>>>,
62
63 pub remote_error_sender: Sender<QueuedItem<TaskErrorEvent>>,
64 pub remote_error_receiver: Arc<Mutex<Receiver<QueuedItem<TaskErrorEvent>>>>,
65
66 pub remote_task_sender: Sender<QueuedItem<TaskEvent>>,
68 pub remote_task_receiver: Arc<Mutex<Receiver<QueuedItem<TaskEvent>>>>,
69}
70impl Clone for Channel {
71 fn clone(&self) -> Self {
72 Channel {
73 task_sender: self.task_sender.clone(),
74 task_receiver: self.task_receiver.clone(),
75
76 request_sender: self.request_sender.clone(),
77 request_receiver: self.request_receiver.clone(),
78
79 download_request_sender: self.download_request_sender.clone(),
80 download_request_receiver: self.download_request_receiver.clone(),
81
82 response_sender: self.response_sender.clone(),
83 response_receiver: self.response_receiver.clone(),
84
85 error_sender: self.error_sender.clone(),
86 error_receiver: self.error_receiver.clone(),
87
88 log_sender: self.log_sender.clone(),
89 log_receiver: self.log_receiver.clone(),
90
91 parser_task_sender: self.parser_task_sender.clone(),
92 parser_task_receiver: self.parser_task_receiver.clone(),
93
94 remote_response_sender: self.remote_response_sender.clone(),
95 remote_response_receiver: self.remote_response_receiver.clone(),
96
97 remote_parser_task_sender: self.remote_parser_task_sender.clone(),
98 remote_parser_task_receiver: self.remote_parser_task_receiver.clone(),
99
100 remote_error_sender: self.remote_error_sender.clone(),
101 remote_error_receiver: self.remote_error_receiver.clone(),
102
103 remote_task_sender: self.remote_task_sender.clone(),
104 remote_task_receiver: self.remote_task_receiver.clone(),
105 }
106 }
107}
108
109impl Channel {
110 pub fn new(capacity: usize) -> Self {
111 let effective_capacity = if capacity < 10000 { 10000 } else { capacity };
115
116 let (task_sender, task_receiver) = channel(effective_capacity);
117 let (request_sender, request_receiver) = channel(effective_capacity);
118 let (download_request_sender, download_request_receiver) = channel(effective_capacity);
119 let (response_sender, response_receiver) = channel(effective_capacity);
120 let (error_sender, error_receiver) = channel(effective_capacity);
121 let (log_sender, log_receiver) = channel(10000); let (parser_task_sender, parser_task_receiver) = channel(effective_capacity);
124
125 let (remote_response_sender, remote_response_receiver) = channel(effective_capacity);
127 let (remote_parser_task_sender, remote_parser_task_receiver) = channel(effective_capacity);
128 let (remote_error_sender, remote_error_receiver) = channel(effective_capacity);
129 let (remote_task_sender, remote_task_receiver) = channel(effective_capacity);
130
131 Channel {
132 task_sender,
133 task_receiver: Arc::new(Mutex::new(task_receiver)),
134
135 request_sender,
136 request_receiver: Arc::new(Mutex::new(request_receiver)),
137
138 download_request_sender,
139 download_request_receiver: Arc::new(Mutex::new(download_request_receiver)),
140
141 response_sender,
142 response_receiver: Arc::new(Mutex::new(response_receiver)),
143
144 error_sender,
145 error_receiver: Arc::new(Mutex::new(error_receiver)),
146
147 log_sender,
148 log_receiver: Arc::new(Mutex::new(log_receiver)),
149
150 parser_task_sender,
151 parser_task_receiver: Arc::new(Mutex::new(parser_task_receiver)),
152
153 remote_response_sender,
154 remote_response_receiver: Arc::new(Mutex::new(remote_response_receiver)),
155
156 remote_parser_task_sender,
157 remote_parser_task_receiver: Arc::new(Mutex::new(remote_parser_task_receiver)),
158
159 remote_error_sender,
160 remote_error_receiver: Arc::new(Mutex::new(remote_error_receiver)),
161
162 remote_task_sender,
163 remote_task_receiver: Arc::new(Mutex::new(remote_task_receiver)),
164 }
165 }
166}