1use std::cell::{Cell, RefCell};
6use std::collections::VecDeque;
7use std::mem;
8use std::ptr::{self};
9use std::rc::Rc;
10
11use dom_struct::dom_struct;
12use js::context::JSContext;
13use js::jsapi::{Heap, JSObject};
14use js::jsval::{JSVal, ObjectValue, UndefinedValue};
15use js::realm::CurrentRealm;
16use js::rust::{
17 HandleObject as SafeHandleObject, HandleValue as SafeHandleValue,
18 MutableHandleValue as SafeMutableHandleValue,
19};
20use rustc_hash::FxHashMap;
21use script_bindings::cell::DomRefCell;
22use script_bindings::codegen::GenericBindings::MessagePortBinding::MessagePortMethods;
23use script_bindings::conversions::SafeToJSValConvertible;
24use script_bindings::reflector::{Reflector, reflect_dom_object_with_proto};
25use servo_base::id::{MessagePortId, MessagePortIndex};
26use servo_constellation_traits::MessagePortImpl;
27
28use crate::dom::bindings::codegen::Bindings::QueuingStrategyBinding::{
29 QueuingStrategy, QueuingStrategySize,
30};
31use crate::dom::bindings::codegen::Bindings::UnderlyingSinkBinding::UnderlyingSink;
32use crate::dom::bindings::codegen::Bindings::WritableStreamBinding::WritableStreamMethods;
33use crate::dom::bindings::conversions::ConversionResult;
34use crate::dom::bindings::error::{Error, Fallible};
35use crate::dom::bindings::reflector::DomGlobal;
36use crate::dom::bindings::root::{Dom, DomRoot, MutNullableDom};
37use crate::dom::bindings::structuredclone::StructuredData;
38use crate::dom::bindings::transferable::Transferable;
39use crate::dom::domexception::{DOMErrorName, DOMException};
40use crate::dom::globalscope::GlobalScope;
41use crate::dom::messageport::MessagePort;
42use crate::dom::promise::Promise;
43use crate::dom::promisenativehandler::{Callback, PromiseNativeHandler};
44use crate::dom::readablestream::{ReadableStream, get_type_and_value_from_message};
45use crate::dom::stream::countqueuingstrategy::{extract_high_water_mark, extract_size_algorithm};
46use crate::dom::stream::writablestreamdefaultcontroller::{
47 UnderlyingSinkType, WritableStreamDefaultController,
48};
49use crate::dom::stream::writablestreamdefaultwriter::WritableStreamDefaultWriter;
50use crate::realms::enter_auto_realm;
51
52impl js::gc::Rootable for AbortAlgorithmFulfillmentHandler {}
53
54#[derive(JSTraceable, MallocSizeOf)]
57#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
58struct AbortAlgorithmFulfillmentHandler {
59 stream: Dom<WritableStream>,
60 #[conditional_malloc_size_of]
61 abort_request_promise: Rc<Promise>,
62}
63
64impl Callback for AbortAlgorithmFulfillmentHandler {
65 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
66 self.abort_request_promise.resolve_native(cx, &());
68
69 self.stream
71 .as_rooted()
72 .reject_close_and_closed_promise_if_needed(cx);
73 }
74}
75
76impl js::gc::Rootable for AbortAlgorithmRejectionHandler {}
77
78#[derive(JSTraceable, MallocSizeOf)]
81#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
82struct AbortAlgorithmRejectionHandler {
83 stream: Dom<WritableStream>,
84 #[conditional_malloc_size_of]
85 abort_request_promise: Rc<Promise>,
86}
87
88impl Callback for AbortAlgorithmRejectionHandler {
89 fn callback(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) {
90 self.abort_request_promise.reject_native(cx, &reason);
92
93 self.stream
95 .as_rooted()
96 .reject_close_and_closed_promise_if_needed(cx);
97 }
98}
99
100impl js::gc::Rootable for PendingAbortRequest {}
101
102#[derive(JSTraceable, MallocSizeOf)]
104#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
105struct PendingAbortRequest {
106 #[conditional_malloc_size_of]
108 promise: Rc<Promise>,
109
110 #[ignore_malloc_size_of = "mozjs"]
112 reason: Box<Heap<JSVal>>,
113
114 was_already_erroring: bool,
116}
117
118#[derive(Clone, Copy, Debug, Default, JSTraceable, MallocSizeOf)]
120pub(crate) enum WritableStreamState {
121 #[default]
122 Writable,
123 Closed,
124 Erroring,
125 Errored,
126}
127
128#[dom_struct]
130pub struct WritableStream {
131 reflector_: Reflector,
132
133 backpressure: Cell<bool>,
135
136 #[conditional_malloc_size_of]
138 close_request: DomRefCell<Option<Rc<Promise>>>,
139
140 controller: MutNullableDom<WritableStreamDefaultController>,
142
143 detached: Cell<bool>,
145
146 #[conditional_malloc_size_of]
148 in_flight_write_request: DomRefCell<Option<Rc<Promise>>>,
149
150 #[conditional_malloc_size_of]
152 in_flight_close_request: DomRefCell<Option<Rc<Promise>>>,
153
154 pending_abort_request: DomRefCell<Option<PendingAbortRequest>>,
156
157 state: Cell<WritableStreamState>,
159
160 #[ignore_malloc_size_of = "mozjs"]
162 stored_error: Heap<JSVal>,
163
164 writer: MutNullableDom<WritableStreamDefaultWriter>,
166
167 #[conditional_malloc_size_of]
169 write_requests: DomRefCell<VecDeque<Rc<Promise>>>,
170}
171
172impl WritableStream {
173 fn new_inherited() -> WritableStream {
175 WritableStream {
176 reflector_: Reflector::new(),
177 backpressure: Default::default(),
178 close_request: Default::default(),
179 controller: Default::default(),
180 detached: Default::default(),
181 in_flight_write_request: Default::default(),
182 in_flight_close_request: Default::default(),
183 pending_abort_request: Default::default(),
184 state: Default::default(),
185 stored_error: Default::default(),
186 writer: Default::default(),
187 write_requests: Default::default(),
188 }
189 }
190
191 pub(crate) fn new_with_proto(
192 cx: &mut JSContext,
193 global: &GlobalScope,
194 proto: Option<SafeHandleObject>,
195 ) -> DomRoot<WritableStream> {
196 reflect_dom_object_with_proto(cx, Box::new(WritableStream::new_inherited()), global, proto)
197 }
198
199 pub(crate) fn assert_no_controller(&self) {
202 assert!(self.controller.get().is_none());
203 }
204
205 pub(crate) fn set_default_controller(&self, controller: &WritableStreamDefaultController) {
208 self.controller.set(Some(controller));
209 }
210
211 pub(crate) fn get_default_controller(&self) -> DomRoot<WritableStreamDefaultController> {
212 self.controller.get().expect("Controller should be set.")
213 }
214
215 pub(crate) fn is_writable(&self) -> bool {
216 matches!(self.state.get(), WritableStreamState::Writable)
217 }
218
219 pub(crate) fn is_erroring(&self) -> bool {
220 matches!(self.state.get(), WritableStreamState::Erroring)
221 }
222
223 pub(crate) fn is_errored(&self) -> bool {
224 matches!(self.state.get(), WritableStreamState::Errored)
225 }
226
227 pub(crate) fn is_closed(&self) -> bool {
228 matches!(self.state.get(), WritableStreamState::Closed)
229 }
230
231 pub(crate) fn has_in_flight_write_request(&self) -> bool {
232 self.in_flight_write_request.borrow().is_some()
233 }
234
235 pub(crate) fn has_operations_marked_inflight(&self) -> bool {
237 let in_flight_write_requested = self.in_flight_write_request.borrow().is_some();
238 let in_flight_close_requested = self.in_flight_close_request.borrow().is_some();
239
240 in_flight_write_requested || in_flight_close_requested
241 }
242
243 pub(crate) fn get_stored_error(&self, mut handle_mut: SafeMutableHandleValue) {
245 handle_mut.set(self.stored_error.get());
246 }
247
248 pub(crate) fn finish_erroring(&self, cx: &mut JSContext, global: &GlobalScope) {
250 assert!(self.is_erroring());
252
253 assert!(!self.has_operations_marked_inflight());
255
256 self.state.set(WritableStreamState::Errored);
258
259 let Some(controller) = self.controller.get() else {
261 unreachable!("Stream should have a controller.");
262 };
263 controller.perform_error_steps();
264
265 rooted!(&in(cx) let mut stored_error = UndefinedValue());
267 self.get_stored_error(stored_error.handle_mut());
268
269 let write_requests = mem::take(&mut *self.write_requests.borrow_mut());
271 for request in write_requests {
272 request.reject(cx, stored_error.handle());
274 }
275
276 if self.pending_abort_request.borrow().is_none() {
281 self.reject_close_and_closed_promise_if_needed(cx);
283
284 return;
286 }
287
288 rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
291 if let Some(pending_abort_request) = &*pending_abort_request {
292 if pending_abort_request.was_already_erroring {
294 pending_abort_request
296 .promise
297 .reject(cx, stored_error.handle());
298
299 self.reject_close_and_closed_promise_if_needed(cx);
301
302 return;
304 }
305
306 rooted!(&in(cx) let mut reason = UndefinedValue());
308 reason.set(pending_abort_request.reason.get());
309 let promise = controller.abort_steps(cx, global, reason.handle());
310
311 rooted!(&in(cx) let mut fulfillment_handler = Some(AbortAlgorithmFulfillmentHandler {
313 stream: Dom::from_ref(self),
314 abort_request_promise: pending_abort_request.promise.clone(),
315 }));
316
317 rooted!(&in(cx) let mut rejection_handler = Some(AbortAlgorithmRejectionHandler {
319 stream: Dom::from_ref(self),
320 abort_request_promise: pending_abort_request.promise.clone(),
321 }));
322
323 let handler = PromiseNativeHandler::new(
324 cx,
325 global,
326 fulfillment_handler.take().map(|h| Box::new(h) as Box<_>),
327 rejection_handler.take().map(|h| Box::new(h) as Box<_>),
328 );
329
330 let mut realm = enter_auto_realm(cx, global);
331 let cx = &mut realm.current_realm();
332 promise.append_native_handler(cx, &handler);
333 }
334 }
335
336 fn reject_close_and_closed_promise_if_needed(&self, cx: &mut JSContext) {
338 assert!(self.is_errored());
340
341 rooted!(&in(cx) let mut stored_error = UndefinedValue());
342 self.get_stored_error(stored_error.handle_mut());
343
344 let close_request = self.close_request.borrow_mut().take();
346 if let Some(close_request) = close_request {
347 assert!(self.in_flight_close_request.borrow().is_none());
349
350 close_request.reject_native(cx, &stored_error.handle())
352
353 }
356
357 if let Some(writer) = self.writer.get() {
360 writer.reject_closed_promise_with_stored_error(cx, &stored_error.handle());
362
363 writer.set_close_promise_is_handled(cx);
365 }
366 }
367
368 pub(crate) fn close_queued_or_in_flight(&self) -> bool {
370 let close_requested = self.close_request.borrow().is_some();
371 let in_flight_close_requested = self.in_flight_close_request.borrow().is_some();
372
373 close_requested || in_flight_close_requested
374 }
375
376 pub(crate) fn finish_in_flight_write(&self, cx: &mut JSContext) {
378 let Some(in_flight_write_request) = self.in_flight_write_request.borrow_mut().take() else {
379 unreachable!("Stream should have a write request");
381 };
382
383 in_flight_write_request.resolve_native(cx, &());
385
386 }
389
390 pub(crate) fn start_erroring(
392 &self,
393 cx: &mut JSContext,
394 global: &GlobalScope,
395 error: SafeHandleValue,
396 ) {
397 assert!(self.stored_error.get().is_undefined());
399
400 assert!(self.is_writable());
402
403 let Some(controller) = self.controller.get() else {
405 unreachable!("Stream should have a controller.");
407 };
408
409 self.state.set(WritableStreamState::Erroring);
411
412 self.stored_error.set(*error);
414
415 if let Some(writer) = self.writer.get() {
417 writer.ensure_ready_promise_rejected(cx, global, error);
419 }
420
421 if !self.has_operations_marked_inflight() && controller.started() {
423 self.finish_erroring(cx, global);
425 }
426 }
427
428 pub(crate) fn deal_with_rejection(
430 &self,
431 cx: &mut JSContext,
432 global: &GlobalScope,
433 error: SafeHandleValue,
434 ) {
435 if self.is_writable() {
439 self.start_erroring(cx, global, error);
441
442 return;
444 }
445
446 assert!(self.is_erroring());
448
449 self.finish_erroring(cx, global);
451 }
452
453 pub(crate) fn mark_first_write_request_in_flight(&self) {
455 let mut in_flight_write_request = self.in_flight_write_request.borrow_mut();
456 let mut write_requests = self.write_requests.borrow_mut();
457
458 assert!(in_flight_write_request.is_none());
460
461 assert!(!write_requests.is_empty());
463
464 let write_request = write_requests.pop_front().unwrap();
467
468 *in_flight_write_request = Some(write_request);
470 }
471
472 pub(crate) fn mark_close_request_in_flight(&self) {
474 let mut in_flight_close_request = self.in_flight_close_request.borrow_mut();
475 let mut close_request = self.close_request.borrow_mut();
476
477 assert!(in_flight_close_request.is_none());
479
480 assert!(close_request.is_some());
482
483 let close_request = close_request.take().unwrap();
486
487 *in_flight_close_request = Some(close_request);
489 }
490
491 pub(crate) fn finish_in_flight_close(&self, cx: &mut JSContext) {
493 let Some(in_flight_close_request) = self.in_flight_close_request.borrow_mut().take() else {
494 unreachable!("in_flight_close_request must be Some");
496 };
497
498 in_flight_close_request.resolve_native(cx, &());
500
501 assert!(self.is_writable() || self.is_erroring());
506
507 if self.is_erroring() {
509 self.stored_error.set(UndefinedValue());
511
512 rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
514 if let Some(pending_abort_request) = &*pending_abort_request {
515 pending_abort_request.promise.resolve_native(cx, &());
517
518 }
521 }
522
523 self.state.set(WritableStreamState::Closed);
525
526 if let Some(writer) = self.writer.get() {
528 writer.resolve_closed_promise_with_undefined(cx);
531 }
532
533 assert!(self.pending_abort_request.borrow().is_none());
535
536 assert!(self.stored_error.get().is_undefined());
538 }
539
540 pub(crate) fn finish_in_flight_close_with_error(
542 &self,
543 cx: &mut JSContext,
544 global: &GlobalScope,
545 error: SafeHandleValue,
546 ) {
547 let Some(in_flight_close_request) = self.in_flight_close_request.borrow_mut().take() else {
548 unreachable!("Inflight close request must be defined.");
550 };
551
552 in_flight_close_request.reject_native(cx, &error);
554
555 assert!(self.is_erroring() || self.is_writable());
560
561 rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
563 if let Some(pending_abort_request) = &*pending_abort_request {
564 pending_abort_request.promise.reject_native(cx, &error);
566
567 }
570
571 self.deal_with_rejection(cx, global, error);
573 }
574
575 pub(crate) fn finish_in_flight_write_with_error(
577 &self,
578 cx: &mut JSContext,
579 global: &GlobalScope,
580 error: SafeHandleValue,
581 ) {
582 let Some(in_flight_write_request) = self.in_flight_write_request.borrow_mut().take() else {
583 unreachable!("Inflight write request must be defined.");
585 };
586
587 in_flight_write_request.reject_native(cx, &error);
589
590 assert!(self.is_erroring() || self.is_writable());
595
596 self.deal_with_rejection(cx, global, error);
598 }
599
600 pub(crate) fn get_writer(&self) -> Option<DomRoot<WritableStreamDefaultWriter>> {
601 self.writer.get()
602 }
603
604 pub(crate) fn set_writer(&self, writer: Option<&WritableStreamDefaultWriter>) {
605 self.writer.set(writer);
606 }
607
608 pub(crate) fn set_backpressure(&self, backpressure: bool) {
609 self.backpressure.set(backpressure);
610 }
611
612 pub(crate) fn get_backpressure(&self) -> bool {
613 self.backpressure.get()
614 }
615
616 pub(crate) fn is_locked(&self) -> bool {
618 self.get_writer().is_some()
621 }
622
623 pub(crate) fn add_write_request(
625 &self,
626 cx: &mut JSContext,
627 global: &GlobalScope,
628 ) -> Rc<Promise> {
629 assert!(self.is_locked());
631
632 assert!(self.is_writable());
634
635 let promise = Promise::new(cx, global);
637
638 self.write_requests.borrow_mut().push_back(promise.clone());
640
641 promise
643 }
644
645 pub(crate) fn get_controller(&self) -> Option<DomRoot<WritableStreamDefaultController>> {
647 self.controller.get()
648 }
649
650 pub(crate) fn abort(
652 &self,
653 cx: &mut CurrentRealm,
654 global: &GlobalScope,
655 provided_reason: SafeHandleValue,
656 ) -> Rc<Promise> {
657 if self.is_closed() || self.is_errored() {
659 return Promise::new_resolved(cx, global, ());
661 }
662
663 self.get_controller()
665 .expect("Stream must have a controller.")
666 .signal_abort(cx, provided_reason);
667
668 let state = self.state.get();
670
671 if matches!(
673 state,
674 WritableStreamState::Closed | WritableStreamState::Errored
675 ) {
676 return Promise::new_resolved(cx, global, ());
677 }
678
679 if self.pending_abort_request.borrow().is_some() {
681 return self
683 .pending_abort_request
684 .borrow()
685 .as_ref()
686 .expect("Pending abort request must be Some.")
687 .promise
688 .clone();
689 }
690
691 assert!(self.is_writable() || self.is_erroring());
693
694 let mut was_already_erroring = false;
696 rooted!(&in(cx) let undefined_reason = UndefinedValue());
697
698 let reason = if self.is_erroring() {
700 was_already_erroring = true;
702
703 undefined_reason.handle()
705 } else {
706 provided_reason
708 };
709
710 let promise = Promise::new(cx, global);
712
713 *self.pending_abort_request.borrow_mut() = Some(PendingAbortRequest {
718 promise: promise.clone(),
719 reason: Heap::boxed(reason.get()),
720 was_already_erroring,
721 });
722
723 if !was_already_erroring {
725 self.start_erroring(cx, global, reason);
727 }
728
729 promise
731 }
732
733 pub(crate) fn close(&self, cx: &mut JSContext, global: &GlobalScope) -> Rc<Promise> {
735 if self.is_closed() || self.is_errored() {
738 let promise = Promise::new(cx, global);
740 promise.reject_error(cx, Error::Type(c"Stream is closed or errored.".to_owned()));
741 return promise;
742 }
743
744 assert!(self.is_writable() || self.is_erroring());
746
747 assert!(!self.close_queued_or_in_flight());
749
750 let promise = Promise::new(cx, global);
752
753 *self.close_request.borrow_mut() = Some(promise.clone());
755
756 if let Some(writer) = self.writer.get() {
759 if self.get_backpressure() && self.is_writable() {
762 writer.resolve_ready_promise_with_undefined(cx);
764 }
765 }
766
767 let Some(controller) = self.controller.get() else {
769 unreachable!("Stream must have a controller.");
770 };
771 controller.close(cx, global);
772
773 promise
775 }
776
777 pub(crate) fn get_desired_size(&self) -> Option<f64> {
780 if self.is_errored() || self.is_erroring() {
786 return None;
787 }
788
789 if self.is_closed() {
791 return Some(0.);
792 }
793
794 let Some(controller) = self.controller.get() else {
795 unreachable!("Stream must have a controller.");
796 };
797 Some(controller.get_desired_size())
798 }
799
800 pub(crate) fn aquire_default_writer(
802 &self,
803 cx: &mut CurrentRealm,
804 global: &GlobalScope,
805 ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> {
806 let writer = WritableStreamDefaultWriter::new(cx, global, None);
808
809 writer.setup(cx, self)?;
811
812 Ok(writer)
814 }
815
816 pub(crate) fn update_backpressure(
818 &self,
819 cx: &mut JSContext,
820 backpressure: bool,
821 global: &GlobalScope,
822 ) {
823 self.is_writable();
825
826 assert!(!self.close_queued_or_in_flight());
828
829 let writer = self.get_writer();
831
832 if let Some(writer) = writer {
833 if backpressure != self.get_backpressure() {
835 if backpressure {
837 let promise = Promise::new(cx, global);
839 writer.set_ready_promise(promise);
840 } else {
841 assert!(!backpressure);
844 writer.resolve_ready_promise_with_undefined(cx);
846 }
847 }
848 }
849
850 self.set_backpressure(backpressure);
852 }
853
854 pub(crate) fn setup_cross_realm_transform_writable(
856 &self,
857 cx: &mut JSContext,
858 port: &MessagePort,
859 ) {
860 let port_id = port.message_port_id();
861 let global = self.global();
862
863 let size_algorithm = extract_size_algorithm(cx, &QueuingStrategy::default());
869
870 let backpressure_promise = Rc::new(RefCell::new(Some(Promise::new(cx, &global))));
874
875 let controller = WritableStreamDefaultController::new(
877 cx,
878 &global,
879 UnderlyingSinkType::Transfer {
880 backpressure_promise: backpressure_promise.clone(),
881 port: Dom::from_ref(port),
882 },
883 1.0,
884 size_algorithm,
885 );
886
887 rooted!(&in(cx) let cross_realm_transform_writable = CrossRealmTransformWritable {
890 controller: Dom::from_ref(&controller),
891 backpressure_promise,
892 });
893 global.note_cross_realm_transform_writable(&cross_realm_transform_writable, port_id);
894
895 port.Start(cx);
897
898 controller
900 .setup(cx, &global, self)
901 .expect("Setup for transfer cannot fail");
902 }
903 #[allow(clippy::too_many_arguments)]
905 fn setup_from_underlying_sink(
906 &self,
907 cx: &mut JSContext,
908 global: &GlobalScope,
909 stream: &WritableStream,
910 underlying_sink_obj: SafeHandleObject,
911 underlying_sink: &UnderlyingSink,
912 strategy_hwm: f64,
913 strategy_size: Rc<QueuingStrategySize>,
914 ) -> Result<(), Error> {
915 let controller = WritableStreamDefaultController::new(
941 cx,
942 global,
943 UnderlyingSinkType::new_js(
944 underlying_sink.abort.clone(),
945 underlying_sink.start.clone(),
946 underlying_sink.close.clone(),
947 underlying_sink.write.clone(),
948 ),
949 strategy_hwm,
950 strategy_size,
951 );
952
953 controller.set_underlying_sink_this_object(underlying_sink_obj);
956
957 controller.setup(cx, global, stream)
959 }
960}
961
962#[cfg_attr(crown, expect(crown::unrooted_must_root))]
964pub(crate) fn create_writable_stream(
965 cx: &mut JSContext,
966 global: &GlobalScope,
967 writable_high_water_mark: f64,
968 writable_size_algorithm: Rc<QueuingStrategySize>,
969 underlying_sink_type: UnderlyingSinkType,
970) -> Fallible<DomRoot<WritableStream>> {
971 assert!(writable_high_water_mark >= 0.0);
973
974 let stream = WritableStream::new_with_proto(cx, global, None);
977
978 let controller = WritableStreamDefaultController::new(
980 cx,
981 global,
982 underlying_sink_type,
983 writable_high_water_mark,
984 writable_size_algorithm,
985 );
986
987 controller.setup(cx, global, &stream)?;
990
991 Ok(stream)
993}
994
995impl WritableStreamMethods<crate::DomTypeHolder> for WritableStream {
996 fn Constructor(
998 cx: &mut JSContext,
999 global: &GlobalScope,
1000 proto: Option<SafeHandleObject>,
1001 underlying_sink: Option<*mut JSObject>,
1002 strategy: &QueuingStrategy,
1003 ) -> Fallible<DomRoot<WritableStream>> {
1004 rooted!(&in(cx) let underlying_sink_obj = underlying_sink.unwrap_or(ptr::null_mut()));
1006
1007 let underlying_sink_dict = if !underlying_sink_obj.is_null() {
1010 rooted!(&in(cx) let obj_val = ObjectValue(underlying_sink_obj.get()));
1011 match UnderlyingSink::new(cx, obj_val.handle()) {
1012 Ok(ConversionResult::Success(val)) => val,
1013 Ok(ConversionResult::Failure(error)) => {
1014 return Err(Error::Type(error.into_owned()));
1015 },
1016 _ => {
1017 return Err(Error::JSFailed);
1018 },
1019 }
1020 } else {
1021 UnderlyingSink::empty()
1022 };
1023
1024 if !underlying_sink_dict.type_.handle().is_undefined() {
1025 return Err(Error::Range(c"type is set".to_owned()));
1027 }
1028
1029 let stream = WritableStream::new_with_proto(cx, global, proto);
1031
1032 let size_algorithm = extract_size_algorithm(cx, strategy);
1034
1035 let high_water_mark = extract_high_water_mark(strategy, 1.0)?;
1037
1038 stream.setup_from_underlying_sink(
1041 cx,
1042 global,
1043 &stream,
1044 underlying_sink_obj.handle(),
1045 &underlying_sink_dict,
1046 high_water_mark,
1047 size_algorithm,
1048 )?;
1049
1050 Ok(stream)
1051 }
1052
1053 fn Locked(&self) -> bool {
1055 self.is_locked()
1057 }
1058
1059 fn Abort(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) -> Rc<Promise> {
1061 let global = GlobalScope::from_current_realm(cx);
1062
1063 if self.is_locked() {
1065 let promise = Promise::new(cx, &global);
1067 promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1068 return promise;
1069 }
1070
1071 self.abort(cx, &global, reason)
1073 }
1074
1075 fn Close(&self, cx: &mut CurrentRealm) -> Rc<Promise> {
1077 let global = GlobalScope::from_current_realm(cx);
1078
1079 if self.is_locked() {
1081 let promise = Promise::new(cx, &global);
1083 promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1084 return promise;
1085 }
1086
1087 if self.close_queued_or_in_flight() {
1089 let promise = Promise::new(cx, &global);
1091 promise.reject_error(
1092 cx,
1093 Error::Type(c"Stream has closed queued or in-flight".to_owned()),
1094 );
1095 return promise;
1096 }
1097
1098 self.close(cx, &global)
1100 }
1101
1102 fn GetWriter(
1104 &self,
1105 realm: &mut CurrentRealm,
1106 ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> {
1107 let global = GlobalScope::from_current_realm(realm);
1108
1109 self.aquire_default_writer(realm, &global)
1111 }
1112}
1113
1114impl js::gc::Rootable for CrossRealmTransformWritable {}
1115
1116#[derive(Clone, JSTraceable, MallocSizeOf)]
1120#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
1121pub(crate) struct CrossRealmTransformWritable {
1122 controller: Dom<WritableStreamDefaultController>,
1124
1125 #[ignore_malloc_size_of = "nested Rc"]
1127 backpressure_promise: Rc<RefCell<Option<Rc<Promise>>>>,
1128}
1129
1130impl CrossRealmTransformWritable {
1131 pub(crate) fn handle_message(
1134 &self,
1135 cx: &mut CurrentRealm,
1136 global: &GlobalScope,
1137 message: SafeHandleValue,
1138 ) {
1139 rooted!(&in(cx) let mut value = UndefinedValue());
1140 let type_string = get_type_and_value_from_message(cx, message, value.handle_mut());
1141
1142 if type_string == "error" {
1147 self.controller.error_if_needed(cx, value.handle(), global);
1149 }
1150
1151 let backpressure_promise = self.backpressure_promise.borrow_mut().take();
1152
1153 if let Some(promise) = backpressure_promise {
1156 promise.resolve_native(cx, &());
1158
1159 }
1162 }
1163
1164 pub(crate) fn handle_error(
1167 &self,
1168 cx: &mut CurrentRealm,
1169 global: &GlobalScope,
1170 port: &MessagePort,
1171 ) {
1172 let error = DOMException::new(cx, global, DOMErrorName::DataCloneError);
1174 rooted!(&in(cx) let mut rooted_error = UndefinedValue());
1175 error.safe_to_jsval(cx, rooted_error.handle_mut());
1176
1177 port.cross_realm_transform_send_error(cx, rooted_error.handle());
1179
1180 self.controller
1182 .error_if_needed(cx, rooted_error.handle(), global);
1183
1184 global.disentangle_port(cx, port);
1186 }
1187}
1188
1189impl Transferable for WritableStream {
1191 type Index = MessagePortIndex;
1192 type Data = MessagePortImpl;
1193
1194 fn transfer(&self, cx: &mut JSContext) -> Fallible<(MessagePortId, MessagePortImpl)> {
1196 if self.is_locked() {
1199 return Err(Error::DataClone(None));
1200 }
1201
1202 let global = self.global();
1203 let mut realm = enter_auto_realm(cx, &*global);
1204 let mut realm = realm.current_realm();
1205 let cx = &mut realm;
1206
1207 let port_1 = MessagePort::new(cx, &global);
1209 global.track_message_port(&port_1, None);
1210
1211 let port_2 = MessagePort::new(cx, &global);
1213 global.track_message_port(&port_2, None);
1214
1215 global.entangle_ports(*port_1.message_port_id(), *port_2.message_port_id());
1217
1218 let readable = ReadableStream::new_with_proto(cx, &global, None);
1220
1221 readable.setup_cross_realm_transform_readable(cx, &port_1);
1223
1224 let promise = readable.pipe_to(cx, &global, self, false, false, false, None);
1226
1227 promise.set_promise_is_handled(cx);
1229
1230 port_2.transfer(cx)
1232 }
1233
1234 fn transfer_receive(
1236 cx: &mut JSContext,
1237 owner: &GlobalScope,
1238 id: MessagePortId,
1239 port_impl: MessagePortImpl,
1240 ) -> Result<DomRoot<Self>, ()> {
1241 let value = WritableStream::new_with_proto(cx, owner, None);
1244
1245 let transferred_port = MessagePort::transfer_receive(cx, owner, id, port_impl)?;
1252
1253 value.setup_cross_realm_transform_writable(cx, &transferred_port);
1255 Ok(value)
1256 }
1257
1258 fn serialized_storage<'a>(
1260 data: StructuredData<'a, '_>,
1261 ) -> &'a mut Option<FxHashMap<MessagePortId, Self::Data>> {
1262 match data {
1263 StructuredData::Reader(r) => &mut r.port_impls,
1264 StructuredData::Writer(w) => &mut w.ports,
1265 }
1266 }
1267}