Skip to main content

script/dom/globalscope/
listeners.rs

1/* This Source Code Form is subject to the terms of the Mozilla Public
2 * License, v. 2.0. If a copy of the MPL was not distributed with this
3 * file, You can obtain one at https://mozilla.org/MPL/2.0/. */
4
5use 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
34/// A wrapper for glue-code between the ipc router and the event-loop.
35pub(super) struct MessageListener {
36    pub(super) task_source: SendableTaskSource,
37    pub(super) context: Trusted<GlobalScope>,
38}
39
40/// A wrapper for broadcasts coming in over IPC, and the event-loop.
41pub(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
49/// A wrapper for the handling of file data received by the ipc router
50pub(super) struct FileListener {
51    /// State should progress as either of:
52    /// - Some(Empty) => Some(Receiving) => None
53    /// - Some(Empty) => None
54    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)]
69/// A holder of a weak reference for a DOM blob or file.
70pub(crate) enum BlobTracker {
71    /// A weak ref to a DOM file.
72    File(WeakRef<File>),
73    /// A weak ref to a DOM blob.
74    Blob(WeakRef<Blob>),
75}
76
77#[derive(JSTraceable, MallocSizeOf)]
78/// The info pertaining to a blob managed by this global.
79pub(crate) struct BlobInfo {
80    /// The weak ref to the corresponding DOM object.
81    pub(super) tracker: BlobTracker,
82    /// The data and logic backing the DOM object.
83    #[no_trace]
84    pub(super) blob_impl: BlobImpl,
85    /// Whether this blob has an outstanding URL,
86    /// <https://w3c.github.io/FileAPI/#url>.
87    pub(super) has_url: bool,
88}
89
90/// The result of looking-up the data for a Blob,
91/// containing either the in-memory bytes,
92/// or the file-id.
93pub(super) enum BlobResult {
94    Bytes(Vec<u8>),
95    File(Uuid, usize),
96}
97
98/// Data representing a message-port managed by this global.
99#[derive(JSTraceable, MallocSizeOf)]
100#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
101pub(crate) struct ManagedMessagePort {
102    /// The DOM port.
103    pub(super) dom_port: Dom<MessagePort>,
104    /// The logic and data backing the DOM port.
105    /// The option is needed to take out the port-impl
106    /// as part of its transferring steps,
107    /// without having to worry about rooting the dom-port.
108    #[no_trace]
109    pub(super) port_impl: Option<MessagePortImpl>,
110    /// We keep ports pending when they are first transfer-received,
111    /// and only add them, and ask the constellation to complete the transfer,
112    /// in a subsequent task if the port hasn't been re-transfered.
113    pub(super) pending: bool,
114    /// Whether the port has been closed by script in this global,
115    /// so it can be removed.
116    pub(super) explicitly_closed: bool,
117    /// The handler for `message` or `messageerror` used in the cross realm transform,
118    /// if any was setup with this port.
119    pub(super) cross_realm_transform: Option<CrossRealmTransform>,
120}
121
122/// State representing whether this global is currently managing broadcast channels.
123#[derive(JSTraceable, MallocSizeOf)]
124#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
125pub(super) enum BroadcastChannelState {
126    /// The broadcast-channel router id for this global, and a queue of managed channels.
127    /// Step 9, "sort destinations"
128    /// of <https://html.spec.whatwg.org/multipage/#dom-broadcastchannel-postmessage>
129    /// requires keeping track of creation order, hence the queue.
130    Managed(
131        #[no_trace] BroadcastChannelRouterId,
132        /// The map of channel-name to queue of channels, in order of creation.
133        HashMap<DOMString, VecDeque<Dom<BroadcastChannel>>>,
134    ),
135    /// This global is not managing any broadcast channels at this time.
136    UnManaged,
137}
138
139/// State representing whether this global is currently managing messageports.
140#[derive(JSTraceable, MallocSizeOf)]
141#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
142pub(super) enum MessagePortState {
143    /// The message-port router id for this global, and a map of managed ports.
144    Managed(
145        #[no_trace] MessagePortRouterId,
146        HashMapTracedValues<MessagePortId, ManagedMessagePort, FxBuildHasher>,
147    ),
148    /// This global is not managing any ports at this time.
149    UnManaged,
150}
151
152impl BroadcastListener {
153    /// Handle a broadcast coming in over IPC,
154    /// by queueing the appropriate task on the relevant event-loop.
155    pub(super) fn handle(&self, event: BroadcastChannelMsg) {
156        let context = self.context.clone();
157
158        // Note: strictly speaking we should just queue the message event tasks,
159        // not queue a task that then queues more tasks.
160        // This however seems to be hard to avoid in the light of the IPC.
161        // One can imagine queueing tasks directly,
162        // for channels that would be in the same script-thread.
163        self.task_source
164            .queue(task!(broadcast_message_event: move || {
165                let global = context.root();
166                // Step 10 of https://html.spec.whatwg.org/multipage/#dom-broadcastchannel-postmessage,
167                // For each BroadcastChannel object destination in destinations, queue a task.
168                global.broadcast_message_event(event, None);
169            }));
170    }
171}
172
173impl MessageListener {
174    /// A new message came in, handle it via a task enqueued on the event-loop.
175    /// A task is required, since we are using a trusted globalscope,
176    /// and we can only access the root from the event-loop.
177    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                                // If not managing any ports, no transfer can succeed,
189                                // so just send back everything.
190                                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
245/// Callback used to enqueue file chunks to streams as part of FileListener.
246fn 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
261/// Callback used to close streams as part of FileListener.
262fn 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}