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::conversions::ToJSValConvertible;
14use js::jsapi::{Heap, JSObject};
15use js::jsval::{JSVal, ObjectValue, UndefinedValue};
16use js::realm::CurrentRealm;
17use js::rust::{
18 HandleObject as SafeHandleObject, HandleValue as SafeHandleValue,
19 MutableHandleValue as SafeMutableHandleValue,
20};
21use rustc_hash::FxHashMap;
22use script_bindings::cell::DomRefCell;
23use script_bindings::codegen::GenericBindings::MessagePortBinding::MessagePortMethods;
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, RootedPromise, TracedPromise};
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 abort_request_promise: TracedPromise,
61}
62
63impl Callback for AbortAlgorithmFulfillmentHandler {
64 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
65 self.abort_request_promise.resolve_native(cx, &());
67
68 self.stream
70 .as_rooted()
71 .reject_close_and_closed_promise_if_needed(cx);
72 }
73}
74
75impl js::gc::Rootable for AbortAlgorithmRejectionHandler {}
76
77#[derive(JSTraceable, MallocSizeOf)]
80#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
81struct AbortAlgorithmRejectionHandler {
82 stream: Dom<WritableStream>,
83 abort_request_promise: TracedPromise,
84}
85
86impl Callback for AbortAlgorithmRejectionHandler {
87 fn callback(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) {
88 self.abort_request_promise.reject_native(cx, &reason);
90
91 self.stream
93 .as_rooted()
94 .reject_close_and_closed_promise_if_needed(cx);
95 }
96}
97
98impl js::gc::Rootable for PendingAbortRequest {}
99
100#[derive(JSTraceable, MallocSizeOf)]
102#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
103struct PendingAbortRequest {
104 promise: TracedPromise,
106
107 #[ignore_malloc_size_of = "mozjs"]
109 reason: Box<Heap<JSVal>>,
110
111 was_already_erroring: bool,
113}
114
115#[derive(Clone, Copy, Debug, Default, JSTraceable, MallocSizeOf)]
117pub(crate) enum WritableStreamState {
118 #[default]
119 Writable,
120 Closed,
121 Erroring,
122 Errored,
123}
124
125#[dom_struct]
127pub struct WritableStream {
128 reflector_: Reflector,
129
130 backpressure: Cell<bool>,
132
133 close_request: DomRefCell<Option<TracedPromise>>,
135
136 controller: MutNullableDom<WritableStreamDefaultController>,
138
139 detached: Cell<bool>,
141
142 in_flight_write_request: DomRefCell<Option<TracedPromise>>,
144
145 in_flight_close_request: DomRefCell<Option<TracedPromise>>,
147
148 pending_abort_request: DomRefCell<Option<PendingAbortRequest>>,
150
151 state: Cell<WritableStreamState>,
153
154 #[ignore_malloc_size_of = "mozjs"]
156 stored_error: Heap<JSVal>,
157
158 writer: MutNullableDom<WritableStreamDefaultWriter>,
160
161 write_requests: DomRefCell<VecDeque<TracedPromise>>,
163}
164
165impl WritableStream {
166 fn new_inherited() -> WritableStream {
168 WritableStream {
169 reflector_: Reflector::new(),
170 backpressure: Default::default(),
171 close_request: Default::default(),
172 controller: Default::default(),
173 detached: Default::default(),
174 in_flight_write_request: Default::default(),
175 in_flight_close_request: Default::default(),
176 pending_abort_request: Default::default(),
177 state: Default::default(),
178 stored_error: Default::default(),
179 writer: Default::default(),
180 write_requests: Default::default(),
181 }
182 }
183
184 pub(crate) fn new_with_proto(
185 cx: &mut JSContext,
186 global: &GlobalScope,
187 proto: Option<SafeHandleObject>,
188 ) -> DomRoot<WritableStream> {
189 reflect_dom_object_with_proto(cx, Box::new(WritableStream::new_inherited()), global, proto)
190 }
191
192 pub(crate) fn assert_no_controller(&self) {
195 assert!(self.controller.get().is_none());
196 }
197
198 pub(crate) fn set_default_controller(&self, controller: &WritableStreamDefaultController) {
201 self.controller.set(Some(controller));
202 }
203
204 pub(crate) fn get_default_controller(&self) -> DomRoot<WritableStreamDefaultController> {
205 self.controller.get().expect("Controller should be set.")
206 }
207
208 pub(crate) fn is_writable(&self) -> bool {
209 matches!(self.state.get(), WritableStreamState::Writable)
210 }
211
212 pub(crate) fn is_erroring(&self) -> bool {
213 matches!(self.state.get(), WritableStreamState::Erroring)
214 }
215
216 pub(crate) fn is_errored(&self) -> bool {
217 matches!(self.state.get(), WritableStreamState::Errored)
218 }
219
220 pub(crate) fn is_closed(&self) -> bool {
221 matches!(self.state.get(), WritableStreamState::Closed)
222 }
223
224 pub(crate) fn has_in_flight_write_request(&self) -> bool {
225 self.in_flight_write_request.borrow().is_some()
226 }
227
228 pub(crate) fn has_operations_marked_inflight(&self) -> bool {
230 let in_flight_write_requested = self.in_flight_write_request.borrow().is_some();
231 let in_flight_close_requested = self.in_flight_close_request.borrow().is_some();
232
233 in_flight_write_requested || in_flight_close_requested
234 }
235
236 pub(crate) fn get_stored_error(&self, mut handle_mut: SafeMutableHandleValue) {
238 handle_mut.set(self.stored_error.get());
239 }
240
241 pub(crate) fn finish_erroring(&self, cx: &mut JSContext, global: &GlobalScope) {
243 assert!(self.is_erroring());
245
246 assert!(!self.has_operations_marked_inflight());
248
249 self.state.set(WritableStreamState::Errored);
251
252 let Some(controller) = self.controller.get() else {
254 unreachable!("Stream should have a controller.");
255 };
256 controller.perform_error_steps();
257
258 rooted!(&in(cx) let mut stored_error = UndefinedValue());
260 self.get_stored_error(stored_error.handle_mut());
261
262 rooted!(&in(cx) let write_requests = mem::take(&mut *self.write_requests.borrow_mut()));
264 for request in write_requests.iter() {
265 request.reject(cx, stored_error.handle());
267 }
268
269 if self.pending_abort_request.borrow().is_none() {
274 self.reject_close_and_closed_promise_if_needed(cx);
276
277 return;
279 }
280
281 rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
284 if let Some(pending_abort_request) = &*pending_abort_request {
285 if pending_abort_request.was_already_erroring {
287 pending_abort_request
289 .promise
290 .reject(cx, stored_error.handle());
291
292 self.reject_close_and_closed_promise_if_needed(cx);
294
295 return;
297 }
298
299 rooted!(&in(cx) let mut reason = UndefinedValue());
301 reason.set(pending_abort_request.reason.get());
302 let promise = controller.abort_steps(cx, global, reason.handle());
303
304 rooted!(&in(cx) let mut fulfillment_handler = Some(AbortAlgorithmFulfillmentHandler {
306 stream: Dom::from_ref(self),
307 abort_request_promise: pending_abort_request.promise.clone(),
308 }));
309
310 rooted!(&in(cx) let mut rejection_handler = Some(AbortAlgorithmRejectionHandler {
312 stream: Dom::from_ref(self),
313 abort_request_promise: pending_abort_request.promise.clone(),
314 }));
315
316 let handler = PromiseNativeHandler::new(
317 cx,
318 global,
319 fulfillment_handler.take().map(|h| Box::new(h) as Box<_>),
320 rejection_handler.take().map(|h| Box::new(h) as Box<_>),
321 );
322
323 let mut realm = enter_auto_realm(cx, global);
324 let cx = &mut realm.current_realm();
325 promise.append_native_handler(cx, &handler);
326 }
327 }
328
329 fn reject_close_and_closed_promise_if_needed(&self, cx: &mut JSContext) {
331 assert!(self.is_errored());
333
334 rooted!(&in(cx) let mut stored_error = UndefinedValue());
335 self.get_stored_error(stored_error.handle_mut());
336
337 rooted!(&in(cx) let close_request = self.close_request.borrow_mut().take());
339 if let Some(ref close_request) = *close_request {
340 assert!(self.in_flight_close_request.borrow().is_none());
342
343 close_request.reject_native(cx, &stored_error.handle())
345
346 }
349
350 if let Some(writer) = self.writer.get() {
353 writer.reject_closed_promise_with_stored_error(cx, &stored_error.handle());
355
356 writer.set_close_promise_is_handled(cx);
358 }
359 }
360
361 pub(crate) fn close_queued_or_in_flight(&self) -> bool {
363 let close_requested = self.close_request.borrow().is_some();
364 let in_flight_close_requested = self.in_flight_close_request.borrow().is_some();
365
366 close_requested || in_flight_close_requested
367 }
368
369 pub(crate) fn finish_in_flight_write(&self, cx: &mut JSContext) {
371 rooted!(&in(cx) let in_flight_write_request = self.in_flight_write_request.borrow_mut().take());
372 let Some(ref in_flight_write_request) = *in_flight_write_request else {
373 unreachable!("Stream should have a write request");
375 };
376
377 in_flight_write_request.resolve_native(cx, &());
379
380 }
383
384 pub(crate) fn start_erroring(
386 &self,
387 cx: &mut JSContext,
388 global: &GlobalScope,
389 error: SafeHandleValue,
390 ) {
391 assert!(self.stored_error.get().is_undefined());
393
394 assert!(self.is_writable());
396
397 let Some(controller) = self.controller.get() else {
399 unreachable!("Stream should have a controller.");
401 };
402
403 self.state.set(WritableStreamState::Erroring);
405
406 self.stored_error.set(*error);
408
409 if let Some(writer) = self.writer.get() {
411 writer.ensure_ready_promise_rejected(cx, global, error);
413 }
414
415 if !self.has_operations_marked_inflight() && controller.started() {
417 self.finish_erroring(cx, global);
419 }
420 }
421
422 pub(crate) fn deal_with_rejection(
424 &self,
425 cx: &mut JSContext,
426 global: &GlobalScope,
427 error: SafeHandleValue,
428 ) {
429 if self.is_writable() {
433 self.start_erroring(cx, global, error);
435
436 return;
438 }
439
440 assert!(self.is_erroring());
442
443 self.finish_erroring(cx, global);
445 }
446
447 pub(crate) fn mark_first_write_request_in_flight(&self) {
449 let mut in_flight_write_request = self.in_flight_write_request.borrow_mut();
450 let mut write_requests = self.write_requests.borrow_mut();
451
452 assert!(in_flight_write_request.is_none());
454
455 assert!(!write_requests.is_empty());
457
458 *in_flight_write_request = write_requests.pop_front();
462 }
463
464 pub(crate) fn mark_close_request_in_flight(&self) {
466 let mut in_flight_close_request = self.in_flight_close_request.borrow_mut();
467 let mut close_request = self.close_request.borrow_mut();
468
469 assert!(in_flight_close_request.is_none());
471
472 assert!(close_request.is_some());
474
475 *in_flight_close_request = close_request.take();
479 }
480
481 pub(crate) fn finish_in_flight_close(&self, cx: &mut JSContext) {
483 rooted!(&in(cx) let in_flight_close_request = self.in_flight_close_request.borrow_mut().take());
484 let Some(ref in_flight_close_request) = *in_flight_close_request else {
485 unreachable!("in_flight_close_request must be Some");
487 };
488
489 in_flight_close_request.resolve_native(cx, &());
491
492 assert!(self.is_writable() || self.is_erroring());
497
498 if self.is_erroring() {
500 self.stored_error.set(UndefinedValue());
502
503 rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
505 if let Some(pending_abort_request) = &*pending_abort_request {
506 pending_abort_request.promise.resolve_native(cx, &());
508
509 }
512 }
513
514 self.state.set(WritableStreamState::Closed);
516
517 if let Some(writer) = self.writer.get() {
519 writer.resolve_closed_promise_with_undefined(cx);
522 }
523
524 assert!(self.pending_abort_request.borrow().is_none());
526
527 assert!(self.stored_error.get().is_undefined());
529 }
530
531 pub(crate) fn finish_in_flight_close_with_error(
533 &self,
534 cx: &mut JSContext,
535 global: &GlobalScope,
536 error: SafeHandleValue,
537 ) {
538 rooted!(&in(cx) let in_flight_close_request = self.in_flight_close_request.borrow_mut().take());
539 let Some(ref in_flight_close_request) = *in_flight_close_request else {
540 unreachable!("Inflight close request must be defined.");
542 };
543
544 in_flight_close_request.reject_native(cx, &error);
546
547 assert!(self.is_erroring() || self.is_writable());
552
553 rooted!(&in(cx) let pending_abort_request = self.pending_abort_request.borrow_mut().take());
555 if let Some(pending_abort_request) = &*pending_abort_request {
556 pending_abort_request.promise.reject_native(cx, &error);
558
559 }
562
563 self.deal_with_rejection(cx, global, error);
565 }
566
567 pub(crate) fn finish_in_flight_write_with_error(
569 &self,
570 cx: &mut JSContext,
571 global: &GlobalScope,
572 error: SafeHandleValue,
573 ) {
574 rooted!(&in(cx) let in_flight_write_request = self.in_flight_write_request.borrow_mut().take());
575 let Some(ref in_flight_write_request) = *in_flight_write_request else {
576 unreachable!("Inflight write request must be defined.");
578 };
579
580 in_flight_write_request.reject_native(cx, &error);
582
583 assert!(self.is_erroring() || self.is_writable());
588
589 self.deal_with_rejection(cx, global, error);
591 }
592
593 pub(crate) fn get_writer(&self) -> Option<DomRoot<WritableStreamDefaultWriter>> {
594 self.writer.get()
595 }
596
597 pub(crate) fn set_writer(&self, writer: Option<&WritableStreamDefaultWriter>) {
598 self.writer.set(writer);
599 }
600
601 pub(crate) fn set_backpressure(&self, backpressure: bool) {
602 self.backpressure.set(backpressure);
603 }
604
605 pub(crate) fn get_backpressure(&self) -> bool {
606 self.backpressure.get()
607 }
608
609 pub(crate) fn is_locked(&self) -> bool {
611 self.get_writer().is_some()
614 }
615
616 pub(crate) fn add_write_request(
618 &self,
619 cx: &mut JSContext,
620 global: &GlobalScope,
621 ) -> RootedPromise {
622 assert!(self.is_locked());
624
625 assert!(self.is_writable());
627
628 let promise = Promise::new_rooted(cx, global);
630
631 self.write_requests
633 .borrow_mut()
634 .push_back(promise.to_traced());
635
636 promise
638 }
639
640 pub(crate) fn get_controller(&self) -> Option<DomRoot<WritableStreamDefaultController>> {
642 self.controller.get()
643 }
644
645 pub(crate) fn abort(
647 &self,
648 cx: &mut CurrentRealm,
649 global: &GlobalScope,
650 provided_reason: SafeHandleValue,
651 ) -> RootedPromise {
652 if self.is_closed() || self.is_errored() {
654 return Promise::new_resolved_rooted(cx, global, ());
656 }
657
658 self.get_controller()
660 .expect("Stream must have a controller.")
661 .signal_abort(cx, provided_reason);
662
663 let state = self.state.get();
665
666 if matches!(
668 state,
669 WritableStreamState::Closed | WritableStreamState::Errored
670 ) {
671 return Promise::new_resolved_rooted(cx, global, ());
672 }
673
674 if self.pending_abort_request.borrow().is_some() {
676 return self
678 .pending_abort_request
679 .borrow()
680 .as_ref()
681 .expect("Pending abort request must be Some.")
682 .promise
683 .root(cx);
684 }
685
686 assert!(self.is_writable() || self.is_erroring());
688
689 let mut was_already_erroring = false;
691 rooted!(&in(cx) let undefined_reason = UndefinedValue());
692
693 let reason = if self.is_erroring() {
695 was_already_erroring = true;
697
698 undefined_reason.handle()
700 } else {
701 provided_reason
703 };
704
705 let promise = Promise::new_rooted(cx, global);
707
708 *self.pending_abort_request.borrow_mut() = Some(PendingAbortRequest {
713 promise: promise.to_traced(),
714 reason: Heap::boxed(reason.get()),
715 was_already_erroring,
716 });
717
718 if !was_already_erroring {
720 self.start_erroring(cx, global, reason);
722 }
723
724 promise
726 }
727
728 pub(crate) fn close(&self, cx: &mut JSContext, global: &GlobalScope) -> RootedPromise {
730 if self.is_closed() || self.is_errored() {
733 let promise = Promise::new_rooted(cx, global);
735 promise.reject_error(cx, Error::Type(c"Stream is closed or errored.".to_owned()));
736 return promise;
737 }
738
739 assert!(self.is_writable() || self.is_erroring());
741
742 assert!(!self.close_queued_or_in_flight());
744
745 let promise = Promise::new_rooted(cx, global);
747
748 *self.close_request.borrow_mut() = Some(promise.to_traced());
750
751 if let Some(writer) = self.writer.get() {
754 if self.get_backpressure() && self.is_writable() {
757 writer.resolve_ready_promise_with_undefined(cx);
759 }
760 }
761
762 let Some(controller) = self.controller.get() else {
764 unreachable!("Stream must have a controller.");
765 };
766 controller.close(cx, global);
767
768 promise
770 }
771
772 pub(crate) fn get_desired_size(&self) -> Option<f64> {
775 if self.is_errored() || self.is_erroring() {
781 return None;
782 }
783
784 if self.is_closed() {
786 return Some(0.);
787 }
788
789 let Some(controller) = self.controller.get() else {
790 unreachable!("Stream must have a controller.");
791 };
792 Some(controller.get_desired_size())
793 }
794
795 pub(crate) fn aquire_default_writer(
797 &self,
798 cx: &mut CurrentRealm,
799 global: &GlobalScope,
800 ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> {
801 let writer = WritableStreamDefaultWriter::new(cx, global, None);
803
804 writer.setup(cx, self)?;
806
807 Ok(writer)
809 }
810
811 pub(crate) fn update_backpressure(
813 &self,
814 cx: &mut JSContext,
815 backpressure: bool,
816 global: &GlobalScope,
817 ) {
818 self.is_writable();
820
821 assert!(!self.close_queued_or_in_flight());
823
824 let writer = self.get_writer();
826
827 if let Some(writer) = writer {
828 if backpressure != self.get_backpressure() {
830 if backpressure {
832 let promise = Promise::new_rooted(cx, global);
834 writer.set_ready_promise(&promise);
835 } else {
836 assert!(!backpressure);
839 writer.resolve_ready_promise_with_undefined(cx);
841 }
842 }
843 }
844
845 self.set_backpressure(backpressure);
847 }
848
849 pub(crate) fn setup_cross_realm_transform_writable(
851 &self,
852 cx: &mut JSContext,
853 port: &MessagePort,
854 ) {
855 let port_id = port.message_port_id();
856 let global = self.global();
857
858 let size_algorithm = extract_size_algorithm(cx, &QueuingStrategy::default());
864
865 let backpressure_promise = Promise::new_rooted(cx, &global);
869 rooted!(&in(cx) let backpressure_promise = RcHolder(Rc::new(RefCell::new(Some(backpressure_promise.to_traced())))));
870
871 let controller = WritableStreamDefaultController::new(
873 cx,
874 &global,
875 UnderlyingSinkType::Transfer {
876 backpressure_promise: backpressure_promise.0.clone(),
877 port: Dom::from_ref(port),
878 },
879 1.0,
880 size_algorithm,
881 );
882
883 rooted!(&in(cx) let cross_realm_transform_writable = CrossRealmTransformWritable {
886 controller: Dom::from_ref(&controller),
887 backpressure_promise: backpressure_promise.0.clone(),
888 });
889 global.note_cross_realm_transform_writable(&cross_realm_transform_writable, port_id);
890
891 port.Start(cx);
893
894 controller
896 .setup(cx, &global, self)
897 .expect("Setup for transfer cannot fail");
898 }
899 #[allow(clippy::too_many_arguments)]
901 fn setup_from_underlying_sink(
902 &self,
903 cx: &mut JSContext,
904 global: &GlobalScope,
905 stream: &WritableStream,
906 underlying_sink_obj: SafeHandleObject,
907 underlying_sink: &UnderlyingSink,
908 strategy_hwm: f64,
909 strategy_size: Rc<QueuingStrategySize>,
910 ) -> Result<(), Error> {
911 let controller = WritableStreamDefaultController::new(
937 cx,
938 global,
939 UnderlyingSinkType::new_js(
940 underlying_sink.abort.clone(),
941 underlying_sink.start.clone(),
942 underlying_sink.close.clone(),
943 underlying_sink.write.clone(),
944 ),
945 strategy_hwm,
946 strategy_size,
947 );
948
949 controller.set_underlying_sink_this_object(underlying_sink_obj);
952
953 controller.setup(cx, global, stream)
955 }
956}
957
958#[cfg_attr(crown, expect(crown::unrooted_must_root))]
960pub(crate) fn create_writable_stream(
961 cx: &mut JSContext,
962 global: &GlobalScope,
963 writable_high_water_mark: f64,
964 writable_size_algorithm: Rc<QueuingStrategySize>,
965 underlying_sink_type: UnderlyingSinkType,
966) -> Fallible<DomRoot<WritableStream>> {
967 assert!(writable_high_water_mark >= 0.0);
969
970 let stream = WritableStream::new_with_proto(cx, global, None);
973
974 let controller = WritableStreamDefaultController::new(
976 cx,
977 global,
978 underlying_sink_type,
979 writable_high_water_mark,
980 writable_size_algorithm,
981 );
982
983 controller.setup(cx, global, &stream)?;
986
987 Ok(stream)
989}
990
991impl WritableStreamMethods<crate::DomTypeHolder> for WritableStream {
992 fn Constructor(
994 cx: &mut JSContext,
995 global: &GlobalScope,
996 proto: Option<SafeHandleObject>,
997 underlying_sink: Option<*mut JSObject>,
998 strategy: &QueuingStrategy,
999 ) -> Fallible<DomRoot<WritableStream>> {
1000 rooted!(&in(cx) let underlying_sink_obj = underlying_sink.unwrap_or(ptr::null_mut()));
1002
1003 let underlying_sink_dict = if !underlying_sink_obj.is_null() {
1006 rooted!(&in(cx) let obj_val = ObjectValue(underlying_sink_obj.get()));
1007 match UnderlyingSink::new(cx, obj_val.handle()) {
1008 Ok(ConversionResult::Success(val)) => val,
1009 Ok(ConversionResult::Failure(error)) => {
1010 return Err(Error::Type(error.into_owned()));
1011 },
1012 _ => {
1013 return Err(Error::JSFailed);
1014 },
1015 }
1016 } else {
1017 UnderlyingSink::empty()
1018 };
1019
1020 if !underlying_sink_dict.type_.handle().is_undefined() {
1021 return Err(Error::Range(c"type is set".to_owned()));
1023 }
1024
1025 let stream = WritableStream::new_with_proto(cx, global, proto);
1027
1028 let size_algorithm = extract_size_algorithm(cx, strategy);
1030
1031 let high_water_mark = extract_high_water_mark(strategy, 1.0)?;
1033
1034 stream.setup_from_underlying_sink(
1037 cx,
1038 global,
1039 &stream,
1040 underlying_sink_obj.handle(),
1041 &underlying_sink_dict,
1042 high_water_mark,
1043 size_algorithm,
1044 )?;
1045
1046 Ok(stream)
1047 }
1048
1049 fn Locked(&self) -> bool {
1051 self.is_locked()
1053 }
1054
1055 fn Abort(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) -> RootedPromise {
1057 let global = GlobalScope::from_current_realm(cx);
1058
1059 if self.is_locked() {
1061 let promise = Promise::new_rooted(cx, &global);
1063 promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1064 return promise;
1065 }
1066
1067 self.abort(cx, &global, reason)
1069 }
1070
1071 fn Close(&self, cx: &mut CurrentRealm) -> RootedPromise {
1073 let global = GlobalScope::from_current_realm(cx);
1074
1075 if self.is_locked() {
1077 let promise = Promise::new_rooted(cx, &global);
1079 promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1080 return promise;
1081 }
1082
1083 if self.close_queued_or_in_flight() {
1085 let promise = Promise::new_rooted(cx, &global);
1087 promise.reject_error(
1088 cx,
1089 Error::Type(c"Stream has closed queued or in-flight".to_owned()),
1090 );
1091 return promise;
1092 }
1093
1094 self.close(cx, &global)
1096 }
1097
1098 fn GetWriter(
1100 &self,
1101 realm: &mut CurrentRealm,
1102 ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> {
1103 let global = GlobalScope::from_current_realm(realm);
1104
1105 self.aquire_default_writer(realm, &global)
1107 }
1108}
1109
1110impl js::gc::Rootable for CrossRealmTransformWritable {}
1111
1112#[derive(Clone, JSTraceable, MallocSizeOf)]
1116#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
1117pub(crate) struct CrossRealmTransformWritable {
1118 controller: Dom<WritableStreamDefaultController>,
1120
1121 #[conditional_malloc_size_of]
1123 backpressure_promise: Rc<RefCell<Option<TracedPromise>>>,
1124}
1125
1126impl CrossRealmTransformWritable {
1127 pub(crate) fn handle_message(
1130 &self,
1131 cx: &mut CurrentRealm,
1132 global: &GlobalScope,
1133 message: SafeHandleValue,
1134 ) {
1135 rooted!(&in(cx) let mut value = UndefinedValue());
1136 let type_string = get_type_and_value_from_message(cx, message, value.handle_mut());
1137
1138 if type_string == "error" {
1143 self.controller.error_if_needed(cx, value.handle(), global);
1145 }
1146
1147 rooted!(&in(cx) let backpressure_promise = self.backpressure_promise.borrow_mut().take());
1148
1149 if let Some(ref promise) = *backpressure_promise {
1152 promise.resolve_native(cx, &());
1154
1155 }
1158 }
1159
1160 pub(crate) fn handle_error(
1163 &self,
1164 cx: &mut CurrentRealm,
1165 global: &GlobalScope,
1166 port: &MessagePort,
1167 ) {
1168 let error = DOMException::new(cx, global, DOMErrorName::DataCloneError);
1170 rooted!(&in(cx) let mut rooted_error = UndefinedValue());
1171 error.to_jsval(cx, rooted_error.handle_mut());
1172
1173 port.cross_realm_transform_send_error(cx, rooted_error.handle());
1175
1176 self.controller
1178 .error_if_needed(cx, rooted_error.handle(), global);
1179
1180 global.disentangle_port(cx, port);
1182 }
1183}
1184
1185impl Transferable for WritableStream {
1187 type Index = MessagePortIndex;
1188 type Data = MessagePortImpl;
1189
1190 fn transfer(&self, cx: &mut JSContext) -> Fallible<(MessagePortId, MessagePortImpl)> {
1192 if self.is_locked() {
1195 return Err(Error::DataClone(None));
1196 }
1197
1198 let global = self.global();
1199 let mut realm = enter_auto_realm(cx, &*global);
1200 let mut realm = realm.current_realm();
1201 let cx = &mut realm;
1202
1203 let port_1 = MessagePort::new(cx, &global);
1205 global.track_message_port(&port_1, None);
1206
1207 let port_2 = MessagePort::new(cx, &global);
1209 global.track_message_port(&port_2, None);
1210
1211 global.entangle_ports(*port_1.message_port_id(), *port_2.message_port_id());
1213
1214 let readable = ReadableStream::new_with_proto(cx, &global, None);
1216
1217 readable.setup_cross_realm_transform_readable(cx, &port_1);
1219
1220 let promise = readable.pipe_to(cx, &global, self, false, false, false, None);
1222
1223 promise.set_promise_is_handled(cx);
1225
1226 port_2.transfer(cx)
1228 }
1229
1230 fn transfer_receive(
1232 cx: &mut JSContext,
1233 owner: &GlobalScope,
1234 id: MessagePortId,
1235 port_impl: MessagePortImpl,
1236 ) -> Result<DomRoot<Self>, ()> {
1237 let value = WritableStream::new_with_proto(cx, owner, None);
1240
1241 let transferred_port = MessagePort::transfer_receive(cx, owner, id, port_impl)?;
1248
1249 value.setup_cross_realm_transform_writable(cx, &transferred_port);
1251 Ok(value)
1252 }
1253
1254 fn serialized_storage<'a>(
1256 data: StructuredData<'a, '_>,
1257 ) -> &'a mut Option<FxHashMap<MessagePortId, Self::Data>> {
1258 match data {
1259 StructuredData::Reader(r) => &mut r.port_impls,
1260 StructuredData::Writer(w) => &mut w.ports,
1261 }
1262 }
1263}
1264
1265#[derive(JSTraceable)]
1266struct RcHolder<T>(Rc<T>);
1267
1268impl<T: js::rust::Traceable> js::gc::Rootable for RcHolder<T> {}