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();
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(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 = Rc::new(RefCell::new(Some(Promise::new(cx, &global))));
869
870 let controller = WritableStreamDefaultController::new(
872 cx,
873 &global,
874 UnderlyingSinkType::Transfer {
875 backpressure_promise: backpressure_promise.clone(),
876 port: Dom::from_ref(port),
877 },
878 1.0,
879 size_algorithm,
880 );
881
882 rooted!(&in(cx) let cross_realm_transform_writable = CrossRealmTransformWritable {
885 controller: Dom::from_ref(&controller),
886 backpressure_promise,
887 });
888 global.note_cross_realm_transform_writable(&cross_realm_transform_writable, port_id);
889
890 port.Start(cx);
892
893 controller
895 .setup(cx, &global, self)
896 .expect("Setup for transfer cannot fail");
897 }
898 #[allow(clippy::too_many_arguments)]
900 fn setup_from_underlying_sink(
901 &self,
902 cx: &mut JSContext,
903 global: &GlobalScope,
904 stream: &WritableStream,
905 underlying_sink_obj: SafeHandleObject,
906 underlying_sink: &UnderlyingSink,
907 strategy_hwm: f64,
908 strategy_size: Rc<QueuingStrategySize>,
909 ) -> Result<(), Error> {
910 let controller = WritableStreamDefaultController::new(
936 cx,
937 global,
938 UnderlyingSinkType::new_js(
939 underlying_sink.abort.clone(),
940 underlying_sink.start.clone(),
941 underlying_sink.close.clone(),
942 underlying_sink.write.clone(),
943 ),
944 strategy_hwm,
945 strategy_size,
946 );
947
948 controller.set_underlying_sink_this_object(underlying_sink_obj);
951
952 controller.setup(cx, global, stream)
954 }
955}
956
957#[cfg_attr(crown, expect(crown::unrooted_must_root))]
959pub(crate) fn create_writable_stream(
960 cx: &mut JSContext,
961 global: &GlobalScope,
962 writable_high_water_mark: f64,
963 writable_size_algorithm: Rc<QueuingStrategySize>,
964 underlying_sink_type: UnderlyingSinkType,
965) -> Fallible<DomRoot<WritableStream>> {
966 assert!(writable_high_water_mark >= 0.0);
968
969 let stream = WritableStream::new_with_proto(cx, global, None);
972
973 let controller = WritableStreamDefaultController::new(
975 cx,
976 global,
977 underlying_sink_type,
978 writable_high_water_mark,
979 writable_size_algorithm,
980 );
981
982 controller.setup(cx, global, &stream)?;
985
986 Ok(stream)
988}
989
990impl WritableStreamMethods<crate::DomTypeHolder> for WritableStream {
991 fn Constructor(
993 cx: &mut JSContext,
994 global: &GlobalScope,
995 proto: Option<SafeHandleObject>,
996 underlying_sink: Option<*mut JSObject>,
997 strategy: &QueuingStrategy,
998 ) -> Fallible<DomRoot<WritableStream>> {
999 rooted!(&in(cx) let underlying_sink_obj = underlying_sink.unwrap_or(ptr::null_mut()));
1001
1002 let underlying_sink_dict = if !underlying_sink_obj.is_null() {
1005 rooted!(&in(cx) let obj_val = ObjectValue(underlying_sink_obj.get()));
1006 match UnderlyingSink::new(cx, obj_val.handle()) {
1007 Ok(ConversionResult::Success(val)) => val,
1008 Ok(ConversionResult::Failure(error)) => {
1009 return Err(Error::Type(error.into_owned()));
1010 },
1011 _ => {
1012 return Err(Error::JSFailed);
1013 },
1014 }
1015 } else {
1016 UnderlyingSink::empty()
1017 };
1018
1019 if !underlying_sink_dict.type_.handle().is_undefined() {
1020 return Err(Error::Range(c"type is set".to_owned()));
1022 }
1023
1024 let stream = WritableStream::new_with_proto(cx, global, proto);
1026
1027 let size_algorithm = extract_size_algorithm(cx, strategy);
1029
1030 let high_water_mark = extract_high_water_mark(strategy, 1.0)?;
1032
1033 stream.setup_from_underlying_sink(
1036 cx,
1037 global,
1038 &stream,
1039 underlying_sink_obj.handle(),
1040 &underlying_sink_dict,
1041 high_water_mark,
1042 size_algorithm,
1043 )?;
1044
1045 Ok(stream)
1046 }
1047
1048 fn Locked(&self) -> bool {
1050 self.is_locked()
1052 }
1053
1054 fn Abort(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) -> Rc<Promise> {
1056 let global = GlobalScope::from_current_realm(cx);
1057
1058 if self.is_locked() {
1060 let promise = Promise::new(cx, &global);
1062 promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1063 return promise;
1064 }
1065
1066 self.abort(cx, &global, reason).into()
1068 }
1069
1070 fn Close(&self, cx: &mut CurrentRealm) -> Rc<Promise> {
1072 let global = GlobalScope::from_current_realm(cx);
1073
1074 if self.is_locked() {
1076 let promise = Promise::new(cx, &global);
1078 promise.reject_error(cx, Error::Type(c"Stream is locked.".to_owned()));
1079 return promise;
1080 }
1081
1082 if self.close_queued_or_in_flight() {
1084 let promise = Promise::new(cx, &global);
1086 promise.reject_error(
1087 cx,
1088 Error::Type(c"Stream has closed queued or in-flight".to_owned()),
1089 );
1090 return promise;
1091 }
1092
1093 self.close(cx, &global).into()
1095 }
1096
1097 fn GetWriter(
1099 &self,
1100 realm: &mut CurrentRealm,
1101 ) -> Result<DomRoot<WritableStreamDefaultWriter>, Error> {
1102 let global = GlobalScope::from_current_realm(realm);
1103
1104 self.aquire_default_writer(realm, &global)
1106 }
1107}
1108
1109impl js::gc::Rootable for CrossRealmTransformWritable {}
1110
1111#[derive(Clone, JSTraceable, MallocSizeOf)]
1115#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
1116pub(crate) struct CrossRealmTransformWritable {
1117 controller: Dom<WritableStreamDefaultController>,
1119
1120 #[ignore_malloc_size_of = "nested Rc"]
1122 backpressure_promise: Rc<RefCell<Option<Rc<Promise>>>>,
1123}
1124
1125impl CrossRealmTransformWritable {
1126 pub(crate) fn handle_message(
1129 &self,
1130 cx: &mut CurrentRealm,
1131 global: &GlobalScope,
1132 message: SafeHandleValue,
1133 ) {
1134 rooted!(&in(cx) let mut value = UndefinedValue());
1135 let type_string = get_type_and_value_from_message(cx, message, value.handle_mut());
1136
1137 if type_string == "error" {
1142 self.controller.error_if_needed(cx, value.handle(), global);
1144 }
1145
1146 let backpressure_promise = self.backpressure_promise.borrow_mut().take();
1147
1148 if let Some(promise) = backpressure_promise {
1151 promise.resolve_native(cx, &());
1153
1154 }
1157 }
1158
1159 pub(crate) fn handle_error(
1162 &self,
1163 cx: &mut CurrentRealm,
1164 global: &GlobalScope,
1165 port: &MessagePort,
1166 ) {
1167 let error = DOMException::new(cx, global, DOMErrorName::DataCloneError);
1169 rooted!(&in(cx) let mut rooted_error = UndefinedValue());
1170 error.to_jsval(cx, rooted_error.handle_mut());
1171
1172 port.cross_realm_transform_send_error(cx, rooted_error.handle());
1174
1175 self.controller
1177 .error_if_needed(cx, rooted_error.handle(), global);
1178
1179 global.disentangle_port(cx, port);
1181 }
1182}
1183
1184impl Transferable for WritableStream {
1186 type Index = MessagePortIndex;
1187 type Data = MessagePortImpl;
1188
1189 fn transfer(&self, cx: &mut JSContext) -> Fallible<(MessagePortId, MessagePortImpl)> {
1191 if self.is_locked() {
1194 return Err(Error::DataClone(None));
1195 }
1196
1197 let global = self.global();
1198 let mut realm = enter_auto_realm(cx, &*global);
1199 let mut realm = realm.current_realm();
1200 let cx = &mut realm;
1201
1202 let port_1 = MessagePort::new(cx, &global);
1204 global.track_message_port(&port_1, None);
1205
1206 let port_2 = MessagePort::new(cx, &global);
1208 global.track_message_port(&port_2, None);
1209
1210 global.entangle_ports(*port_1.message_port_id(), *port_2.message_port_id());
1212
1213 let readable = ReadableStream::new_with_proto(cx, &global, None);
1215
1216 readable.setup_cross_realm_transform_readable(cx, &port_1);
1218
1219 let promise = readable.pipe_to(cx, &global, self, false, false, false, None);
1221
1222 promise.set_promise_is_handled(cx);
1224
1225 port_2.transfer(cx)
1227 }
1228
1229 fn transfer_receive(
1231 cx: &mut JSContext,
1232 owner: &GlobalScope,
1233 id: MessagePortId,
1234 port_impl: MessagePortImpl,
1235 ) -> Result<DomRoot<Self>, ()> {
1236 let value = WritableStream::new_with_proto(cx, owner, None);
1239
1240 let transferred_port = MessagePort::transfer_receive(cx, owner, id, port_impl)?;
1247
1248 value.setup_cross_realm_transform_writable(cx, &transferred_port);
1250 Ok(value)
1251 }
1252
1253 fn serialized_storage<'a>(
1255 data: StructuredData<'a, '_>,
1256 ) -> &'a mut Option<FxHashMap<MessagePortId, Self::Data>> {
1257 match data {
1258 StructuredData::Reader(r) => &mut r.port_impls,
1259 StructuredData::Writer(w) => &mut w.ports,
1260 }
1261 }
1262}