script/tasks/
task_queue.rs1use std::cell::Cell;
8use std::collections::VecDeque;
9use std::default::Default;
10
11use crossbeam_channel::{self, Receiver, Sender};
12use rustc_hash::{FxHashMap, FxHashSet};
13use script_bindings::cell::DomRefCell;
14use servo_base::id::PipelineId;
15use strum::VariantArray;
16
17use crate::dom::worker::TrustedWorkerAddress;
18use crate::runtime::script_runtime::ScriptThreadEventCategory;
19use crate::tasks::task::TaskBox;
20use crate::tasks::task_source::TaskSourceName;
21
22#[derive(MallocSizeOf)]
23pub(crate) struct QueuedTask {
24 pub(crate) worker: Option<TrustedWorkerAddress>,
25 pub(crate) event_category: ScriptThreadEventCategory,
26 #[ignore_malloc_size_of = "TaskBox is difficult"]
27 pub(crate) task: Box<dyn TaskBox>,
28 pub(crate) pipeline_id: Option<PipelineId>,
29 pub(crate) task_source: TaskSourceName,
30}
31
32pub(crate) trait QueuedTaskConversion {
34 fn task_source_name(&self) -> Option<&TaskSourceName>;
35 fn pipeline_id(&self) -> Option<PipelineId>;
36 fn into_queued_task(self) -> Option<QueuedTask>;
37 fn from_queued_task(queued_task: QueuedTask) -> Self;
38 fn inactive_msg() -> Self;
39 fn wake_up_msg() -> Self;
40 fn is_wake_up(&self) -> bool;
41}
42
43#[derive(MallocSizeOf)]
44pub(crate) struct TaskQueue<T> {
45 port: Receiver<T>,
47 wake_up_sender: Sender<T>,
49 msg_queue: DomRefCell<VecDeque<T>>,
51 taken_task_counter: Cell<u64>,
53 throttled: DomRefCell<FxHashMap<TaskSourceName, VecDeque<QueuedTask>>>,
55 inactive: DomRefCell<FxHashMap<PipelineId, VecDeque<QueuedTask>>>,
57}
58
59impl<T: QueuedTaskConversion> TaskQueue<T> {
60 pub(crate) fn new(port: Receiver<T>, wake_up_sender: Sender<T>) -> TaskQueue<T> {
61 TaskQueue {
62 port,
63 wake_up_sender,
64 msg_queue: DomRefCell::new(VecDeque::new()),
65 taken_task_counter: Default::default(),
66 throttled: Default::default(),
67 inactive: Default::default(),
68 }
69 }
70
71 pub(crate) fn remove_tasks_for_exiting_pipeline(&self, pipeline_id: &PipelineId) {
75 self.inactive.borrow_mut().remove(pipeline_id);
76 }
77
78 fn release_tasks_for_fully_active_documents(
81 &self,
82 fully_active: &FxHashSet<PipelineId>,
83 ) -> Vec<T> {
84 self.inactive
85 .borrow_mut()
86 .iter_mut()
87 .filter(|(pipeline_id, _)| fully_active.contains(pipeline_id))
88 .flat_map(|(_, inactive_queue)| {
89 inactive_queue
90 .drain(0..)
91 .map(|queued_task| T::from_queued_task(queued_task))
92 })
93 .collect()
94 }
95
96 fn store_task_for_inactive_pipeline(&self, msg: T, pipeline_id: &PipelineId) {
99 let mut inactive = self.inactive.borrow_mut();
100 let inactive_queue = inactive.entry(*pipeline_id).or_default();
101 inactive_queue.push_back(
102 msg.into_queued_task()
103 .expect("Incoming messages should always be convertible into queued tasks"),
104 );
105 let mut msg_queue = self.msg_queue.borrow_mut();
106 if msg_queue.is_empty() {
107 msg_queue.push_back(T::inactive_msg());
112 }
113 }
114
115 fn process_incoming_tasks(&self, first_msg: T, fully_active: &FxHashSet<PipelineId>) {
118 let mut incoming = self.release_tasks_for_fully_active_documents(fully_active);
120
121 if !first_msg.is_wake_up() {
123 incoming.push(first_msg);
124 }
125
126 while let Ok(msg) = self.port.try_recv() {
128 if !msg.is_wake_up() {
129 incoming.push(msg);
130 }
131 }
132
133 let mut to_be_throttled = Vec::new();
136 let mut index = 0;
137 while index != incoming.len() {
138 index += 1; let task_source = match incoming[index - 1].task_source_name() {
141 Some(task_source) => task_source,
142 None => continue,
143 };
144
145 match task_source {
146 TaskSourceName::PerformanceTimeline => {
147 to_be_throttled.push(incoming.remove(index - 1));
148 index -= 1; },
150 _ => {
151 self.taken_task_counter
153 .set(self.taken_task_counter.get() + 1);
154 },
155 }
156 }
157
158 for msg in incoming {
159 if let Some(TaskSourceName::Rendering) = msg.task_source_name() {
162 self.msg_queue.borrow_mut().push_back(msg);
163 continue;
164 }
165 if let Some(pipeline_id) = msg.pipeline_id() &&
166 !fully_active.contains(&pipeline_id)
167 {
168 self.store_task_for_inactive_pipeline(msg, &pipeline_id);
169 continue;
170 }
171 self.msg_queue.borrow_mut().push_back(msg);
173 }
174
175 for msg in to_be_throttled {
176 let Some(queued_task) = msg.into_queued_task() else {
178 unreachable!(
179 "A message to be throttled should always be convertible into a queued task"
180 );
181 };
182 let mut throttled_tasks = self.throttled.borrow_mut();
183 throttled_tasks
184 .entry(queued_task.task_source)
185 .or_default()
186 .push_back(queued_task);
187 }
188 }
189
190 pub(crate) fn select(&self) -> &crossbeam_channel::Receiver<T> {
193 self.taken_task_counter.set(0);
195 &self.port
198 }
199
200 pub(crate) fn recv(&self) -> Result<T, ()> {
202 self.msg_queue.borrow_mut().pop_front().ok_or(())
203 }
204
205 pub(crate) fn take_tasks_and_recv(
207 &self,
208 fully_active: &FxHashSet<PipelineId>,
209 ) -> Result<T, ()> {
210 self.take_tasks(T::wake_up_msg(), fully_active);
211 self.recv()
212 }
213
214 pub(crate) fn take_tasks(&self, first_msg: T, fully_active: &FxHashSet<PipelineId>) {
217 const PER_ITERATION_MAX: u64 = 5;
219 self.process_incoming_tasks(first_msg, fully_active);
221 let mut throttled = self.throttled.borrow_mut();
222 let mut throttled_length: usize = throttled.values().map(|queue| queue.len()).sum();
223 let mut task_source_cycler = TaskSourceName::VARIANTS.iter().cycle();
224 loop {
227 let max_reached = self.taken_task_counter.get() > PER_ITERATION_MAX;
228 let none_left = throttled_length == 0;
229 match (max_reached, none_left) {
230 (_, true) => break,
231 (true, false) => {
232 let _ = self.wake_up_sender.send(T::wake_up_msg());
236 break;
237 },
238 (false, false) => {
239 let task_source = task_source_cycler.next().unwrap();
241 let throttled_queue = match throttled.get_mut(task_source) {
242 Some(queue) => queue,
243 None => continue,
244 };
245 let queued_task = match throttled_queue.pop_front() {
246 Some(queued_task) => queued_task,
247 None => continue,
248 };
249 let msg = T::from_queued_task(queued_task);
250
251 if let Some(pipeline_id) = msg.pipeline_id() &&
253 !fully_active.contains(&pipeline_id)
254 {
255 self.store_task_for_inactive_pipeline(msg, &pipeline_id);
256 throttled_length -= 1;
260 continue;
261 }
262
263 self.msg_queue.borrow_mut().push_back(msg);
265 self.taken_task_counter
266 .set(self.taken_task_counter.get() + 1);
267 throttled_length -= 1;
268 },
269 }
270 }
271 }
272}