Skip to main content

script/dom/stream/
readablestreambyobreader.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
5use std::cell::Cell;
6use std::collections::VecDeque;
7use std::mem;
8use std::rc::Rc;
9
10use dom_struct::dom_struct;
11use js::context::JSContext;
12use js::gc::CustomAutoRooterGuard;
13use js::jsapi::Heap;
14use js::jsval::{JSVal, UndefinedValue};
15use js::realm::CurrentRealm;
16use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue};
17use js::typedarray::{ArrayBufferView, ArrayBufferViewU8};
18use script_bindings::cell::DomRefCell;
19use script_bindings::reflector::{
20    Reflector, reflect_dom_object_with_cx, reflect_dom_object_with_proto,
21};
22use script_bindings::root::Dom;
23
24use super::byteteereadintorequest::ByteTeeReadIntoRequest;
25use super::readablebytestreamcontroller::ReadableByteStreamController;
26use super::readablestreamgenericreader::ReadableStreamGenericReader;
27use crate::dom::bindings::buffer_source::HeapBufferSource;
28use crate::dom::bindings::codegen::Bindings::ReadableStreamBYOBReaderBinding::{
29    ReadableStreamBYOBReaderMethods, ReadableStreamBYOBReaderReadOptions,
30};
31use crate::dom::bindings::codegen::Bindings::ReadableStreamDefaultReaderBinding::ReadableStreamReadResult;
32use crate::dom::bindings::error::{Error, ErrorToJsval, Fallible};
33use crate::dom::bindings::reflector::DomGlobal;
34use crate::dom::bindings::root::{DomRoot, MutNullableDom};
35use crate::dom::bindings::trace::RootedTraceableBox;
36use crate::dom::globalscope::GlobalScope;
37use crate::dom::promise::Promise;
38use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler};
39use crate::dom::stream::readablestream::ReadableStream;
40use crate::realms::enter_auto_realm;
41
42/// <https://streams.spec.whatwg.org/#read-into-request>
43#[derive(Clone, JSTraceable, MallocSizeOf)]
44pub enum ReadIntoRequest {
45    /// <https://streams.spec.whatwg.org/#byob-reader-read>
46    Read(#[conditional_malloc_size_of] Rc<Promise>),
47    ByteTee {
48        byte_tee_read_into_request: Dom<ByteTeeReadIntoRequest>,
49    },
50}
51
52impl ReadIntoRequest {
53    /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-chunk-steps%E2%91%A0>
54    pub fn chunk_steps(&self, cx: &mut JSContext, chunk: RootedTraceableBox<Heap<JSVal>>) {
55        match self {
56            ReadIntoRequest::Read(promise) => {
57                // chunk steps, given chunk
58                // Resolve promise with «[ "value" → chunk, "done" → false ]».
59                promise.resolve_native(
60                    cx,
61                    &ReadableStreamReadResult {
62                        done: Some(false),
63                        value: chunk,
64                    },
65                );
66            },
67            ReadIntoRequest::ByteTee {
68                byte_tee_read_into_request,
69            } => {
70                rooted!(&in(cx) let chunk_object = chunk.get().to_object());
71                byte_tee_read_into_request.enqueue_chunk_steps(
72                    cx,
73                    RootedTraceableBox::new(HeapBufferSource::<ArrayBufferViewU8>::new(
74                        chunk_object.handle(),
75                    )),
76                )
77            },
78        }
79    }
80
81    /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-close-steps%E2%91%A0>
82    pub fn close_steps(&self, cx: &mut JSContext, chunk: Option<RootedTraceableBox<Heap<JSVal>>>) {
83        match self {
84            ReadIntoRequest::Read(promise) => match chunk {
85                // close steps, given chunk
86                // Resolve promise with «[ "value" → chunk, "done" → true ]».
87                Some(chunk) => promise.resolve_native(
88                    cx,
89                    &ReadableStreamReadResult {
90                        done: Some(true),
91                        value: chunk,
92                    },
93                ),
94                None => {
95                    let result = RootedTraceableBox::new(Heap::default());
96                    result.set(UndefinedValue());
97                    promise.resolve_native(
98                        cx,
99                        &ReadableStreamReadResult {
100                            done: Some(true),
101                            value: result,
102                        },
103                    );
104                },
105            },
106            ReadIntoRequest::ByteTee {
107                byte_tee_read_into_request,
108            } => match chunk {
109                Some(chunk) => {
110                    rooted!(&in(cx) let chunk_object = chunk.get().to_object());
111                    byte_tee_read_into_request
112                        .close_steps(
113                            cx,
114                            Some(RootedTraceableBox::new(
115                                HeapBufferSource::<ArrayBufferViewU8>::new(chunk_object.handle()),
116                            )),
117                        )
118                        .expect("close steps should not fail")
119                },
120                None => byte_tee_read_into_request
121                    .close_steps(cx, None)
122                    .expect("close steps should not fail"),
123            },
124        }
125    }
126
127    /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-error-steps%E2%91%A0>
128    pub(crate) fn error_steps(&self, cx: &mut JSContext, e: SafeHandleValue) {
129        match self {
130            ReadIntoRequest::Read(promise) => {
131                // error steps, given e
132                // Reject promise with e.
133                promise.reject_native(cx, &e)
134            },
135            ReadIntoRequest::ByteTee {
136                byte_tee_read_into_request,
137            } => {
138                byte_tee_read_into_request.error_steps();
139            },
140        }
141    }
142}
143
144/// The rejection handler for
145/// <https://streams.spec.whatwg.org/#abstract-opdef-readablebytestreamtee>
146#[derive(Clone, JSTraceable, MallocSizeOf)]
147#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
148struct ByteTeeClosedPromiseRejectionHandler {
149    branch_1_controller: Dom<ReadableByteStreamController>,
150    branch_2_controller: Dom<ReadableByteStreamController>,
151    #[conditional_malloc_size_of]
152    canceled_1: Rc<Cell<bool>>,
153    #[conditional_malloc_size_of]
154    canceled_2: Rc<Cell<bool>>,
155    #[conditional_malloc_size_of]
156    cancel_promise: Rc<Promise>,
157    #[conditional_malloc_size_of]
158    reader_version: Rc<Cell<u64>>,
159    expected_version: u64,
160}
161
162impl Callback for ByteTeeClosedPromiseRejectionHandler {
163    /// Continuation of <https://streams.spec.whatwg.org/#abstract-opdef-readablebytestreamtee>
164    /// Upon rejection of `reader.closedPromise` with reason `r``,
165    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
166        // If thisReader is not reader, return.
167        if self.reader_version.get() != self.expected_version {
168            return;
169        }
170
171        // Perform ! ReadableByteStreamControllerError(branch1.[[controller]], r).
172        self.branch_1_controller.error(cx, v);
173
174        // Perform ! ReadableByteStreamControllerError(branch2.[[controller]], r).
175        self.branch_2_controller.error(cx, v);
176
177        // If canceled1 is false or canceled2 is false, resolve cancelPromise with undefined.
178        if !self.canceled_1.get() || !self.canceled_2.get() {
179            self.cancel_promise.resolve_native(cx, &());
180        }
181    }
182}
183
184/// <https://streams.spec.whatwg.org/#readablestreambyobreader>
185#[dom_struct]
186pub(crate) struct ReadableStreamBYOBReader {
187    reflector_: Reflector,
188
189    /// <https://streams.spec.whatwg.org/#readablestreamgenericreader-stream>
190    stream: MutNullableDom<ReadableStream>,
191
192    read_into_requests: DomRefCell<VecDeque<ReadIntoRequest>>,
193
194    /// <https://streams.spec.whatwg.org/#readablestreamgenericreader-closedpromise>
195    #[conditional_malloc_size_of]
196    closed_promise: DomRefCell<Rc<Promise>>,
197}
198
199impl ReadableStreamBYOBReader {
200    fn new_with_proto(
201        cx: &mut JSContext,
202        global: &GlobalScope,
203        proto: Option<SafeHandleObject>,
204    ) -> DomRoot<ReadableStreamBYOBReader> {
205        let closed_promise = Promise::new(cx, global);
206        reflect_dom_object_with_proto(
207            cx,
208            Box::new(ReadableStreamBYOBReader::new_inherited(closed_promise)),
209            global,
210            proto,
211        )
212    }
213
214    fn new_inherited(promise: Rc<Promise>) -> ReadableStreamBYOBReader {
215        ReadableStreamBYOBReader {
216            reflector_: Reflector::new(),
217            stream: MutNullableDom::new(None),
218            read_into_requests: DomRefCell::new(Default::default()),
219            closed_promise: DomRefCell::new(promise),
220        }
221    }
222
223    pub(crate) fn new(
224        cx: &mut JSContext,
225        global: &GlobalScope,
226    ) -> DomRoot<ReadableStreamBYOBReader> {
227        let closed_promise = Promise::new(cx, global);
228        reflect_dom_object_with_cx(Box::new(Self::new_inherited(closed_promise)), global, cx)
229    }
230
231    /// <https://streams.spec.whatwg.org/#set-up-readable-stream-byob-reader>
232    pub(crate) fn set_up(
233        &self,
234        cx: &mut JSContext,
235        stream: &ReadableStream,
236        global: &GlobalScope,
237    ) -> Fallible<()> {
238        // If ! IsReadableStreamLocked(stream) is true, throw a TypeError exception.
239        if stream.is_locked() {
240            return Err(Error::Type(c"stream is locked".to_owned()));
241        }
242
243        // If stream.[[controller]] does not implement ReadableByteStreamController, throw a TypeError exception.
244        if !stream.has_byte_controller() {
245            return Err(Error::Type(
246                c"stream controller is not a byte stream controller".to_owned(),
247            ));
248        }
249
250        // Perform ! ReadableStreamReaderGenericInitialize(reader, stream).
251        self.generic_initialize(cx, global, stream);
252
253        // Set reader.[[readIntoRequests]] to a new empty list.
254        self.read_into_requests.borrow_mut().clear();
255
256        Ok(())
257    }
258
259    /// <https://streams.spec.whatwg.org/#abstract-opdef-readablestreambyobreaderrelease>
260    pub(crate) fn release(&self, cx: &mut JSContext) -> Fallible<()> {
261        // Perform ! ReadableStreamReaderGenericRelease(reader).
262        self.generic_release(cx).expect("Generic release failed");
263        // Let e be a new TypeError exception.
264        rooted!(&in(cx) let mut error = UndefinedValue());
265        Error::Type(c"Reader is released".to_owned()).to_jsval(
266            cx,
267            &self.global(),
268            error.handle_mut(),
269        );
270
271        // Perform ! ReadableStreamBYOBReaderErrorReadIntoRequests(reader, e).
272        self.error_read_into_requests(cx, error.handle());
273        Ok(())
274    }
275
276    /// <https://streams.spec.whatwg.org/#abstract-opdef-readablestreambyobreadererrorreadintorequests>
277    pub(crate) fn error_read_into_requests(&self, cx: &mut JSContext, e: SafeHandleValue) {
278        // Reject reader.[[closedPromise]] with e.
279        self.closed_promise.borrow().reject_native(cx, &e);
280
281        // Set reader.[[closedPromise]].[[PromiseIsHandled]] to true.
282        self.closed_promise.borrow().set_promise_is_handled(cx);
283
284        // Let readRequests be reader.[[readRequests]].
285        let mut read_into_requests = self.take_read_into_requests();
286
287        // Set reader.[[readIntoRequests]] to a new empty list.
288        for request in read_into_requests.drain(0..) {
289            // Perform readIntoRequest’s error steps, given e.
290            request.error_steps(cx, e);
291        }
292    }
293
294    fn take_read_into_requests(&self) -> VecDeque<ReadIntoRequest> {
295        mem::take(&mut *self.read_into_requests.borrow_mut())
296    }
297
298    /// <https://streams.spec.whatwg.org/#readable-stream-add-read-into-request>
299    pub(crate) fn add_read_into_request(&self, read_request: &ReadIntoRequest) {
300        self.read_into_requests
301            .borrow_mut()
302            .push_back(read_request.clone());
303    }
304
305    /// <https://streams.spec.whatwg.org/#readable-stream-cancel>
306    pub(crate) fn cancel(&self, cx: &mut JSContext) {
307        // If reader is not undefined and reader implements ReadableStreamBYOBReader,
308        // Let readIntoRequests be reader.[[readIntoRequests]].
309        let mut read_into_requests = self.take_read_into_requests();
310        // Set reader.[[readIntoRequests]] to an empty list.
311        // Perform readIntoRequest’s close steps, given undefined.
312        for request in read_into_requests.drain(0..) {
313            // Perform readIntoRequest’s close steps, given undefined.
314            request.close_steps(cx, None);
315        }
316    }
317
318    pub(crate) fn close(&self, cx: &mut JSContext) {
319        // Resolve reader.[[closedPromise]] with undefined.
320        self.closed_promise.borrow().resolve_native(cx, &());
321    }
322
323    /// <https://streams.spec.whatwg.org/#readable-stream-byob-reader-read>
324    pub(crate) fn read(
325        &self,
326        cx: &mut JSContext,
327        view: &HeapBufferSource<ArrayBufferViewU8>,
328        min: u64,
329        read_into_request: &ReadIntoRequest,
330    ) {
331        // Let stream be reader.[[stream]].
332
333        // Assert: stream is not undefined.
334        assert!(self.stream.get().is_some());
335
336        let stream = self.stream.get().unwrap();
337
338        // Set stream.[[disturbed]] to true.
339        stream.set_is_disturbed(true);
340        // If stream.[[state]] is "errored", perform readIntoRequest’s error steps given stream.[[storedError]].
341        if stream.is_errored() {
342            rooted!(&in(cx) let mut error = UndefinedValue());
343            stream.get_stored_error(error.handle_mut());
344
345            read_into_request.error_steps(cx, error.handle());
346        } else {
347            // Otherwise,
348            // perform ! ReadableByteStreamControllerPullInto(stream.[[controller]], view, min, readIntoRequest).
349            stream.perform_pull_into(cx, read_into_request, view, min);
350        }
351    }
352
353    pub(crate) fn get_num_read_into_requests(&self) -> usize {
354        self.read_into_requests.borrow().len()
355    }
356
357    pub(crate) fn remove_read_into_request(&self) -> ReadIntoRequest {
358        self.read_into_requests
359            .borrow_mut()
360            .pop_front()
361            .expect("read into requests is empty")
362    }
363
364    #[allow(clippy::too_many_arguments)]
365    pub(crate) fn byte_tee_append_native_handler_to_closed_promise(
366        &self,
367        cx: &mut JSContext,
368        branch_1: &ReadableStream,
369        branch_2: &ReadableStream,
370        canceled_1: Rc<Cell<bool>>,
371        canceled_2: Rc<Cell<bool>>,
372        cancel_promise: Rc<Promise>,
373        reader_version: Rc<Cell<u64>>,
374        expected_version: u64,
375    ) {
376        let branch_1_controller = branch_1.get_byte_controller();
377
378        let branch_2_controller = branch_2.get_byte_controller();
379
380        let global = self.global();
381        let handler = PromiseNativeHandler::new(
382            cx,
383            &global,
384            None,
385            Some(Box::new(ByteTeeClosedPromiseRejectionHandler {
386                branch_1_controller: Dom::from_ref(&branch_1_controller),
387                branch_2_controller: Dom::from_ref(&branch_2_controller),
388                canceled_1,
389                canceled_2,
390                cancel_promise,
391                reader_version,
392                expected_version,
393            })),
394        );
395
396        let mut realm = enter_auto_realm(cx, &*global);
397        let cx = &mut realm.current_realm();
398
399        self.closed_promise
400            .borrow()
401            .append_native_handler(cx, &handler);
402    }
403}
404
405impl ReadableStreamBYOBReaderMethods<crate::DomTypeHolder> for ReadableStreamBYOBReader {
406    /// <https://streams.spec.whatwg.org/#byob-reader-constructor>
407    fn Constructor(
408        cx: &mut JSContext,
409        global: &GlobalScope,
410        proto: Option<SafeHandleObject>,
411        stream: &ReadableStream,
412    ) -> Fallible<DomRoot<Self>> {
413        let reader = Self::new_with_proto(cx, global, proto);
414
415        // Perform ? SetUpReadableStreamBYOBReader(this, stream).
416        reader.set_up(cx, stream, global)?;
417
418        Ok(reader)
419    }
420
421    /// <https://streams.spec.whatwg.org/#byob-reader-read>
422    fn Read(
423        &self,
424        cx: &mut JSContext,
425        view: CustomAutoRooterGuard<ArrayBufferView>,
426        options: &ReadableStreamBYOBReaderReadOptions,
427    ) -> Rc<Promise> {
428        let view = HeapBufferSource::<ArrayBufferViewU8>::from_view(cx, view);
429        let min = options.min;
430        // Let promise be a new promise.
431        let promise = Promise::new(cx, &self.global());
432
433        // If view.[[ByteLength]] is 0, return a promise rejected with a TypeError exception.
434        if view.byte_length() == 0 {
435            promise.reject_error(cx, Error::Type(c"view byte length is 0".to_owned()));
436            return promise;
437        }
438        // If view.[[ViewedArrayBuffer]].[[ArrayBufferByteLength]] is 0,
439        // return a promise rejected with a TypeError exception.
440        if view.viewed_buffer_array_byte_length(cx) == 0 {
441            promise.reject_error(
442                cx,
443                Error::Type(c"viewed buffer byte length is 0".to_owned()),
444            );
445            return promise;
446        }
447
448        // If ! IsDetachedBuffer(view.[[ViewedArrayBuffer]]) is true,
449        // return a promise rejected with a TypeError exception.
450        if view.is_detached_buffer(cx) {
451            promise.reject_error(cx, Error::Type(c"view is detached".to_owned()));
452            return promise;
453        }
454
455        // If options["min"] is 0, return a promise rejected with a TypeError exception.
456        if min == 0 {
457            promise.reject_error(cx, Error::Type(c"min is 0".to_owned()));
458            return promise;
459        }
460
461        // If view has a [[TypedArrayName]] internal slot,
462        if view.has_typed_array_name() {
463            // If options["min"] > view.[[ArrayLength]], return a promise rejected with a RangeError exception.
464            if min > (view.get_typed_array_length() as u64) {
465                promise.reject_error(
466                    cx,
467                    Error::Range(c"min is greater than array length".to_owned()),
468                );
469                return promise;
470            }
471        } else {
472            // Otherwise (i.e., it is a DataView),
473            // If options["min"] > view.[[ByteLength]], return a promise rejected with a RangeError exception.
474            if min > (view.byte_length() as u64) {
475                promise.reject_error(
476                    cx,
477                    Error::Range(c"min is greater than byte length".to_owned()),
478                );
479                return promise;
480            }
481        }
482
483        // If this.[[stream]] is undefined, return a promise rejected with a TypeError exception.
484        if self.stream.get().is_none() {
485            promise.reject_error(
486                cx,
487                Error::Type(c"min is greater than byte length".to_owned()),
488            );
489            return promise;
490        }
491
492        // Let readIntoRequest be a new read-into request with the following items:
493        //
494        // chunk steps, given chunk
495        // Resolve promise with «[ "value" → chunk, "done" → false ]».
496        //
497        // close steps, given chunk
498        // Resolve promise with «[ "value" → chunk, "done" → true ]».
499        //
500        // error steps, given e
501        // Reject promise with e
502        let read_into_request = ReadIntoRequest::Read(promise.clone());
503
504        // Perform ! ReadableStreamBYOBReaderRead(this, view, options["min"], readIntoRequest).
505        self.read(cx, &view, min, &read_into_request);
506
507        // Return promise.
508        promise
509    }
510
511    /// <https://streams.spec.whatwg.org/#byob-reader-release-lock>
512    fn ReleaseLock(&self, cx: &mut JSContext) -> Fallible<()> {
513        if self.stream.get().is_none() {
514            // If this.[[stream]] is undefined, return.
515            return Ok(());
516        }
517
518        // Perform !ReadableStreamBYOBReaderRelease(this).
519        self.release(cx)
520    }
521
522    /// <https://streams.spec.whatwg.org/#generic-reader-closed>
523    fn Closed(&self) -> Rc<Promise> {
524        self.closed()
525    }
526
527    /// <https://streams.spec.whatwg.org/#generic-reader-cancel>
528    fn Cancel(&self, cx: &mut JSContext, reason: SafeHandleValue) -> Rc<Promise> {
529        self.generic_cancel(cx, &self.global(), reason)
530    }
531}
532
533impl ReadableStreamGenericReader for ReadableStreamBYOBReader {
534    fn get_closed_promise(&self) -> Rc<Promise> {
535        self.closed_promise.borrow().clone()
536    }
537
538    fn set_closed_promise(&self, promise: Rc<Promise>) {
539        *self.closed_promise.borrow_mut() = promise;
540    }
541
542    fn set_stream(&self, stream: Option<&ReadableStream>) {
543        self.stream.set(stream);
544    }
545
546    fn get_stream(&self) -> Option<DomRoot<ReadableStream>> {
547        self.stream.get()
548    }
549
550    fn as_byob_reader(&self) -> Option<&ReadableStreamBYOBReader> {
551        Some(self)
552    }
553}