/
githubmirror
/
servo
Обзор
Документация
Войти
/
githubmirror
/
servo
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
components/script/dom/stream/readablestreamdefaultreader.rs
719 строк
27 KB
Euclid Ye
script: Reduce rooting in Stream (#46548)
16 июл 2026, 06:37
Не верифицирован
16 июл 2026, 06:37
ca1ea55
Код
Авторство
О чём код?
/* 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::jsapi::Heap; use js::jsval::{JSVal, UndefinedValue}; use js::realm::CurrentRealm; use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue}; use script_bindings::cell::DomRefCell; use script_bindings::reflector::{ Reflector, reflect_dom_object_with_cx, reflect_dom_object_with_proto, }; use super::byteteereadrequest::ByteTeeReadRequest; use super::readablebytestreamcontroller::ReadableByteStreamController; use crate::dom::bindings::codegen::Bindings::ReadableStreamDefaultReaderBinding::{ ReadableStreamDefaultReaderMethods, ReadableStreamReadResult, }; use crate::dom::bindings::error::{Error, ErrorToJsval, Fallible}; use crate::dom::bindings::reflector::DomGlobal; use crate::dom::bindings::root::{Dom, 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::readablestream::{ReadableStream, bytes_from_chunk_jsval}; use crate::dom::stream::defaultteereadrequest::DefaultTeeReadRequest; use crate::dom::stream::readablestreamgenericreader::ReadableStreamGenericReader; use crate::dom::types::ReadableStreamDefaultController; use crate::realms::enter_auto_realm; type ReadAllBytesSuccessSteps = dyn Fn(&mut js::context::JSContext, &[u8]); type ReadAllBytesFailureSteps = dyn Fn(&mut js::context::JSContext, SafeHandleValue); impl js::gc::Rootable for ContinueReadMicrotask {} /// Microtask handler to continue the read loop without recursion. /// Spec note: "This recursion could potentially cause a stack overflow /// if implemented directly. Implementations will need to mitigate this, /// e.g. by using a non-recursive variant of this algorithm, or queuing /// a microtask…" #[derive(Clone, JSTraceable, MallocSizeOf)] #[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)] struct ContinueReadMicrotask { reader: Dom<ReadableStreamDefaultReader>, request: ReadRequest, } impl Callback for ContinueReadMicrotask { fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) { // https://streams.spec.whatwg.org/#ref-for-read-loop%E2%91%A0 // Note: continuing the read-loop from inside a micro-task to break recursion. self.reader.read(cx, &self.request); } } /// <https://streams.spec.whatwg.org/#read-loop> fn read_loop( cx: &mut js::context::JSContext, reader: &ReadableStreamDefaultReader, success_steps: Rc<ReadAllBytesSuccessSteps>, failure_steps: Rc<ReadAllBytesFailureSteps>, ) { // For the purposes of the above algorithm, to read-loop given reader, // bytes, successSteps, and failureSteps: // Step 1 .Let readRequest be a new read request with the following items: let req = ReadRequest::ReadLoop { success_steps, failure_steps, reader: Dom::from_ref(reader), bytes: Rc::new(DomRefCell::new(Vec::new())), }; // Step 2 .Perform ! ReadableStreamDefaultReaderRead(reader, readRequest). reader.read(cx, &req); } /// <https://streams.spec.whatwg.org/#read-request> #[derive(Clone, JSTraceable, MallocSizeOf)] pub(crate) enum ReadRequest { /// <https://streams.spec.whatwg.org/#default-reader-read> Read(#[conditional_malloc_size_of] Rc<Promise>), /// <https://streams.spec.whatwg.org/#ref-for-read-request%E2%91%A2> DefaultTee { tee_read_request: Dom<DefaultTeeReadRequest>, }, /// Spec read loop variant, driven by read-request steps (no Promise). /// <https://streams.spec.whatwg.org/#read-loop> ReadLoop { #[ignore_malloc_size_of = "dyn Fn"] #[no_trace] success_steps: Rc<ReadAllBytesSuccessSteps>, #[ignore_malloc_size_of = "dyn Fn"] #[no_trace] failure_steps: Rc<ReadAllBytesFailureSteps>, reader: Dom<ReadableStreamDefaultReader>, #[conditional_malloc_size_of] bytes: Rc<DomRefCell<Vec<u8>>>, }, ByteTee { byte_tee_read_request: Dom<ByteTeeReadRequest>, }, } impl ReadRequest { /// <https://streams.spec.whatwg.org/#read-request-chunk-steps> pub(crate) fn chunk_steps( &self, cx: &mut js::context::JSContext, chunk: RootedTraceableBox<Heap<JSVal>>, global: &GlobalScope, ) { match self { ReadRequest::Read(promise) => { // chunk steps, given chunk // Resolve promise with «[ "value" → chunk, "done" → false ]». promise.resolve_native( cx, &ReadableStreamReadResult { done: Some(false), value: chunk, }, ); }, ReadRequest::DefaultTee { tee_read_request } => { tee_read_request.enqueue_chunk_steps(cx, chunk); }, ReadRequest::ByteTee { byte_tee_read_request, } => { byte_tee_read_request.enqueue_chunk_steps(cx, global, chunk); }, ReadRequest::ReadLoop { success_steps: _, failure_steps, reader, bytes, } => { // Spec: chunk steps, given chunk let global = reader.global(); match bytes_from_chunk_jsval(cx, &chunk) { Ok(vec) => { // Step 2. Append the bytes represented by chunk to bytes. bytes.borrow_mut().extend_from_slice(&vec); // Step 3. Read-loop given reader, bytes, successSteps, and failureSteps. // Spec note: Avoid direct recursion; queue into a microtask. // Resolving the promise will queue a microtask to call into the native handler. let tick = Promise::new(cx, &global); tick.resolve_native(cx, &()); let handler = PromiseNativeHandler::new( cx, &global, Some(Box::new(ContinueReadMicrotask { reader: Dom::from_ref(reader), request: self.clone(), })), None, ); let mut realm = enter_auto_realm(cx, &*global); let cx = &mut realm.current_realm(); tick.append_native_handler(cx, &handler); }, Err(err) => { // Step 1. If chunk is not a Uint8Array object, call failureSteps with a TypeError and abort. rooted!(&in(cx) let mut v = UndefinedValue()); err.to_jsval(cx, &global, v.handle_mut()); (failure_steps)(cx, v.handle()); }, } }, } } /// <https://streams.spec.whatwg.org/#read-request-close-steps> pub(crate) fn close_steps(&self, cx: &mut js::context::JSContext) { match self { ReadRequest::Read(promise) => { // close steps // Resolve promise with «[ "value" → undefined, "done" → true ]». let result = RootedTraceableBox::new(Heap::default()); result.set(UndefinedValue()); promise.resolve_native( cx, &ReadableStreamReadResult { done: Some(true), value: result, }, ); }, ReadRequest::DefaultTee { tee_read_request } => { tee_read_request.close_steps(cx); }, ReadRequest::ByteTee { byte_tee_read_request, } => { byte_tee_read_request .close_steps(cx) .expect("ByteTeeReadRequest close steps should not fail"); }, ReadRequest::ReadLoop { success_steps, reader, bytes, .. } => { // Step 1. Call successSteps with bytes. (success_steps)(cx, &bytes.borrow()); reader .release(cx) .expect("Releasing the read-all-bytes reader should succeed"); }, } } /// <https://streams.spec.whatwg.org/#read-request-error-steps> pub(crate) fn error_steps(&self, cx: &mut js::context::JSContext, e: SafeHandleValue) { match self { ReadRequest::Read(promise) => { // error steps, given e // Reject promise with e. promise.reject_native(cx, &e) }, ReadRequest::DefaultTee { tee_read_request } => { tee_read_request.error_steps(); }, ReadRequest::ByteTee { byte_tee_read_request, } => { byte_tee_read_request.error_steps(); }, ReadRequest::ReadLoop { failure_steps, reader, .. } => { // Step 1. Call failureSteps with e. (failure_steps)(cx, e); reader .release(cx) .expect("Releasing the read-all-bytes reader should succeed"); }, } } } /// 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 the current `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, &()); } } } /// The rejection handler for /// <https://streams.spec.whatwg.org/#readable-stream-tee> #[derive(Clone, JSTraceable, MallocSizeOf)] #[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)] struct DefaultTeeClosedPromiseRejectionHandler { branch_1_controller: Dom<ReadableStreamDefaultController>, branch_2_controller: Dom<ReadableStreamDefaultController>, #[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>, } impl Callback for DefaultTeeClosedPromiseRejectionHandler { /// Continuation of <https://streams.spec.whatwg.org/#abstract-opdef-readablestreamdefaulttee> /// Upon rejection of reader.[[closedPromise]] with reason r, fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) { // Perform ! ReadableStreamDefaultControllerError(branch_1.[[controller]], r). self.branch_1_controller.error(cx, v); // Perform ! ReadableStreamDefaultControllerError(branch_2.[[controller]], r). self.branch_2_controller.error(cx, v); // If canceled_1 is false or canceled_2 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/#readablestreamdefaultreader> #[dom_struct] pub(crate) struct ReadableStreamDefaultReader { reflector_: Reflector, /// <https://streams.spec.whatwg.org/#readablestreamgenericreader-stream> stream: MutNullableDom<ReadableStream>, read_requests: DomRefCell<VecDeque<ReadRequest>>, /// <https://streams.spec.whatwg.org/#readablestreamgenericreader-closedpromise> #[conditional_malloc_size_of] closed_promise: DomRefCell<Rc<Promise>>, } impl ReadableStreamDefaultReader { fn new_with_proto( cx: &mut JSContext, global: &GlobalScope, proto: Option<SafeHandleObject>, ) -> DomRoot<ReadableStreamDefaultReader> { let closed_promise = Promise::new(cx, global); reflect_dom_object_with_proto( cx, Box::new(ReadableStreamDefaultReader::new_inherited(closed_promise)), global, proto, ) } fn new_inherited(promise: Rc<Promise>) -> ReadableStreamDefaultReader { ReadableStreamDefaultReader { reflector_: Reflector::new(), stream: MutNullableDom::new(None), read_requests: DomRefCell::new(Default::default()), closed_promise: DomRefCell::new(promise), } } pub(crate) fn new( cx: &mut JSContext, global: &GlobalScope, ) -> DomRoot<ReadableStreamDefaultReader> { 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-default-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())); } // Perform ! ReadableStreamReaderGenericInitialize(reader, stream). self.generic_initialize(cx, global, stream); // Set reader.[[readRequests]] to a new empty list. self.read_requests.borrow_mut().clear(); Ok(()) } /// <https://streams.spec.whatwg.org/#readable-stream-close> pub(crate) fn close(&self, cx: &mut js::context::JSContext) { // Resolve reader.[[closedPromise]] with undefined. self.closed_promise.borrow().resolve_native(cx, &()); // If reader implements ReadableStreamDefaultReader, // Let readRequests be reader.[[readRequests]]. let mut read_requests = self.take_read_requests(); // Set reader.[[readRequests]] to an empty list. // For each readRequest of readRequests, for request in read_requests.drain(0..) { // Perform readRequest’s close steps. request.close_steps(cx); } } /// <https://streams.spec.whatwg.org/#readable-stream-add-read-request> pub(crate) fn add_read_request(&self, read_request: &ReadRequest) { self.read_requests .borrow_mut() .push_back(read_request.clone()); } /// <https://streams.spec.whatwg.org/#readable-stream-get-num-read-requests> pub(crate) fn get_num_read_requests(&self) -> usize { self.read_requests.borrow().len() } /// <https://streams.spec.whatwg.org/#readable-stream-error> pub(crate) fn error(&self, cx: &mut js::context::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); // Perform ! ReadableStreamDefaultReaderErrorReadRequests(reader, e). self.error_read_requests(cx, e); } /// The removal steps of <https://streams.spec.whatwg.org/#readable-stream-fulfill-read-request> pub(crate) fn remove_read_request(&self) -> ReadRequest { self.read_requests .borrow_mut() .pop_front() .expect("Reader must have read request when remove is called into.") } /// <https://streams.spec.whatwg.org/#abstract-opdef-readablestreamdefaultreaderrelease> pub(crate) fn release(&self, cx: &mut js::context::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 ! ReadableStreamDefaultReaderErrorReadRequests(reader, e). self.error_read_requests(cx, error.handle()); Ok(()) } fn take_read_requests(&self) -> VecDeque<ReadRequest> { mem::take(&mut *self.read_requests.borrow_mut()) } /// <https://streams.spec.whatwg.org/#abstract-opdef-readablestreamdefaultreadererrorreadrequests> fn error_read_requests(&self, cx: &mut js::context::JSContext, rval: SafeHandleValue) { // step 1 let mut read_requests = self.take_read_requests(); // step 2 & 3 for request in read_requests.drain(0..) { request.error_steps(cx, rval); } } /// <https://streams.spec.whatwg.org/#readable-stream-default-reader-read> pub(crate) fn read(&self, cx: &mut js::context::JSContext, read_request: &ReadRequest) { // 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 "closed", perform readRequest’s close steps. if stream.is_closed() { read_request.close_steps(cx); } else if stream.is_errored() { // Otherwise, if stream.[[state]] is "errored", // perform readRequest’s error steps given stream.[[storedError]]. rooted!(&in(cx) let mut error = UndefinedValue()); stream.get_stored_error(error.handle_mut()); read_request.error_steps(cx, error.handle()); } else { // Otherwise // Assert: stream.[[state]] is "readable". assert!(stream.is_readable()); // Perform ! stream.[[controller]].[[PullSteps]](readRequest). stream.perform_pull_steps(cx, read_request); } } /// Attach the byte-tee error handler to this reader's closedPromise. /// Used by ReadableByteStreamTee. #[allow(clippy::too_many_arguments)] pub(crate) fn byte_tee_append_native_handler_to_closed_promise( &self, cx: &mut js::context::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, ) { // Note: for byte tee we always operate on *byte controllers*. 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); } /// <https://streams.spec.whatwg.org/#ref-for-readablestreamgenericreader-closedpromise%E2%91%A1> pub(crate) fn default_tee_append_native_handler_to_closed_promise( &self, cx: &mut js::context::JSContext, branch_1: &ReadableStream, branch_2: &ReadableStream, canceled_1: Rc<Cell<bool>>, canceled_2: Rc<Cell<bool>>, cancel_promise: Rc<Promise>, ) { let branch_1_controller = branch_1.get_default_controller(); let branch_2_controller = branch_2.get_default_controller(); let global = self.global(); let handler = PromiseNativeHandler::new( cx, &global, None, Some(Box::new(DefaultTeeClosedPromiseRejectionHandler { branch_1_controller: Dom::from_ref(&branch_1_controller), branch_2_controller: Dom::from_ref(&branch_2_controller), canceled_1, canceled_2, cancel_promise, })), ); let mut realm = enter_auto_realm(cx, &*global); let cx = &mut realm.current_realm(); self.closed_promise .borrow() .append_native_handler(cx, &handler); } /// <https://streams.spec.whatwg.org/#readablestreamdefaultreader-read-all-bytes> pub(crate) fn read_all_bytes( &self, cx: &mut js::context::JSContext, success_steps: Rc<ReadAllBytesSuccessSteps>, failure_steps: Rc<ReadAllBytesFailureSteps>, ) { // To read all bytes from a ReadableStreamDefaultReader reader, // given successSteps, which is an algorithm accepting a byte sequence, // and failureSteps, which is an algorithm accepting a JavaScript value: // read-loop given reader, a new byte sequence, successSteps, and failureSteps. read_loop(cx, self, success_steps, failure_steps); } /// step 3 of <https://streams.spec.whatwg.org/#abstract-opdef-readablebytestreamcontrollerprocessreadrequestsusingqueue> pub(crate) fn process_read_requests( &self, cx: &mut js::context::JSContext, controller: &ReadableByteStreamController, ) -> Fallible<()> { // While reader.[[readRequests]] is not empty, while !self.read_requests.borrow().is_empty() { // If controller.[[queueTotalSize]] is 0, return. if controller.get_queue_total_size() == 0.0 { return Ok(()); } // Let readRequest be reader.[[readRequests]][0]. // Remove entry from controller.[[queue]]. let read_request = self.remove_read_request(); // Perform ! ReadableByteStreamControllerFillReadRequestFromQueue(controller, readRequest). controller .fill_read_request_from_queue(cx, &read_request) .expect("Fill read request from queue failed"); } Ok(()) } } impl ReadableStreamDefaultReaderMethods<crate::DomTypeHolder> for ReadableStreamDefaultReader { /// <https://streams.spec.whatwg.org/#default-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 ? SetUpReadableStreamDefaultReader(this, stream). reader.set_up(cx, stream, global)?; Ok(reader) } /// <https://streams.spec.whatwg.org/#default-reader-read> fn Read(&self, cx: &mut js::context::JSContext) -> Rc<Promise> { // If this.[[stream]] is undefined, return a promise rejected with a TypeError exception. if self.stream.get().is_none() { rooted!(&in(cx) let mut error = UndefinedValue()); Error::Type(c"stream is undefined".to_owned()).to_jsval( cx, &self.global(), error.handle_mut(), ); return Promise::new_rejected(cx, &self.global(), error.handle()); } // Let promise be a new promise. let promise = Promise::new(cx, &self.global()); // Let readRequest be a new read request with the following items: // chunk steps, given chunk // Resolve promise with «[ "value" → chunk, "done" → false ]». // // close steps // Resolve promise with «[ "value" → undefined, "done" → true ]». // // error steps, given e // Reject promise with e. // Rooting(unrooted_must_root): the read request contains only a promise, // which does not need to be rooted, // as it is safely managed natively via an Rc. let read_request = ReadRequest::Read(promise.clone()); // Perform ! ReadableStreamDefaultReaderRead(this, readRequest). self.read(cx, &read_request); // Return promise. promise } /// <https://streams.spec.whatwg.org/#default-reader-release-lock> fn ReleaseLock(&self, cx: &mut js::context::JSContext) -> Fallible<()> { if self.stream.get().is_none() { // Step 1: If this.[[stream]] is undefined, return. return Ok(()); } // Step 2: Perform !ReadableStreamDefaultReaderRelease(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 js::context::JSContext, reason: SafeHandleValue) -> Rc<Promise> { self.generic_cancel(cx, &self.global(), reason) } } impl ReadableStreamGenericReader for ReadableStreamDefaultReader { 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_default_reader(&self) -> Option<&ReadableStreamDefaultReader> { Some(self) } }