1use std::collections::{HashMap, VecDeque};
6
7use net_traits::filemanager_thread::{FileManagerResult, ReadFileProgress};
8use rustc_hash::{FxBuildHasher, FxHashMap};
9use script_bindings::refcounted::Trusted;
10use script_bindings::root::Dom;
11use servo_base::id::{BroadcastChannelRouterId, MessagePortId, MessagePortRouterId};
12use servo_constellation_traits::{
13 BlobImpl, BroadcastChannelMsg, MessagePortImpl, MessagePortMsg, ScriptToConstellationMessage,
14};
15use uuid::Uuid;
16
17use crate::dom::GlobalScope;
18use crate::dom::bindings::error::{Error, Fallible};
19use crate::dom::bindings::reflector::DomGlobal;
20use crate::dom::bindings::str::DOMString;
21use crate::dom::bindings::trace::HashMapTracedValues;
22use crate::dom::bindings::weakref::WeakRef;
23use crate::dom::blob::Blob;
24use crate::dom::file::File;
25use crate::dom::globalscope::broadcastchannel::BroadcastChannel;
26use crate::dom::messageport::MessagePort;
27use crate::dom::promise::RootedPromise;
28use crate::dom::readablestream::ReadableStream;
29use crate::dom::transformstream::CrossRealmTransform;
30use crate::realms::enter_auto_realm;
31use crate::tasks::task_source::SendableTaskSource;
32use crate::test::TrustedPromise;
33
34pub(super) struct MessageListener {
36 pub(super) task_source: SendableTaskSource,
37 pub(super) context: Trusted<GlobalScope>,
38}
39
40pub(super) struct BroadcastListener {
42 pub(super) task_source: SendableTaskSource,
43 pub(super) context: Trusted<GlobalScope>,
44}
45
46pub(super) type FileListenerCallback =
47 Box<dyn Fn(&mut js::context::JSContext, &RootedPromise, Fallible<Vec<u8>>) + Send>;
48
49pub(super) struct FileListener {
51 pub(super) state: Option<FileListenerState>,
55 pub(super) task_source: SendableTaskSource,
56}
57
58pub(super) enum FileListenerTarget {
59 Promise(TrustedPromise, FileListenerCallback),
60 Stream(Trusted<ReadableStream>),
61}
62
63pub(super) enum FileListenerState {
64 Empty(FileListenerTarget),
65 Receiving(Vec<u8>, FileListenerTarget),
66}
67
68#[derive(JSTraceable, MallocSizeOf)]
69pub(crate) enum BlobTracker {
71 File(WeakRef<File>),
73 Blob(WeakRef<Blob>),
75}
76
77#[derive(JSTraceable, MallocSizeOf)]
78pub(crate) struct BlobInfo {
80 pub(super) tracker: BlobTracker,
82 #[no_trace]
84 pub(super) blob_impl: BlobImpl,
85 pub(super) has_url: bool,
88}
89
90pub(super) enum BlobResult {
94 Bytes(Vec<u8>),
95 File(Uuid, usize),
96}
97
98#[derive(JSTraceable, MallocSizeOf)]
100#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
101pub(crate) struct ManagedMessagePort {
102 pub(super) dom_port: Dom<MessagePort>,
104 #[no_trace]
109 pub(super) port_impl: Option<MessagePortImpl>,
110 pub(super) pending: bool,
114 pub(super) explicitly_closed: bool,
117 pub(super) cross_realm_transform: Option<CrossRealmTransform>,
120}
121
122#[derive(JSTraceable, MallocSizeOf)]
124#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
125pub(super) enum BroadcastChannelState {
126 Managed(
131 #[no_trace] BroadcastChannelRouterId,
132 HashMap<DOMString, VecDeque<Dom<BroadcastChannel>>>,
134 ),
135 UnManaged,
137}
138
139#[derive(JSTraceable, MallocSizeOf)]
141#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
142pub(super) enum MessagePortState {
143 Managed(
145 #[no_trace] MessagePortRouterId,
146 HashMapTracedValues<MessagePortId, ManagedMessagePort, FxBuildHasher>,
147 ),
148 UnManaged,
150}
151
152impl BroadcastListener {
153 pub(super) fn handle(&self, event: BroadcastChannelMsg) {
156 let context = self.context.clone();
157
158 self.task_source
164 .queue(task!(broadcast_message_event: move || {
165 let global = context.root();
166 global.broadcast_message_event(event, None);
169 }));
170 }
171}
172
173impl MessageListener {
174 pub(super) fn notify(&self, msg: MessagePortMsg) {
178 match msg {
179 MessagePortMsg::CompleteTransfer(ports) => {
180 let context = self.context.clone();
181 self.task_source.queue(
182 task!(process_complete_transfer: move |cx| {
183 let global = context.root();
184
185 let router_id = match global.port_router_id() {
186 Some(router_id) => router_id,
187 None => {
188 let _ = global.script_to_constellation_chan().send(
191 ScriptToConstellationMessage::MessagePortTransferResult(None, vec![], ports),
192 );
193 return;
194 }
195 };
196
197 let mut succeeded = vec![];
198 let mut failed = FxHashMap::default();
199
200 for (id, info) in ports.into_iter() {
201 if global.is_managing_port(&id) {
202 succeeded.push(id);
203 global.complete_port_transfer(
204 cx,
205 id,
206 info.port_message_queue,
207 info.disentangled,
208 );
209 } else {
210 failed.insert(id, info);
211 }
212 }
213 let _ = global.script_to_constellation_chan().send(
214 ScriptToConstellationMessage::MessagePortTransferResult(Some(router_id), succeeded, failed),
215 );
216 })
217 );
218 },
219 MessagePortMsg::CompletePendingTransfer(port_id, info) => {
220 let context = self.context.clone();
221 self.task_source.queue(task!(complete_pending: move |cx| {
222 let global = context.root();
223 global.complete_port_transfer(cx, port_id, info.port_message_queue, info.disentangled);
224 }));
225 },
226 MessagePortMsg::CompleteDisentanglement(port_id) => {
227 let context = self.context.clone();
228 self.task_source
229 .queue(task!(try_complete_disentanglement: move |cx| {
230 let global = context.root();
231 global.try_complete_disentanglement(cx, port_id);
232 }));
233 },
234 MessagePortMsg::NewTask(port_id, task) => {
235 let context = self.context.clone();
236 self.task_source.queue(task!(process_new_task: move |cx| {
237 let global = context.root();
238 global.route_task_to_port(cx, port_id, task);
239 }));
240 },
241 }
242 }
243}
244
245fn stream_handle_incoming(
247 cx: &mut js::context::JSContext,
248 stream: &ReadableStream,
249 bytes: Fallible<Vec<u8>>,
250) {
251 match bytes {
252 Ok(b) => {
253 stream.enqueue_native(cx, b);
254 },
255 Err(e) => {
256 stream.error_native(cx, e);
257 },
258 }
259}
260
261fn stream_handle_eof(cx: &mut js::context::JSContext, stream: &ReadableStream) {
263 stream.controller_close_native(cx);
264}
265
266impl FileListener {
267 pub(super) fn handle(&mut self, msg: FileManagerResult<ReadFileProgress>) {
268 match msg {
269 Ok(ReadFileProgress::Meta(blob_buf)) => match self.state.take() {
270 Some(FileListenerState::Empty(target)) => {
271 let bytes = if let FileListenerTarget::Stream(ref trusted_stream) = target {
272 let trusted = trusted_stream.clone();
273
274 let task = task!(enqueue_stream_chunk: move |cx| {
275 let stream = trusted.root();
276 stream_handle_incoming(cx, &stream, Ok(blob_buf.bytes));
277 });
278 self.task_source.queue(task);
279
280 Vec::with_capacity(0)
281 } else {
282 blob_buf.bytes
283 };
284
285 self.state = Some(FileListenerState::Receiving(bytes, target));
286 },
287 _ => panic!(
288 "Unexpected FileListenerState when receiving ReadFileProgress::Meta msg."
289 ),
290 },
291 Ok(ReadFileProgress::Partial(mut bytes_in)) => match self.state.take() {
292 Some(FileListenerState::Receiving(mut bytes, target)) => {
293 if let FileListenerTarget::Stream(ref trusted_stream) = target {
294 let trusted = trusted_stream.clone();
295
296 let task = task!(enqueue_stream_chunk: move |cx| {
297 let stream = trusted.root();
298 stream_handle_incoming(cx, &stream, Ok(bytes_in));
299 });
300
301 self.task_source.queue(task);
302 } else {
303 bytes.append(&mut bytes_in);
304 };
305
306 self.state = Some(FileListenerState::Receiving(bytes, target));
307 },
308 _ => panic!(
309 "Unexpected FileListenerState when receiving ReadFileProgress::Partial msg."
310 ),
311 },
312 Ok(ReadFileProgress::EOF) => match self.state.take() {
313 Some(FileListenerState::Receiving(bytes, target)) => match target {
314 FileListenerTarget::Promise(trusted_promise, callback) => {
315 let task = task!(resolve_promise: move |cx| {
316 let promise = trusted_promise.root(cx);
317 let mut realm = enter_auto_realm(cx, &*promise.global());
318 callback(&mut realm, &promise, Ok(bytes));
319 });
320
321 self.task_source.queue(task);
322 },
323 FileListenerTarget::Stream(trusted_stream) => {
324 let task = task!(enqueue_stream_chunk: move |cx| {
325 let stream = trusted_stream.root();
326 stream_handle_eof(cx, &stream);
327 });
328
329 self.task_source.queue(task);
330 },
331 },
332 _ => {
333 panic!("Unexpected FileListenerState when receiving ReadFileProgress::EOF msg.")
334 },
335 },
336 Err(_) => match self.state.take() {
337 Some(FileListenerState::Receiving(_, target)) |
338 Some(FileListenerState::Empty(target)) => {
339 let error = Err(Error::Network(Some("No more file to receive.".into())));
340
341 match target {
342 FileListenerTarget::Promise(trusted_promise, callback) => {
343 self.task_source.queue(task!(reject_promise: move |cx| {
344 let promise = trusted_promise.root(cx);
345 let mut realm = enter_auto_realm(cx, &*promise.global());
346 callback(&mut realm, &promise, error);
347 }));
348 },
349 FileListenerTarget::Stream(trusted_stream) => {
350 self.task_source.queue(task!(error_stream: move |cx| {
351 let stream = trusted_stream.root();
352 stream_handle_incoming(cx, &stream, error);
353 }));
354 },
355 }
356 },
357 _ => panic!("Unexpected FileListenerState when receiving Err msg."),
358 },
359 }
360 }
361}