Skip to main content

script/dom/stream/
byteteereadintorequest.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 http://mozilla.org/MPL/2.0/. */
4
5#![cfg_attr(crown, allow(crown::jscontext_first_arg))]
6
7use std::cell::Cell;
8use std::rc::Rc;
9
10use dom_struct::dom_struct;
11use js::context::JSContext;
12use js::jsval::UndefinedValue;
13use js::typedarray::ArrayBufferViewU8;
14use script_bindings::reflector::{Reflector, reflect_dom_object};
15use script_bindings::trace::RootedTraceableBox;
16
17use super::byteteeunderlyingsource::ByteTeePullAlgorithm;
18use crate::dom::bindings::buffer_source::HeapBufferSource;
19use crate::dom::bindings::error::{ErrorToJsval, Fallible};
20use crate::dom::bindings::refcounted::Trusted;
21use crate::dom::bindings::reflector::DomGlobal;
22use crate::dom::bindings::root::{Dom, DomRoot};
23use crate::dom::globalscope::GlobalScope;
24use crate::dom::promise::{RootedPromise, TracedPromise};
25use crate::dom::stream::byteteeunderlyingsource::ByteTeeUnderlyingSource;
26use crate::dom::stream::readablestream::ReadableStream;
27use crate::runtime::job_queue::MicrotaskRunnable;
28
29#[derive(JSTraceable, MallocSizeOf)]
30pub(crate) struct ByteTeeReadIntoRequestMicrotask {
31    #[ignore_malloc_size_of = "mozjs"]
32    chunk: RootedTraceableBox<HeapBufferSource<ArrayBufferViewU8>>,
33    tee_read_request: Trusted<ByteTeeReadIntoRequest>,
34}
35
36impl MicrotaskRunnable for ByteTeeReadIntoRequestMicrotask {
37    fn handler(&self, cx: &mut JSContext) {
38        self.tee_read_request
39            .root()
40            .chunk_steps(&self.chunk, cx)
41            .expect("Failed to enqueue chunk");
42    }
43}
44
45#[dom_struct]
46pub(crate) struct ByteTeeReadIntoRequest {
47    reflector_: Reflector,
48    for_branch2: bool,
49    byob_branch: Dom<ReadableStream>,
50    other_branch: Dom<ReadableStream>,
51    stream: Dom<ReadableStream>,
52    #[conditional_malloc_size_of]
53    read_again_for_branch_1: Rc<Cell<bool>>,
54    #[conditional_malloc_size_of]
55    read_again_for_branch_2: Rc<Cell<bool>>,
56    #[conditional_malloc_size_of]
57    reading: Rc<Cell<bool>>,
58    #[conditional_malloc_size_of]
59    canceled_1: Rc<Cell<bool>>,
60    #[conditional_malloc_size_of]
61    canceled_2: Rc<Cell<bool>>,
62    cancel_promise: TracedPromise,
63    tee_underlying_source: Dom<ByteTeeUnderlyingSource>,
64}
65impl ByteTeeReadIntoRequest {
66    #[allow(clippy::too_many_arguments)]
67    pub(crate) fn new(
68        cx: &mut JSContext,
69        for_branch2: bool,
70        byob_branch: &ReadableStream,
71        other_branch: &ReadableStream,
72        stream: &ReadableStream,
73        read_again_for_branch_1: Rc<Cell<bool>>,
74        read_again_for_branch_2: Rc<Cell<bool>>,
75        reading: Rc<Cell<bool>>,
76        canceled_1: Rc<Cell<bool>>,
77        canceled_2: Rc<Cell<bool>>,
78        cancel_promise: &RootedPromise,
79        tee_underlying_source: &ByteTeeUnderlyingSource,
80        global: &GlobalScope,
81    ) -> DomRoot<Self> {
82        reflect_dom_object(
83            cx,
84            Box::new(ByteTeeReadIntoRequest {
85                reflector_: Reflector::new(),
86                for_branch2,
87                byob_branch: Dom::from_ref(byob_branch),
88                other_branch: Dom::from_ref(other_branch),
89                stream: Dom::from_ref(stream),
90                read_again_for_branch_1,
91                read_again_for_branch_2,
92                reading,
93                canceled_1,
94                canceled_2,
95                cancel_promise: cancel_promise.to_traced(),
96                tee_underlying_source: Dom::from_ref(tee_underlying_source),
97            }),
98            global,
99        )
100    }
101
102    pub(crate) fn enqueue_chunk_steps(
103        &self,
104        cx: &mut JSContext,
105        chunk: RootedTraceableBox<HeapBufferSource<ArrayBufferViewU8>>,
106    ) {
107        // Queue a microtask to perform the following steps:
108        let byte_tee_read_request_chunk = ByteTeeReadIntoRequestMicrotask {
109            chunk,
110            tee_read_request: Trusted::new(self),
111        };
112
113        self.global()
114            .enqueue_microtask(cx, Box::new(byte_tee_read_request_chunk));
115    }
116
117    /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-chunk-steps%E2%91%A0>
118    #[allow(clippy::borrowed_box)]
119    pub(crate) fn chunk_steps(
120        &self,
121        chunk: &HeapBufferSource<ArrayBufferViewU8>,
122        cx: &mut JSContext,
123    ) -> Fallible<()> {
124        // Set readAgainForBranch1 to false.
125        self.read_again_for_branch_1.set(false);
126
127        // Set readAgainForBranch2 to false.
128        self.read_again_for_branch_2.set(false);
129
130        // Let byobCanceled be canceled2 if forBranch2 is true, and canceled1 otherwise.
131        let byob_canceled = if self.for_branch2 {
132            self.canceled_2.get()
133        } else {
134            self.canceled_1.get()
135        };
136
137        // Let otherCanceled be canceled2 if forBranch2 is false, and canceled1 otherwise.
138        let other_canceled = if self.for_branch2 {
139            self.canceled_1.get()
140        } else {
141            self.canceled_2.get()
142        };
143
144        // If otherCanceled is false,
145        if !other_canceled {
146            // Let cloneResult be CloneAsUint8Array(chunk).
147            let clone_result = chunk.clone_as_uint8_array(cx);
148
149            // If cloneResult is an abrupt completion,
150            if let Err(error) = clone_result {
151                rooted!(&in(cx) let mut error_value = UndefinedValue());
152                error.to_jsval(cx, &self.global(), error_value.handle_mut());
153
154                // Perform ! ReadableByteStreamControllerError(byobBranch.[[controller]], cloneResult.[[Value]]).
155                let byob_branch_controller = self.byob_branch.get_byte_controller();
156                byob_branch_controller.error(cx, error_value.handle());
157
158                // Perform ! ReadableByteStreamControllerError(otherBranch.[[controller]], cloneResult.[[Value]]).
159                let other_branch_controller = self.other_branch.get_byte_controller();
160                other_branch_controller.error(cx, error_value.handle());
161
162                // Resolve cancelPromise with ! ReadableStreamCancel(stream, cloneResult.[[Value]]).
163                let cancel_result =
164                    self.stream
165                        .cancel(cx, &self.stream.global(), error_value.handle());
166                self.cancel_promise.resolve_native(cx, &cancel_result);
167
168                // Return.
169                return Ok(());
170            } else {
171                // Otherwise, let clonedChunk be cloneResult.[[Value]].
172                let cloned_chunk = clone_result.unwrap();
173
174                // If byobCanceled is false, perform !
175                // ReadableByteStreamControllerRespondWithNewView(byobBranch.[[controller]], chunk).
176                if !byob_canceled {
177                    let byob_branch_controller = self.byob_branch.get_byte_controller();
178                    byob_branch_controller.respond_with_new_view(cx, chunk)?;
179                }
180
181                // Perform ! ReadableByteStreamControllerEnqueue(otherBranch.[[controller]], clonedChunk).
182                let other_branch_controller = self.other_branch.get_byte_controller();
183                other_branch_controller.enqueue(cx, cloned_chunk)?;
184            }
185        } else if !byob_canceled {
186            // Otherwise, if byobCanceled is false, perform
187            // ! ReadableByteStreamControllerRespondWithNewView(byobBranch.[[controller]], chunk).
188
189            let byob_branch_controller = self.byob_branch.get_byte_controller();
190            byob_branch_controller.respond_with_new_view(cx, chunk)?;
191        }
192
193        // Set reading to false.
194        self.reading.set(false);
195
196        // If readAgainForBranch1 is true, perform pull1Algorithm.
197        if self.read_again_for_branch_1.get() {
198            self.pull_algorithm(cx, Some(ByteTeePullAlgorithm::Pull1Algorithm));
199        } else if self.read_again_for_branch_2.get() {
200            // Otherwise, if readAgainForBranch2 is true, perform pull2Algorithm.
201            self.pull_algorithm(cx, Some(ByteTeePullAlgorithm::Pull2Algorithm));
202        }
203
204        Ok(())
205    }
206
207    /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-close-steps%E2%91%A1>
208    pub(crate) fn close_steps(
209        &self,
210        cx: &mut JSContext,
211        chunk: Option<RootedTraceableBox<HeapBufferSource<ArrayBufferViewU8>>>,
212    ) -> Fallible<()> {
213        // Set reading to false.
214        self.reading.set(false);
215
216        // Let byobCanceled be canceled2 if forBranch2 is true, and canceled1 otherwise.
217        let byob_canceled = if self.for_branch2 {
218            self.canceled_2.get()
219        } else {
220            self.canceled_1.get()
221        };
222
223        // Let otherCanceled be canceled2 if forBranch2 is false, and canceled1 otherwise.
224        let other_canceled = if self.for_branch2 {
225            self.canceled_1.get()
226        } else {
227            self.canceled_2.get()
228        };
229
230        // If byobCanceled is false, perform ! ReadableByteStreamControllerClose(byobBranch.[[controller]]).
231        if !byob_canceled {
232            let byob_branch_controller = self.byob_branch.get_byte_controller();
233            byob_branch_controller.close(cx)?;
234        }
235
236        // If otherCanceled is false, perform ! ReadableByteStreamControllerClose(otherBranch.[[controller]]).
237        if !other_canceled {
238            let other_branch_controller = self.other_branch.get_byte_controller();
239            other_branch_controller.close(cx)?;
240        }
241
242        // If chunk is not undefined,
243        if let Some(chunk_value) = chunk {
244            if chunk_value.is_undefined() {
245                // Nothing to respond with if the provided chunk is undefined.
246                // Continue with the remaining close steps.
247            } else {
248                let chunk = chunk_value;
249                // Assert: chunk.[[ByteLength]] is 0.
250                assert_eq!(chunk.byte_length(), 0);
251
252                // If byobCanceled is false, perform !
253                // ReadableByteStreamControllerRespondWithNewView(byobBranch.[[controller]], chunk).
254                if !byob_canceled {
255                    let byob_branch_controller = self.byob_branch.get_byte_controller();
256                    byob_branch_controller.respond_with_new_view(cx, &chunk)?;
257                }
258
259                // If otherCanceled is false and otherBranch.[[controller]].[[pendingPullIntos]] is not empty,
260                // perform ! ReadableByteStreamControllerRespond(otherBranch.[[controller]], 0).
261                if !other_canceled {
262                    let other_branch_controller = self.other_branch.get_byte_controller();
263                    if other_branch_controller.get_pending_pull_intos_size() > 0 {
264                        other_branch_controller.respond(cx, 0)?;
265                    }
266                }
267            }
268        }
269
270        // If byobCanceled is false or otherCanceled is false, resolve cancelPromise with undefined.
271        if !byob_canceled || !other_canceled {
272            self.cancel_promise.resolve_native(cx, &());
273        }
274
275        Ok(())
276    }
277    /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-error-steps%E2%91%A0>
278    pub(crate) fn error_steps(&self) {
279        // Set reading to false.
280        self.reading.set(false);
281    }
282
283    pub(crate) fn pull_algorithm(
284        &self,
285        cx: &mut JSContext,
286        byte_tee_pull_algorithm: Option<ByteTeePullAlgorithm>,
287    ) {
288        self.tee_underlying_source
289            .pull_algorithm(cx, byte_tee_pull_algorithm);
290    }
291}