/
githubmirror
/
servo
Обзор
Документация
Войти
/
githubmirror
/
servo
Код
Запросы
0
Пакеты
0
Релизы
0
Аналитика
Безопасность
main
components/script/dom/stream/writablestreamdefaultwriter.rs
539 строк
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 https://mozilla.org/MPL/2.0/. */ use std::cell::RefCell; use std::rc::Rc; use dom_struct::dom_struct; use js::context::JSContext; use js::jsval::UndefinedValue; use js::realm::CurrentRealm; use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue}; use script_bindings::reflector::{Reflector, reflect_dom_object_with_proto}; use crate::dom::bindings::codegen::Bindings::WritableStreamDefaultWriterBinding::WritableStreamDefaultWriterMethods; use crate::dom::bindings::error::{Error, ErrorToJsval}; use crate::dom::bindings::reflector::DomGlobal; use crate::dom::bindings::root::{DomRoot, MutNullableDom}; use crate::dom::globalscope::GlobalScope; use crate::dom::promise::Promise; use crate::dom::stream::writablestream::WritableStream; /// <https://streams.spec.whatwg.org/#writablestreamdefaultwriter> #[dom_struct] pub struct WritableStreamDefaultWriter { reflector_: Reflector, #[conditional_malloc_size_of] ready_promise: RefCell<Rc<Promise>>, /// <https://streams.spec.whatwg.org/#writablestreamdefaultwriter-closedpromise> #[conditional_malloc_size_of] closed_promise: RefCell<Rc<Promise>>, /// <https://streams.spec.whatwg.org/#writablestreamdefaultwriter-stream> stream: MutNullableDom<WritableStream>, } impl WritableStreamDefaultWriter { /// <https://streams.spec.whatwg.org/#set-up-writable-stream-default-writer> /// The parts that create a new promise. fn new_inherited( closed_promise: Rc<Promise>, ready_promise: Rc<Promise>, ) -> WritableStreamDefaultWriter { WritableStreamDefaultWriter { reflector_: Reflector::new(), stream: Default::default(), closed_promise: RefCell::new(closed_promise), ready_promise: RefCell::new(ready_promise), } } pub(crate) fn new( cx: &mut CurrentRealm, global: &GlobalScope, proto: Option<SafeHandleObject>, ) -> DomRoot<WritableStreamDefaultWriter> { let closed_promise = Promise::new_in_realm(cx); let ready_promise = Promise::new_in_realm(cx); reflect_dom_object_with_proto( cx, Box::new(WritableStreamDefaultWriter::new_inherited( closed_promise, ready_promise, )), global, proto, ) } /// <https://streams.spec.whatwg.org/#set-up-writable-stream-default-writer> /// Continuing from `new_inherited`, the rest. pub(crate) fn setup(&self, cx: &mut JSContext, stream: &WritableStream) -> Result<(), Error> { // If ! IsWritableStreamLocked(stream) is true, throw a TypeError exception. if stream.is_locked() { return Err(Error::Type(c"Stream is locked".to_owned())); } // Set writer.[[stream]] to stream. self.stream.set(Some(stream)); // Set stream.[[writer]] to writer. stream.set_writer(Some(self)); // Let state be stream.[[state]]. // If state is "writable", if stream.is_writable() { // If ! WritableStreamCloseQueuedOrInFlight(stream) is false // and stream.[[backpressure]] is true, if !stream.close_queued_or_in_flight() && stream.get_backpressure() { // set writer.[[readyPromise]] to a new promise. // Done in `new_inherited`. } else { // Otherwise, set writer.[[readyPromise]] to a promise resolved with undefined. // Note: new promise created in `new_inherited`. self.ready_promise.borrow().resolve_native(cx, &()); } // Set writer.[[closedPromise]] to a new promise. // Done in `new_inherited`. return Ok(()); } // Otherwise, if state is "erroring", if stream.is_erroring() { rooted!(&in(cx) let mut error = UndefinedValue()); stream.get_stored_error(error.handle_mut()); // Set writer.[[readyPromise]] to a promise rejected with stream.[[storedError]]. // Set writer.[[readyPromise]].[[PromiseIsHandled]] to true. // Note: new promise created in `new_inherited`. let ready_promise = self.ready_promise.borrow(); ready_promise.reject_native(cx, &error.handle()); ready_promise.set_promise_is_handled(cx); // Set writer.[[closedPromise]] to a new promise. // Done in `new_inherited`. return Ok(()); } // Otherwise, if state is "closed", if stream.is_closed() { // Set writer.[[readyPromise]] to a promise resolved with undefined. // Note: new promise created in `new_inherited`. self.ready_promise.borrow().resolve_native(cx, &()); // Set writer.[[closedPromise]] to a promise resolved with undefined. // Note: new promise created in `new_inherited`. self.closed_promise.borrow().resolve_native(cx, &()); return Ok(()); } // Otherwise, // Assert: state is "errored". assert!(stream.is_errored()); // Let storedError be stream.[[storedError]]. rooted!(&in(cx) let mut error = UndefinedValue()); stream.get_stored_error(error.handle_mut()); // Set writer.[[readyPromise]] to a promise rejected with stream.[[storedError]]. // Set writer.[[readyPromise]].[[PromiseIsHandled]] to true. // Note: new promise created in `new_inherited`. let ready_promise = self.ready_promise.borrow(); ready_promise.reject_native(cx, &error.handle()); ready_promise.set_promise_is_handled(cx); // Set writer.[[closedPromise]] to a promise rejected with storedError. // Set writer.[[closedPromise]].[[PromiseIsHandled]] to true. // Note: new promise created in `new_inherited`. let ready_promise = self.closed_promise.borrow(); ready_promise.reject_native(cx, &error.handle()); ready_promise.set_promise_is_handled(cx); Ok(()) } pub(crate) fn reject_closed_promise_with_stored_error( &self, cx: &mut JSContext, error: &SafeHandleValue, ) { self.closed_promise.borrow().reject_native(cx, error); } pub(crate) fn set_close_promise_is_handled(&self, cx: &mut JSContext) { self.closed_promise.borrow().set_promise_is_handled(cx); } pub(crate) fn set_ready_promise(&self, promise: Rc<Promise>) { *self.ready_promise.borrow_mut() = promise; } pub(crate) fn resolve_ready_promise_with_undefined(&self, cx: &mut JSContext) { self.ready_promise.borrow().resolve_native(cx, &()); } pub(crate) fn resolve_closed_promise_with_undefined(&self, cx: &mut JSContext) { self.closed_promise.borrow().resolve_native(cx, &()); } /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-ensure-ready-promise-rejected> pub(crate) fn ensure_ready_promise_rejected( &self, cx: &mut JSContext, global: &GlobalScope, error: SafeHandleValue, ) { let ready_promise = self.ready_promise.borrow().clone(); // If writer.[[readyPromise]].[[PromiseState]] is "pending", if ready_promise.is_pending() { // reject writer.[[readyPromise]] with error. ready_promise.reject_native(cx, &error); // Set writer.[[readyPromise]].[[PromiseIsHandled]] to true. ready_promise.set_promise_is_handled(cx); } else { // Otherwise, set writer.[[readyPromise]] to a promise rejected with error. let promise = Promise::new_rejected(cx, global, error); // Set writer.[[readyPromise]].[[PromiseIsHandled]] to true. promise.set_promise_is_handled(cx); *self.ready_promise.borrow_mut() = promise; } } /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-ensure-closed-promise-rejected> fn ensure_closed_promise_rejected( &self, cx: &mut JSContext, global: &GlobalScope, error: SafeHandleValue, ) { let closed_promise = self.closed_promise.borrow().clone(); // If writer.[[closedPromise]].[[PromiseState]] is "pending", if closed_promise.is_pending() { // reject writer.[[closedPromise]] with error. closed_promise.reject_native(cx, &error); // Set writer.[[closedPromise]].[[PromiseIsHandled]] to true. closed_promise.set_promise_is_handled(cx); } else { // Otherwise, set writer.[[closedPromise]] to a promise rejected with error. let promise = Promise::new_rejected(cx, global, error); // Set writer.[[closedPromise]].[[PromiseIsHandled]] to true. promise.set_promise_is_handled(cx); *self.closed_promise.borrow_mut() = promise; } } /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-abort> fn abort( &self, cx: &mut CurrentRealm, global: &GlobalScope, reason: SafeHandleValue, ) -> Rc<Promise> { // Let stream be writer.[[stream]]. let Some(stream) = self.stream.get() else { // Assert: stream is not undefined. unreachable!("Stream should be set."); }; // Return ! WritableStreamAbort(stream, reason). stream.abort(cx, global, reason) } /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-close> fn close(&self, cx: &mut JSContext, global: &GlobalScope) -> Rc<Promise> { // Let stream be writer.[[stream]]. let Some(stream) = self.stream.get() else { // Assert: stream is not undefined. unreachable!("Stream should be set."); }; // Return ! WritableStreamClose(stream). stream.close(cx, global) } /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-write> pub(crate) fn write( &self, cx: &mut JSContext, global: &GlobalScope, chunk: SafeHandleValue, ) -> Rc<Promise> { // Let stream be writer.[[stream]]. let Some(stream) = self.stream.get() else { // Assert: stream is not undefined. unreachable!("Stream should be set."); }; // Let controller be stream.[[controller]]. // Note: asserting controller is some. let Some(controller) = stream.get_controller() else { unreachable!("Controller should be set."); }; // Let chunkSize be ! WritableStreamDefaultControllerGetChunkSize(controller, chunk). let chunk_size = controller.get_chunk_size(cx, global, chunk); // If stream is not equal to writer.[[stream]], // return a promise rejected with a TypeError exception. if !self .stream .get() .is_some_and(|current_stream| current_stream == stream) { let promise = Promise::new(cx, global); promise.reject_error( cx, Error::Type(c"Stream is not equal to writer stream".to_owned()), ); return promise; } // Let state be stream.[[state]]. // If state is "errored", if stream.is_errored() { // return a promise rejected with stream.[[storedError]]. rooted!(&in(cx) let mut error = UndefinedValue()); stream.get_stored_error(error.handle_mut()); let promise = Promise::new(cx, global); promise.reject_native(cx, &error.handle()); return promise; } // If ! WritableStreamCloseQueuedOrInFlight(stream) is true // or state is "closed", if stream.close_queued_or_in_flight() || stream.is_closed() { // return a promise rejected with a TypeError exception // indicating that the stream is closing or closed let promise = Promise::new(cx, global); promise.reject_error( cx, Error::Type(c"Stream has been closed, or has close queued or in-flight".to_owned()), ); return promise; } // If state is "erroring", if stream.is_erroring() { // return a promise rejected with stream.[[storedError]]. rooted!(&in(cx) let mut error = UndefinedValue()); stream.get_stored_error(error.handle_mut()); let promise = Promise::new(cx, global); promise.reject_native(cx, &error.handle()); return promise; } // Assert: state is "writable". assert!(stream.is_writable()); // Let promise be ! WritableStreamAddWriteRequest(stream). let promise = stream.add_write_request(cx, global); // Perform ! WritableStreamDefaultControllerWrite(controller, chunk, chunkSize). controller.write(cx, global, chunk, chunk_size); // Return promise. promise } /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-release> pub(crate) fn release(&self, cx: &mut JSContext, global: &GlobalScope) { // Let stream be this.[[stream]]. let Some(stream) = self.stream.get() else { // Assert: stream is not undefined. unreachable!("Stream should be set."); }; // Assert: stream.[[writer]] is writer. assert!(stream.get_writer().is_some_and(|writer| &*writer == self)); // Let releasedError be a new TypeError. let released_error = Error::Type(c"Writer has been released".to_owned()); // Root the js val of the error. rooted!(&in(cx) let mut error = UndefinedValue()); released_error.to_jsval(cx, global, error.handle_mut()); // Perform ! WritableStreamDefaultWriterEnsureReadyPromiseRejected(writer, releasedError). self.ensure_ready_promise_rejected(cx, global, error.handle()); // Perform ! WritableStreamDefaultWriterEnsureClosedPromiseRejected(writer, releasedError). self.ensure_closed_promise_rejected(cx, global, error.handle()); // Set stream.[[writer]] to undefined. stream.set_writer(None); // Set this.[[stream]] to undefined. self.stream.set(None); } /// <https://streams.spec.whatwg.org/#writable-stream-default-writer-close-with-error-propagation> pub(crate) fn close_with_error_propagation( &self, cx: &mut JSContext, global: &GlobalScope, ) -> Rc<Promise> { // Let stream be writer.[[stream]]. let Some(stream) = self.stream.get() else { // Assert: stream is not undefined. unreachable!("Stream should be set."); }; // Let state be stream.[[state]]. // Used via stream method calls. // If ! WritableStreamCloseQueuedOrInFlight(stream) is true // or state is "closed", if stream.close_queued_or_in_flight() || stream.is_closed() { // return a promise resolved with undefined. let promise = Promise::new(cx, global); promise.resolve_native(cx, &()); return promise; } // If state is "errored", if stream.is_errored() { // return a promise rejected with stream.[[storedError]]. rooted!(&in(cx) let mut error = UndefinedValue()); stream.get_stored_error(error.handle_mut()); let promise = Promise::new(cx, global); promise.reject_native(cx, &error.handle()); return promise; } // Assert: state is "writable" or "erroring". assert!(stream.is_writable() || stream.is_erroring()); // Return ! WritableStreamDefaultWriterClose(writer). self.close(cx, global) } pub(crate) fn get_stream(&self) -> Option<DomRoot<WritableStream>> { self.stream.get() } } impl WritableStreamDefaultWriterMethods<crate::DomTypeHolder> for WritableStreamDefaultWriter { /// <https://streams.spec.whatwg.org/#default-writer-closed> fn Closed(&self) -> Rc<Promise> { // Return this.[[closedPromise]]. return self.closed_promise.borrow().clone(); } /// <https://streams.spec.whatwg.org/#default-writer-desired-size> fn GetDesiredSize(&self) -> Result<Option<f64>, Error> { // If this.[[stream]] is undefined, throw a TypeError exception. let Some(stream) = self.stream.get() else { return Err(Error::Type(c"Stream is undefined".to_owned())); }; // Return ! WritableStreamDefaultWriterGetDesiredSize(this). Ok(stream.get_desired_size()) } /// <https://streams.spec.whatwg.org/#default-writer-ready> fn Ready(&self) -> Rc<Promise> { // Return this.[[readyPromise]]. return self.ready_promise.borrow().clone(); } /// <https://streams.spec.whatwg.org/#default-writer-abort> fn Abort(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) -> Rc<Promise> { let global = GlobalScope::from_current_realm(cx); // If this.[[stream]] is undefined, if self.stream.get().is_none() { // return a promise rejected with a TypeError exception. let promise = Promise::new(cx, &global); promise.reject_error(cx, Error::Type(c"Stream is undefined".to_owned())); return promise; } // Return ! WritableStreamDefaultWriterAbort(this, reason). self.abort(cx, &global, reason) } /// <https://streams.spec.whatwg.org/#default-writer-close> fn Close(&self, cx: &mut CurrentRealm) -> Rc<Promise> { let global = GlobalScope::from_current_realm(cx); let promise = Promise::new(cx, &global); // Let stream be this.[[stream]]. let Some(stream) = self.stream.get() else { // If stream is undefined, // return a promise rejected with a TypeError exception. promise.reject_error(cx, Error::Type(c"Stream is undefined".to_owned())); return promise; }; // If ! WritableStreamCloseQueuedOrInFlight(stream) is true if stream.close_queued_or_in_flight() { // return a promise rejected with a TypeError exception. promise.reject_error( cx, Error::Type(c"Stream has closed queued or in-flight".to_owned()), ); return promise; } self.close(cx, &global) } /// <https://streams.spec.whatwg.org/#default-writer-release-lock> fn ReleaseLock(&self, cx: &mut JSContext) { // Let stream be this.[[stream]]. let Some(stream) = self.stream.get() else { // If stream is undefined, return. return; }; // Assert: stream.[[writer]] is not undefined. assert!(stream.get_writer().is_some()); let global = self.global(); // Perform ! WritableStreamDefaultWriterRelease(this). self.release(cx, &global); } /// <https://streams.spec.whatwg.org/#default-writer-write> fn Write(&self, cx: &mut CurrentRealm, chunk: SafeHandleValue) -> Rc<Promise> { let global = GlobalScope::from_current_realm(cx); // If this.[[stream]] is undefined, if self.stream.get().is_none() { // return a promise rejected with a TypeError exception. let promise = Promise::new(cx, &global); promise.reject_error(cx, Error::Type(c"Stream is undefined".to_owned())); return promise; } // Return ! WritableStreamDefaultWriterWrite(this, chunk). self.write(cx, &global, chunk) } /// <https://streams.spec.whatwg.org/#default-writer-constructor> fn Constructor( cx: &mut CurrentRealm, global: &GlobalScope, proto: Option<SafeHandleObject>, stream: &WritableStream, ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> { let writer = WritableStreamDefaultWriter::new(cx, global, proto); // Perform ? SetUpWritableStreamDefaultWriter(this, stream). writer.setup(cx, stream)?; Ok(writer) } }