Skip to main content

mocra_core/queue/
channel.rs

1use 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
10/// Message queue channel manager.
11///
12/// Data flow:
13/// 1. Task Flow: watch TaskModel -> produce Request
14/// 2. Request Flow: watch Request -> download -> produce Response
15/// 3. Response Flow: watch Response -> parse -> finish, or produce a new Request
16pub struct Channel {
17    // --- 1. Initial task queue (Task Queue) ---
18    // Producer: external API / seed generator
19    // Consumer: task processor (TaskProcessor) -> produces Request
20    pub task_sender: Sender<QueuedItem<TaskEvent>>,
21    pub task_receiver: Arc<Mutex<Receiver<QueuedItem<TaskEvent>>>>,
22
23    // --- 2. Request queue (Request Queue) ---
24    // Producer: task processor / response processor (parse results)
25    // Consumer: downloader (Downloader) -> performs the download -> produces Response
26    pub request_sender: Sender<QueuedItem<Request>>,
27    pub request_receiver: Arc<Mutex<Receiver<QueuedItem<Request>>>>,
28
29    // (Optional) download-only queue, distinguishing requests before and after scheduling
30    pub download_request_sender: Sender<QueuedItem<Request>>,
31    pub download_request_receiver: Arc<Mutex<Receiver<QueuedItem<Request>>>>,
32
33    // --- 3. Response queue (Response Queue) ---
34    // Producer: downloader (Downloader)
35    // Consumer: response processor (ResponseProcessor) -> parses data -> (store / produce a new
36    // Request)
37    pub response_sender: Sender<QueuedItem<Response>>,
38    pub response_receiver: Arc<Mutex<Receiver<QueuedItem<Response>>>>,
39
40    // --- 4. Auxiliary Queues ---
41
42    // Error handling
43    pub error_sender: Sender<QueuedItem<TaskErrorEvent>>,
44    pub error_receiver: Arc<Mutex<Receiver<QueuedItem<TaskErrorEvent>>>>,
45
46    // Log handling
47    pub log_sender: Sender<QueuedItem<LogModel>>,
48    pub log_receiver: Arc<Mutex<Receiver<QueuedItem<LogModel>>>>,
49
50    // --- 5. Other / extension queues (Others) ---
51
52    // Parser task queue (when parsing is decoupled from response handling)
53    pub parser_task_sender: Sender<QueuedItem<TaskParserEvent>>,
54    pub parser_task_receiver: Arc<Mutex<Receiver<QueuedItem<TaskParserEvent>>>>,
55
56    // Remote / distributed node communication queues
57    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    // Distributed Task Queue (Outbound)
67    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        // Optimization: Use a much larger buffer in Single Node Mode to reduce backpressure.
112        // If "capacity" passed is 1000, it's too small for high-throughput single node.
113        // We override or multiply it here locally.
114        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); // Increased log channel
122
123        let (parser_task_sender, parser_task_receiver) = channel(effective_capacity);
124
125        // Remote channels
126        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}