/
githubmirror
/
servo
Обзор
Документация
Войти
/
githubmirror
/
servo
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
components/script/dom/stream/readablestreambyobreader.rs
553 строки
20 KB
Tim van der Lippe
script: Mechanically migrate to `reflect_dom_object_with_proto` (#46529)
15 июл 2026, 23:09
Не верифицирован
15 июл 2026, 23:09
2557009
Код
Авторство
О чём код?
/* This Source Code Form is subject to the terms of the Mozilla Public * License, v. 2.0. If a copy of the MPL was not distributed with this * file, You can obtain one at http://mozilla.org/MPL/2.0/. */ use std::cell::Cell; use std::collections::VecDeque; use std::mem; use std::rc::Rc; use dom_struct::dom_struct; use js::context::JSContext; use js::gc::CustomAutoRooterGuard; use js::jsapi::Heap; use js::jsval::{JSVal, UndefinedValue}; use js::realm::CurrentRealm; use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue}; use js::typedarray::{ArrayBufferView, ArrayBufferViewU8}; use script_bindings::cell::DomRefCell; use script_bindings::reflector::{ Reflector, reflect_dom_object_with_cx, reflect_dom_object_with_proto, }; use script_bindings::root::Dom; use super::byteteereadintorequest::ByteTeeReadIntoRequest; use super::readablebytestreamcontroller::ReadableByteStreamController; use super::readablestreamgenericreader::ReadableStreamGenericReader; use crate::dom::bindings::buffer_source::HeapBufferSource; use crate::dom::bindings::codegen::Bindings::ReadableStreamBYOBReaderBinding::{ ReadableStreamBYOBReaderMethods, ReadableStreamBYOBReaderReadOptions, }; use crate::dom::bindings::codegen::Bindings::ReadableStreamDefaultReaderBinding::ReadableStreamReadResult; use crate::dom::bindings::error::{Error, ErrorToJsval, Fallible}; use crate::dom::bindings::reflector::DomGlobal; use crate::dom::bindings::root::{DomRoot, MutNullableDom}; use crate::dom::bindings::trace::RootedTraceableBox; use crate::dom::globalscope::GlobalScope; use crate::dom::promise::Promise; use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler}; use crate::dom::stream::readablestream::ReadableStream; use crate::realms::enter_auto_realm; /// <https://streams.spec.whatwg.org/#read-into-request> #[derive(Clone, JSTraceable, MallocSizeOf)] pub enum ReadIntoRequest { /// <https://streams.spec.whatwg.org/#byob-reader-read> Read(#[conditional_malloc_size_of] Rc<Promise>), ByteTee { byte_tee_read_into_request: Dom<ByteTeeReadIntoRequest>, }, } impl ReadIntoRequest { /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-chunk-steps%E2%91%A0> pub fn chunk_steps(&self, cx: &mut JSContext, chunk: RootedTraceableBox<Heap<JSVal>>) { match self { ReadIntoRequest::Read(promise) => { // chunk steps, given chunk // Resolve promise with «[ "value" → chunk, "done" → false ]». promise.resolve_native( cx, &ReadableStreamReadResult { done: Some(false), value: chunk, }, ); }, ReadIntoRequest::ByteTee { byte_tee_read_into_request, } => { rooted!(&in(cx) let chunk_object = chunk.get().to_object()); byte_tee_read_into_request.enqueue_chunk_steps( cx, RootedTraceableBox::new(HeapBufferSource::<ArrayBufferViewU8>::new( chunk_object.handle(), )), ) }, } } /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-close-steps%E2%91%A0> pub fn close_steps(&self, cx: &mut JSContext, chunk: Option<RootedTraceableBox<Heap<JSVal>>>) { match self { ReadIntoRequest::Read(promise) => match chunk { // close steps, given chunk // Resolve promise with «[ "value" → chunk, "done" → true ]». Some(chunk) => promise.resolve_native( cx, &ReadableStreamReadResult { done: Some(true), value: chunk, }, ), None => { let result = RootedTraceableBox::new(Heap::default()); result.set(UndefinedValue()); promise.resolve_native( cx, &ReadableStreamReadResult { done: Some(true), value: result, }, ); }, }, ReadIntoRequest::ByteTee { byte_tee_read_into_request, } => match chunk { Some(chunk) => { rooted!(&in(cx) let chunk_object = chunk.get().to_object()); byte_tee_read_into_request .close_steps( cx, Some(RootedTraceableBox::new( HeapBufferSource::<ArrayBufferViewU8>::new(chunk_object.handle()), )), ) .expect("close steps should not fail") }, None => byte_tee_read_into_request .close_steps(cx, None) .expect("close steps should not fail"), }, } } /// <https://streams.spec.whatwg.org/#ref-for-read-into-request-error-steps%E2%91%A0> pub(crate) fn error_steps(&self, cx: &mut JSContext, e: SafeHandleValue) { match self { ReadIntoRequest::Read(promise) => { // error steps, given e // Reject promise with e. promise.reject_native(cx, &e) }, ReadIntoRequest::ByteTee { byte_tee_read_into_request, } => { byte_tee_read_into_request.error_steps(); }, } } } /// The rejection handler for /// <https://streams.spec.whatwg.org/#abstract-opdef-readablebytestreamtee> #[derive(Clone, JSTraceable, MallocSizeOf)] #[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)] struct ByteTeeClosedPromiseRejectionHandler { branch_1_controller: Dom<ReadableByteStreamController>, branch_2_controller: Dom<ReadableByteStreamController>, #[conditional_malloc_size_of] canceled_1: Rc<Cell<bool>>, #[conditional_malloc_size_of] canceled_2: Rc<Cell<bool>>, #[conditional_malloc_size_of] cancel_promise: Rc<Promise>, #[conditional_malloc_size_of] reader_version: Rc<Cell<u64>>, expected_version: u64, } impl Callback for ByteTeeClosedPromiseRejectionHandler { /// Continuation of <https://streams.spec.whatwg.org/#abstract-opdef-readablebytestreamtee> /// Upon rejection of `reader.closedPromise` with reason `r``, fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) { // If thisReader is not reader, return. if self.reader_version.get() != self.expected_version { return; } // Perform ! ReadableByteStreamControllerError(branch1.[[controller]], r). self.branch_1_controller.error(cx, v); // Perform ! ReadableByteStreamControllerError(branch2.[[controller]], r). self.branch_2_controller.error(cx, v); // If canceled1 is false or canceled2 is false, resolve cancelPromise with undefined. if !self.canceled_1.get() || !self.canceled_2.get() { self.cancel_promise.resolve_native(cx, &()); } } } /// <https://streams.spec.whatwg.org/#readablestreambyobreader> #[dom_struct] pub(crate) struct ReadableStreamBYOBReader { reflector_: Reflector, /// <https://streams.spec.whatwg.org/#readablestreamgenericreader-stream> stream: MutNullableDom<ReadableStream>, read_into_requests: DomRefCell<VecDeque<ReadIntoRequest>>, /// <https://streams.spec.whatwg.org/#readablestreamgenericreader-closedpromise> #[conditional_malloc_size_of] closed_promise: DomRefCell<Rc<Promise>>, } impl ReadableStreamBYOBReader { fn new_with_proto( cx: &mut JSContext, global: &GlobalScope, proto: Option<SafeHandleObject>, ) -> DomRoot<ReadableStreamBYOBReader> { let closed_promise = Promise::new(cx, global); reflect_dom_object_with_proto( cx, Box::new(ReadableStreamBYOBReader::new_inherited(closed_promise)), global, proto, ) } fn new_inherited(promise: Rc<Promise>) -> ReadableStreamBYOBReader { ReadableStreamBYOBReader { reflector_: Reflector::new(), stream: MutNullableDom::new(None), read_into_requests: DomRefCell::new(Default::default()), closed_promise: DomRefCell::new(promise), } } pub(crate) fn new( cx: &mut JSContext, global: &GlobalScope, ) -> DomRoot<ReadableStreamBYOBReader> { let closed_promise = Promise::new(cx, global); reflect_dom_object_with_cx(Box::new(Self::new_inherited(closed_promise)), global, cx) } /// <https://streams.spec.whatwg.org/#set-up-readable-stream-byob-reader> pub(crate) fn set_up( &self, cx: &mut JSContext, stream: &ReadableStream, global: &GlobalScope, ) -> Fallible<()> { // If ! IsReadableStreamLocked(stream) is true, throw a TypeError exception. if stream.is_locked() { return Err(Error::Type(c"stream is locked".to_owned())); } // If stream.[[controller]] does not implement ReadableByteStreamController, throw a TypeError exception. if !stream.has_byte_controller() { return Err(Error::Type( c"stream controller is not a byte stream controller".to_owned(), )); } // Perform ! ReadableStreamReaderGenericInitialize(reader, stream). self.generic_initialize(cx, global, stream); // Set reader.[[readIntoRequests]] to a new empty list. self.read_into_requests.borrow_mut().clear(); Ok(()) } /// <https://streams.spec.whatwg.org/#abstract-opdef-readablestreambyobreaderrelease> pub(crate) fn release(&self, cx: &mut JSContext) -> Fallible<()> { // Perform ! ReadableStreamReaderGenericRelease(reader). self.generic_release(cx).expect("Generic release failed"); // Let e be a new TypeError exception. rooted!(&in(cx) let mut error = UndefinedValue()); Error::Type(c"Reader is released".to_owned()).to_jsval( cx, &self.global(), error.handle_mut(), ); // Perform ! ReadableStreamBYOBReaderErrorReadIntoRequests(reader, e). self.error_read_into_requests(cx, error.handle()); Ok(()) } /// <https://streams.spec.whatwg.org/#abstract-opdef-readablestreambyobreadererrorreadintorequests> pub(crate) fn error_read_into_requests(&self, cx: &mut JSContext, e: SafeHandleValue) { // Reject reader.[[closedPromise]] with e. self.closed_promise.borrow().reject_native(cx, &e); // Set reader.[[closedPromise]].[[PromiseIsHandled]] to true. self.closed_promise.borrow().set_promise_is_handled(cx); // Let readRequests be reader.[[readRequests]]. let mut read_into_requests = self.take_read_into_requests(); // Set reader.[[readIntoRequests]] to a new empty list. for request in read_into_requests.drain(0..) { // Perform readIntoRequest’s error steps, given e. request.error_steps(cx, e); } } fn take_read_into_requests(&self) -> VecDeque<ReadIntoRequest> { mem::take(&mut *self.read_into_requests.borrow_mut()) } /// <https://streams.spec.whatwg.org/#readable-stream-add-read-into-request> pub(crate) fn add_read_into_request(&self, read_request: &ReadIntoRequest) { self.read_into_requests .borrow_mut() .push_back(read_request.clone()); } /// <https://streams.spec.whatwg.org/#readable-stream-cancel> pub(crate) fn cancel(&self, cx: &mut JSContext) { // If reader is not undefined and reader implements ReadableStreamBYOBReader, // Let readIntoRequests be reader.[[readIntoRequests]]. let mut read_into_requests = self.take_read_into_requests(); // Set reader.[[readIntoRequests]] to an empty list. // Perform readIntoRequest’s close steps, given undefined. for request in read_into_requests.drain(0..) { // Perform readIntoRequest’s close steps, given undefined. request.close_steps(cx, None); } } pub(crate) fn close(&self, cx: &mut JSContext) { // Resolve reader.[[closedPromise]] with undefined. self.closed_promise.borrow().resolve_native(cx, &()); } /// <https://streams.spec.whatwg.org/#readable-stream-byob-reader-read> pub(crate) fn read( &self, cx: &mut JSContext, view: &HeapBufferSource<ArrayBufferViewU8>, min: u64, read_into_request: &ReadIntoRequest, ) { // Let stream be reader.[[stream]]. // Assert: stream is not undefined. assert!(self.stream.get().is_some()); let stream = self.stream.get().unwrap(); // Set stream.[[disturbed]] to true. stream.set_is_disturbed(true); // If stream.[[state]] is "errored", perform readIntoRequest’s error steps given stream.[[storedError]]. if stream.is_errored() { rooted!(&in(cx) let mut error = UndefinedValue()); stream.get_stored_error(error.handle_mut()); read_into_request.error_steps(cx, error.handle()); } else { // Otherwise, // perform ! ReadableByteStreamControllerPullInto(stream.[[controller]], view, min, readIntoRequest). stream.perform_pull_into(cx, read_into_request, view, min); } } pub(crate) fn get_num_read_into_requests(&self) -> usize { self.read_into_requests.borrow().len() } pub(crate) fn remove_read_into_request(&self) -> ReadIntoRequest { self.read_into_requests .borrow_mut() .pop_front() .expect("read into requests is empty") } #[allow(clippy::too_many_arguments)] pub(crate) fn byte_tee_append_native_handler_to_closed_promise( &self, cx: &mut JSContext, branch_1: &ReadableStream, branch_2: &ReadableStream, canceled_1: Rc<Cell<bool>>, canceled_2: Rc<Cell<bool>>, cancel_promise: Rc<Promise>, reader_version: Rc<Cell<u64>>, expected_version: u64, ) { let branch_1_controller = branch_1.get_byte_controller(); let branch_2_controller = branch_2.get_byte_controller(); let global = self.global(); let handler = PromiseNativeHandler::new( cx, &global, None, Some(Box::new(ByteTeeClosedPromiseRejectionHandler { branch_1_controller: Dom::from_ref(&branch_1_controller), branch_2_controller: Dom::from_ref(&branch_2_controller), canceled_1, canceled_2, cancel_promise, reader_version, expected_version, })), ); let mut realm = enter_auto_realm(cx, &*global); let cx = &mut realm.current_realm(); self.closed_promise .borrow() .append_native_handler(cx, &handler); } } impl ReadableStreamBYOBReaderMethods<crate::DomTypeHolder> for ReadableStreamBYOBReader { /// <https://streams.spec.whatwg.org/#byob-reader-constructor> fn Constructor( cx: &mut JSContext, global: &GlobalScope, proto: Option<SafeHandleObject>, stream: &ReadableStream, ) -> Fallible<DomRoot<Self>> { let reader = Self::new_with_proto(cx, global, proto); // Perform ? SetUpReadableStreamBYOBReader(this, stream). reader.set_up(cx, stream, global)?; Ok(reader) } /// <https://streams.spec.whatwg.org/#byob-reader-read> fn Read( &self, cx: &mut JSContext, view: CustomAutoRooterGuard<ArrayBufferView>, options: &ReadableStreamBYOBReaderReadOptions, ) -> Rc<Promise> { let view = HeapBufferSource::<ArrayBufferViewU8>::from_view(cx, view); let min = options.min; // Let promise be a new promise. let promise = Promise::new(cx, &self.global()); // If view.[[ByteLength]] is 0, return a promise rejected with a TypeError exception. if view.byte_length() == 0 { promise.reject_error(cx, Error::Type(c"view byte length is 0".to_owned())); return promise; } // If view.[[ViewedArrayBuffer]].[[ArrayBufferByteLength]] is 0, // return a promise rejected with a TypeError exception. if view.viewed_buffer_array_byte_length(cx) == 0 { promise.reject_error( cx, Error::Type(c"viewed buffer byte length is 0".to_owned()), ); return promise; } // If ! IsDetachedBuffer(view.[[ViewedArrayBuffer]]) is true, // return a promise rejected with a TypeError exception. if view.is_detached_buffer(cx) { promise.reject_error(cx, Error::Type(c"view is detached".to_owned())); return promise; } // If options["min"] is 0, return a promise rejected with a TypeError exception. if min == 0 { promise.reject_error(cx, Error::Type(c"min is 0".to_owned())); return promise; } // If view has a [[TypedArrayName]] internal slot, if view.has_typed_array_name() { // If options["min"] > view.[[ArrayLength]], return a promise rejected with a RangeError exception. if min > (view.get_typed_array_length() as u64) { promise.reject_error( cx, Error::Range(c"min is greater than array length".to_owned()), ); return promise; } } else { // Otherwise (i.e., it is a DataView), // If options["min"] > view.[[ByteLength]], return a promise rejected with a RangeError exception. if min > (view.byte_length() as u64) { promise.reject_error( cx, Error::Range(c"min is greater than byte length".to_owned()), ); return promise; } } // If this.[[stream]] is undefined, return a promise rejected with a TypeError exception. if self.stream.get().is_none() { promise.reject_error( cx, Error::Type(c"min is greater than byte length".to_owned()), ); return promise; } // Let readIntoRequest be a new read-into request with the following items: // // chunk steps, given chunk // Resolve promise with «[ "value" → chunk, "done" → false ]». // // close steps, given chunk // Resolve promise with «[ "value" → chunk, "done" → true ]». // // error steps, given e // Reject promise with e let read_into_request = ReadIntoRequest::Read(promise.clone()); // Perform ! ReadableStreamBYOBReaderRead(this, view, options["min"], readIntoRequest). self.read(cx, &view, min, &read_into_request); // Return promise. promise } /// <https://streams.spec.whatwg.org/#byob-reader-release-lock> fn ReleaseLock(&self, cx: &mut JSContext) -> Fallible<()> { if self.stream.get().is_none() { // If this.[[stream]] is undefined, return. return Ok(()); } // Perform !ReadableStreamBYOBReaderRelease(this). self.release(cx) } /// <https://streams.spec.whatwg.org/#generic-reader-closed> fn Closed(&self) -> Rc<Promise> { self.closed() } /// <https://streams.spec.whatwg.org/#generic-reader-cancel> fn Cancel(&self, cx: &mut JSContext, reason: SafeHandleValue) -> Rc<Promise> { self.generic_cancel(cx, &self.global(), reason) } } impl ReadableStreamGenericReader for ReadableStreamBYOBReader { fn get_closed_promise(&self) -> Rc<Promise> { self.closed_promise.borrow().clone() } fn set_closed_promise(&self, promise: Rc<Promise>) { *self.closed_promise.borrow_mut() = promise; } fn set_stream(&self, stream: Option<&ReadableStream>) { self.stream.set(stream); } fn get_stream(&self) -> Option<DomRoot<ReadableStream>> { self.stream.get() } fn as_byob_reader(&self) -> Option<&ReadableStreamBYOBReader> { Some(self) } }