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