1use std::cell::Cell;
6use std::ptr::{self};
7use std::rc::Rc;
8
9use dom_struct::dom_struct;
10use js::context::JSContext;
11use js::jsapi::{Heap, IsPromiseObject, JSObject};
12use js::jsval::{JSVal, ObjectValue, UndefinedValue};
13use js::realm::CurrentRealm;
14use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue, IntoHandle};
15use rustc_hash::FxHashMap;
16use script_bindings::callback::ExceptionHandling;
17use script_bindings::cell::DomRefCell;
18use script_bindings::reflector::{Reflector, reflect_dom_object_with_proto};
19use servo_base::id::{MessagePortId, MessagePortIndex};
20use servo_constellation_traits::TransformStreamData;
21
22use super::readablestream::CrossRealmTransformReadable;
23use super::writablestream::CrossRealmTransformWritable;
24use crate::dom::bindings::codegen::Bindings::QueuingStrategyBinding::{
25 QueuingStrategy, QueuingStrategySize,
26};
27use crate::dom::bindings::codegen::Bindings::TransformStreamBinding::TransformStreamMethods;
28use crate::dom::bindings::codegen::Bindings::TransformerBinding::Transformer;
29use crate::dom::bindings::conversions::ConversionResult;
30use crate::dom::bindings::error::{Error, Fallible};
31use crate::dom::bindings::reflector::DomGlobal;
32use crate::dom::bindings::root::{Dom, DomRoot, MutNullableDom};
33use crate::dom::bindings::structuredclone::StructuredData;
34use crate::dom::bindings::transferable::Transferable;
35use crate::dom::globalscope::GlobalScope;
36use crate::dom::messageport::MessagePort;
37use crate::dom::promise::Promise;
38use crate::dom::promisenativehandler::Callback;
39use crate::dom::readablestream::{ReadableStream, create_readable_stream};
40use crate::dom::stream::countqueuingstrategy::{extract_high_water_mark, extract_size_algorithm};
41use crate::dom::stream::transformstreamdefaultcontroller::TransformerType;
42use crate::dom::stream::underlyingsourcecontainer::UnderlyingSourceType;
43use crate::dom::stream::writablestream::create_writable_stream;
44use crate::dom::stream::writablestreamdefaultcontroller::UnderlyingSinkType;
45use crate::dom::types::{PromiseNativeHandler, TransformStreamDefaultController, WritableStream};
46use crate::realms::enter_auto_realm;
47
48impl js::gc::Rootable for TransformBackPressureChangePromiseFulfillment {}
49
50#[derive(JSTraceable, MallocSizeOf)]
53#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
54struct TransformBackPressureChangePromiseFulfillment {
55 #[conditional_malloc_size_of]
57 result_promise: Rc<Promise>,
58
59 #[ignore_malloc_size_of = "mozjs"]
60 chunk: Box<Heap<JSVal>>,
61
62 writable: Dom<WritableStream>,
64
65 controller: Dom<TransformStreamDefaultController>,
66}
67
68impl Callback for TransformBackPressureChangePromiseFulfillment {
69 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
71 if self.writable.is_erroring() {
75 rooted!(&in(cx) let mut error = UndefinedValue());
76 self.writable.get_stored_error(error.handle_mut());
77 self.result_promise.reject(cx, error.handle());
78 return;
79 }
80
81 assert!(self.writable.is_writable());
83
84 rooted!(&in(cx) let mut chunk = UndefinedValue());
86 chunk.set(self.chunk.get());
87 let transform_result = self
88 .controller
89 .transform_stream_default_controller_perform_transform(
90 cx,
91 &self.writable.global(),
92 chunk.handle(),
93 )
94 .expect("perform transform failed");
95
96 let handler = PromiseNativeHandler::new(
99 cx,
100 &self.writable.global(),
101 Some(Box::new(PerformTransformFulfillment {
102 result_promise: self.result_promise.clone(),
103 })),
104 Some(Box::new(PerformTransformRejection {
105 result_promise: self.result_promise.clone(),
106 })),
107 );
108
109 let mut realm = enter_auto_realm(cx, &*self.writable.global());
110 let realm = &mut realm.current_realm();
111 transform_result.append_native_handler(realm, &handler);
112 }
113}
114
115#[derive(JSTraceable, MallocSizeOf)]
116#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
117struct PerformTransformFulfillment {
120 #[conditional_malloc_size_of]
121 result_promise: Rc<Promise>,
122}
123
124impl Callback for PerformTransformFulfillment {
125 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
126 self.result_promise.resolve_native(cx, &());
128 }
129}
130
131#[derive(JSTraceable, MallocSizeOf)]
132#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
133struct PerformTransformRejection {
136 #[conditional_malloc_size_of]
137 result_promise: Rc<Promise>,
138}
139
140impl Callback for PerformTransformRejection {
141 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
142 self.result_promise.reject(cx, v);
144 }
145}
146
147#[derive(JSTraceable, MallocSizeOf)]
148#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
149struct BackpressureChangeRejection {
152 #[conditional_malloc_size_of]
153 result_promise: Rc<Promise>,
154}
155
156impl Callback for BackpressureChangeRejection {
157 fn callback(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) {
158 self.result_promise.reject(cx, reason);
159 }
160}
161
162impl js::gc::Rootable for CancelPromiseFulfillment {}
163
164#[derive(JSTraceable, MallocSizeOf)]
167#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
168struct CancelPromiseFulfillment {
169 readable: Dom<ReadableStream>,
170 controller: Dom<TransformStreamDefaultController>,
171 #[ignore_malloc_size_of = "mozjs"]
172 reason: Box<Heap<JSVal>>,
173}
174
175impl Callback for CancelPromiseFulfillment {
176 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
178 if self.readable.is_errored() {
180 rooted!(&in(cx) let mut error = UndefinedValue());
181 self.readable.get_stored_error(error.handle_mut());
182 self.controller
183 .get_finish_promise()
184 .expect("finish promise is not set")
185 .reject_native(cx, &error.handle());
186 } else {
187 rooted!(&in(cx) let mut reason = UndefinedValue());
190 reason.set(self.reason.get());
191 self.readable
192 .get_default_controller()
193 .error(cx, reason.handle());
194
195 self.controller
197 .get_finish_promise()
198 .expect("finish promise is not set")
199 .resolve_native(cx, &());
200 }
201 }
202}
203
204impl js::gc::Rootable for CancelPromiseRejection {}
205
206#[derive(JSTraceable, MallocSizeOf)]
209#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
210struct CancelPromiseRejection {
211 readable: Dom<ReadableStream>,
212 controller: Dom<TransformStreamDefaultController>,
213}
214
215impl Callback for CancelPromiseRejection {
216 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
218 self.readable.get_default_controller().error(cx, v);
220
221 self.controller
223 .get_finish_promise()
224 .expect("finish promise is not set")
225 .reject(cx, v);
226 }
227}
228
229impl js::gc::Rootable for SourceCancelPromiseFulfillment {}
230
231#[derive(JSTraceable, MallocSizeOf)]
234#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
235struct SourceCancelPromiseFulfillment {
236 writeable: Dom<WritableStream>,
237 controller: Dom<TransformStreamDefaultController>,
238 stream: Dom<TransformStream>,
239 #[ignore_malloc_size_of = "mozjs"]
240 reason: Box<Heap<JSVal>>,
241}
242
243impl Callback for SourceCancelPromiseFulfillment {
244 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
246 let finish_promise = self
248 .controller
249 .get_finish_promise()
250 .expect("finish promise is not set");
251
252 let global = &self.writeable.global();
253 if self.writeable.is_errored() {
255 rooted!(&in(cx) let mut error = UndefinedValue());
256 self.writeable.get_stored_error(error.handle_mut());
257 finish_promise.reject(cx, error.handle());
258 } else {
259 rooted!(&in(cx) let mut reason = UndefinedValue());
262 reason.set(self.reason.get());
263 self.writeable
264 .get_default_controller()
265 .error_if_needed(cx, reason.handle(), global);
266
267 self.stream.unblock_write(cx, global);
269
270 finish_promise.resolve_native(cx, &());
272 }
273 }
274}
275
276impl js::gc::Rootable for SourceCancelPromiseRejection {}
277
278#[derive(JSTraceable, MallocSizeOf)]
281#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
282struct SourceCancelPromiseRejection {
283 writeable: Dom<WritableStream>,
284 controller: Dom<TransformStreamDefaultController>,
285 stream: Dom<TransformStream>,
286}
287
288impl Callback for SourceCancelPromiseRejection {
289 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
291 let global = &self.writeable.global();
293
294 self.writeable
295 .get_default_controller()
296 .error_if_needed(cx, v, global);
297
298 self.stream.unblock_write(cx, global);
300
301 self.controller
303 .get_finish_promise()
304 .expect("finish promise is not set")
305 .reject(cx, v);
306 }
307}
308
309impl js::gc::Rootable for FlushPromiseFulfillment {}
310
311#[derive(JSTraceable, MallocSizeOf)]
314#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
315struct FlushPromiseFulfillment {
316 readable: Dom<ReadableStream>,
317 controller: Dom<TransformStreamDefaultController>,
318}
319
320impl Callback for FlushPromiseFulfillment {
321 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
323 let finish_promise = self
325 .controller
326 .get_finish_promise()
327 .expect("finish promise is not set");
328
329 if self.readable.is_errored() {
331 rooted!(&in(cx) let mut error = UndefinedValue());
332 self.readable.get_stored_error(error.handle_mut());
333 finish_promise.reject(cx, error.handle());
334 } else {
335 self.readable.get_default_controller().close(cx);
338
339 finish_promise.resolve_native(cx, &());
341 }
342 }
343}
344
345impl js::gc::Rootable for FlushPromiseRejection {}
346#[derive(JSTraceable, MallocSizeOf)]
350#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
351struct FlushPromiseRejection {
352 readable: Dom<ReadableStream>,
353 controller: Dom<TransformStreamDefaultController>,
354}
355
356impl Callback for FlushPromiseRejection {
357 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
359 self.readable.get_default_controller().error(cx, v);
362
363 self.controller
365 .get_finish_promise()
366 .expect("finish promise is not set")
367 .reject(cx, v);
368 }
369}
370
371impl js::gc::Rootable for CrossRealmTransform {}
372
373#[derive(Clone, JSTraceable, MallocSizeOf)]
376#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
377pub(crate) enum CrossRealmTransform {
378 Readable(CrossRealmTransformReadable),
380 Writable(CrossRealmTransformWritable),
382}
383
384#[dom_struct]
386pub struct TransformStream {
387 reflector_: Reflector,
388
389 backpressure: Cell<bool>,
391
392 #[conditional_malloc_size_of]
394 backpressure_change_promise: DomRefCell<Option<Rc<Promise>>>,
395
396 controller: MutNullableDom<TransformStreamDefaultController>,
398
399 detached: Cell<bool>,
401
402 readable: MutNullableDom<ReadableStream>,
404
405 writable: MutNullableDom<WritableStream>,
407}
408
409impl TransformStream {
410 fn new_inherited() -> TransformStream {
412 TransformStream {
413 reflector_: Reflector::new(),
414 backpressure: Default::default(),
415 backpressure_change_promise: DomRefCell::new(None),
416 controller: MutNullableDom::new(None),
417 detached: Cell::new(false),
418 readable: MutNullableDom::new(None),
419 writable: MutNullableDom::new(None),
420 }
421 }
422
423 pub(crate) fn new_with_proto(
424 cx: &mut JSContext,
425 global: &GlobalScope,
426 proto: Option<SafeHandleObject>,
427 ) -> DomRoot<TransformStream> {
428 reflect_dom_object_with_proto(
429 cx,
430 Box::new(TransformStream::new_inherited()),
431 global,
432 proto,
433 )
434 }
435
436 #[cfg_attr(crown, expect(crown::unrooted_must_root))]
439 pub(crate) fn set_up(
440 &self,
441 cx: &mut JSContext,
442 global: &GlobalScope,
443 transformer_type: TransformerType,
444 ) -> Fallible<()> {
445 let writable_high_water_mark = 1.0;
447
448 let writable_size_algorithm = extract_size_algorithm(cx, &Default::default());
450
451 let readable_high_water_mark = 0.0;
453
454 let readable_size_algorithm = extract_size_algorithm(cx, &Default::default());
456
457 let start_promise = Promise::new_resolved(cx, global, ());
464
465 self.initialize(
469 cx,
470 global,
471 start_promise,
472 writable_high_water_mark,
473 writable_size_algorithm,
474 readable_high_water_mark,
475 readable_size_algorithm,
476 )?;
477
478 let controller = TransformStreamDefaultController::new(cx, global, transformer_type);
480
481 self.set_up_transform_stream_default_controller(&controller);
485
486 Ok(())
487 }
488
489 pub(crate) fn get_controller(&self) -> DomRoot<TransformStreamDefaultController> {
490 self.controller.get().expect("controller is not set")
491 }
492
493 pub(crate) fn get_writable(&self) -> DomRoot<WritableStream> {
494 self.writable.get().expect("writable stream is not set")
495 }
496
497 pub(crate) fn get_readable(&self) -> DomRoot<ReadableStream> {
498 self.readable.get().expect("readable stream is not set")
499 }
500
501 pub(crate) fn get_backpressure(&self) -> bool {
502 self.backpressure.get()
503 }
504
505 #[expect(clippy::too_many_arguments)]
507 fn initialize(
508 &self,
509 cx: &mut JSContext,
510 global: &GlobalScope,
511 start_promise: Rc<Promise>,
512 writable_high_water_mark: f64,
513 writable_size_algorithm: Rc<QueuingStrategySize>,
514 readable_high_water_mark: f64,
515 readable_size_algorithm: Rc<QueuingStrategySize>,
516 ) -> Fallible<()> {
517 let writable = create_writable_stream(
529 cx,
530 global,
531 writable_high_water_mark,
532 writable_size_algorithm,
533 UnderlyingSinkType::Transform(Dom::from_ref(self), start_promise.clone()),
534 )?;
535 self.writable.set(Some(&writable));
536
537 let readable = create_readable_stream(
549 cx,
550 global,
551 UnderlyingSourceType::Transform(self, start_promise),
552 Some(readable_size_algorithm),
553 Some(readable_high_water_mark),
554 );
555 self.readable.set(Some(&readable));
556
557 self.set_backpressure(cx, global, true);
562
563 self.controller.set(None);
565
566 Ok(())
567 }
568
569 pub(crate) fn set_backpressure(
571 &self,
572 cx: &mut JSContext,
573 global: &GlobalScope,
574 backpressure: bool,
575 ) {
576 assert!(self.backpressure.get() != backpressure);
578
579 if let Some(promise) = self.backpressure_change_promise.borrow_mut().take() {
582 promise.resolve_native(cx, &());
583 }
584
585 *self.backpressure_change_promise.borrow_mut() = Some(Promise::new(cx, global));
587
588 self.backpressure.set(backpressure);
590 }
591
592 fn set_up_transform_stream_default_controller(
594 &self,
595 controller: &TransformStreamDefaultController,
596 ) {
597 assert!(self.controller.get().is_none());
602
603 controller.set_stream(self);
605
606 self.controller.set(Some(controller));
608
609 }
614
615 fn set_up_transform_stream_default_controller_from_transformer(
617 &self,
618 cx: &mut JSContext,
619 global: &GlobalScope,
620 transformer_obj: SafeHandleObject,
621 transformer: &Transformer,
622 ) {
623 let controller = TransformStreamDefaultController::new(
625 cx,
626 global,
627 TransformerType::new_from_js_transformer(transformer),
628 );
629
630 controller.set_transform_obj(transformer_obj);
653
654 self.set_up_transform_stream_default_controller(&controller);
657 }
658
659 pub(crate) fn transform_stream_default_sink_write_algorithm(
661 &self,
662 cx: &mut JSContext,
663 global: &GlobalScope,
664 chunk: SafeHandleValue,
665 ) -> Fallible<Rc<Promise>> {
666 assert!(self.writable.get().is_some());
668
669 let controller = self.controller.get().expect("controller is not set");
671
672 if self.backpressure.get() {
674 let backpressure_change_promise = self.backpressure_change_promise.borrow();
676
677 assert!(backpressure_change_promise.is_some());
679
680 let result_promise = Promise::new(cx, global);
682 rooted!(&in(cx) let mut fulfillment_handler = Some(TransformBackPressureChangePromiseFulfillment {
683 controller: Dom::from_ref(&controller),
684 writable: Dom::from_ref(&self.writable.get().expect("writable stream")),
685 chunk: Heap::boxed(chunk.get()),
686 result_promise: result_promise.clone(),
687 }));
688
689 let handler = PromiseNativeHandler::new(
690 cx,
691 global,
692 fulfillment_handler.take().map(|h| Box::new(h) as Box<_>),
693 Some(Box::new(BackpressureChangeRejection {
694 result_promise: result_promise.clone(),
695 })),
696 );
697 let mut realm = enter_auto_realm(cx, global);
698 let realm = &mut realm.current_realm();
699 backpressure_change_promise
700 .as_ref()
701 .expect("Promise must be some by now.")
702 .append_native_handler(realm, &handler);
703
704 return Ok(result_promise);
705 }
706
707 controller.transform_stream_default_controller_perform_transform(cx, global, chunk)
709 }
710
711 pub(crate) fn transform_stream_default_sink_abort_algorithm(
713 &self,
714 cx: &mut JSContext,
715 global: &GlobalScope,
716 reason: SafeHandleValue,
717 ) -> Fallible<Rc<Promise>> {
718 let controller = self.controller.get().expect("controller is not set");
720
721 if let Some(finish_promise) = controller.get_finish_promise() {
723 return Ok(finish_promise);
724 }
725
726 let readable = self.readable.get().expect("readable stream is not set");
728
729 controller.set_finish_promise(Promise::new(cx, global));
731
732 let cancel_promise = controller.perform_cancel(cx, global, reason)?;
734
735 controller.clear_algorithms();
737
738 let handler = PromiseNativeHandler::new(
740 cx,
741 global,
742 Some(Box::new(CancelPromiseFulfillment {
743 readable: Dom::from_ref(&readable),
744 controller: Dom::from_ref(&controller),
745 reason: Heap::boxed(reason.get()),
746 })),
747 Some(Box::new(CancelPromiseRejection {
748 readable: Dom::from_ref(&readable),
749 controller: Dom::from_ref(&controller),
750 })),
751 );
752 let mut realm = enter_auto_realm(cx, global);
753 let cx = &mut realm.current_realm();
754 cancel_promise.append_native_handler(cx, &handler);
755
756 let finish_promise = controller
758 .get_finish_promise()
759 .expect("finish promise is not set");
760 Ok(finish_promise)
761 }
762
763 pub(crate) fn transform_stream_default_sink_close_algorithm(
765 &self,
766 cx: &mut JSContext,
767 global: &GlobalScope,
768 ) -> Fallible<Rc<Promise>> {
769 let controller = self
771 .controller
772 .get()
773 .ok_or(Error::Type(c"controller is not set".to_owned()))?;
774
775 if let Some(finish_promise) = controller.get_finish_promise() {
777 return Ok(finish_promise);
778 }
779
780 let readable = self
782 .readable
783 .get()
784 .ok_or(Error::Type(c"readable stream is not set".to_owned()))?;
785
786 controller.set_finish_promise(Promise::new(cx, global));
788
789 let flush_promise = controller.perform_flush(cx, global)?;
791
792 controller.clear_algorithms();
794
795 let handler = PromiseNativeHandler::new(
797 cx,
798 global,
799 Some(Box::new(FlushPromiseFulfillment {
800 readable: Dom::from_ref(&readable),
801 controller: Dom::from_ref(&controller),
802 })),
803 Some(Box::new(FlushPromiseRejection {
804 readable: Dom::from_ref(&readable),
805 controller: Dom::from_ref(&controller),
806 })),
807 );
808
809 let mut realm = enter_auto_realm(cx, global);
810 let realm = &mut realm.current_realm();
811 flush_promise.append_native_handler(realm, &handler);
812 let finish_promise = controller
814 .get_finish_promise()
815 .expect("finish promise is not set");
816 Ok(finish_promise)
817 }
818
819 pub(crate) fn transform_stream_default_source_cancel(
821 &self,
822 cx: &mut JSContext,
823 global: &GlobalScope,
824 reason: SafeHandleValue,
825 ) -> Fallible<Rc<Promise>> {
826 let controller = self
828 .controller
829 .get()
830 .ok_or(Error::Type(c"controller is not set".to_owned()))?;
831
832 if let Some(finish_promise) = controller.get_finish_promise() {
834 return Ok(finish_promise);
835 }
836
837 let writable = self
839 .writable
840 .get()
841 .ok_or(Error::Type(c"writable stream is not set".to_owned()))?;
842
843 controller.set_finish_promise(Promise::new(cx, global));
845
846 let cancel_promise = controller.perform_cancel(cx, global, reason)?;
848
849 controller.clear_algorithms();
851
852 let handler = PromiseNativeHandler::new(
854 cx,
855 global,
856 Some(Box::new(SourceCancelPromiseFulfillment {
857 writeable: Dom::from_ref(&writable),
858 controller: Dom::from_ref(&controller),
859 stream: Dom::from_ref(self),
860 reason: Heap::boxed(reason.get()),
861 })),
862 Some(Box::new(SourceCancelPromiseRejection {
863 writeable: Dom::from_ref(&writable),
864 controller: Dom::from_ref(&controller),
865 stream: Dom::from_ref(self),
866 })),
867 );
868
869 let finish_promise = controller
871 .get_finish_promise()
872 .expect("finish promise is not set");
873 let mut realm = enter_auto_realm(cx, global);
874 let cx = &mut realm.current_realm();
875 cancel_promise.append_native_handler(cx, &handler);
876 Ok(finish_promise)
877 }
878
879 pub(crate) fn transform_stream_default_source_pull(
881 &self,
882 cx: &mut JSContext,
883 global: &GlobalScope,
884 ) -> Fallible<Rc<Promise>> {
885 assert!(self.backpressure.get());
887
888 assert!(self.backpressure_change_promise.borrow().is_some());
890
891 self.set_backpressure(cx, global, false);
893
894 Ok(self
896 .backpressure_change_promise
897 .borrow()
898 .clone()
899 .expect("Promise must be some by now."))
900 }
901
902 pub(crate) fn error_writable_and_unblock_write(
904 &self,
905 cx: &mut JSContext,
906 global: &GlobalScope,
907 error: SafeHandleValue,
908 ) {
909 self.get_controller().clear_algorithms();
911
912 self.get_writable()
914 .get_default_controller()
915 .error_if_needed(cx, error, global);
916
917 self.unblock_write(cx, global)
919 }
920
921 pub(crate) fn unblock_write(&self, cx: &mut JSContext, global: &GlobalScope) {
923 if self.backpressure.get() {
925 self.set_backpressure(cx, global, false);
926 }
927 }
928
929 pub(crate) fn error(&self, cx: &mut JSContext, global: &GlobalScope, error: SafeHandleValue) {
931 self.get_readable()
933 .get_default_controller()
934 .error(cx, error);
935
936 self.error_writable_and_unblock_write(cx, global, error);
938 }
939}
940
941impl TransformStreamMethods<crate::DomTypeHolder> for TransformStream {
942 #[expect(unsafe_code)]
944 fn Constructor(
945 cx: &mut JSContext,
946 global: &GlobalScope,
947 proto: Option<SafeHandleObject>,
948 transformer: Option<*mut JSObject>,
949 writable_strategy: &QueuingStrategy,
950 readable_strategy: &QueuingStrategy,
951 ) -> Fallible<DomRoot<TransformStream>> {
952 rooted!(&in(cx) let transformer_obj = transformer.unwrap_or(ptr::null_mut()));
954
955 let transformer_dict = if !transformer_obj.is_null() {
958 rooted!(&in(cx) let obj_val = ObjectValue(transformer_obj.get()));
959 match Transformer::new(cx, obj_val.handle()) {
960 Ok(ConversionResult::Success(val)) => val,
961 Ok(ConversionResult::Failure(error)) => {
962 return Err(Error::Type(error.into_owned()));
963 },
964 _ => {
965 return Err(Error::JSFailed);
966 },
967 }
968 } else {
969 Transformer::empty()
970 };
971
972 if !transformer_dict.readableType.handle().is_undefined() {
974 return Err(Error::Range(c"readableType is set".to_owned()));
975 }
976
977 if !transformer_dict.writableType.handle().is_undefined() {
979 return Err(Error::Range(c"writableType is set".to_owned()));
980 }
981
982 let readable_high_water_mark = extract_high_water_mark(readable_strategy, 0.0)?;
984
985 let readable_size_algorithm = extract_size_algorithm(cx, readable_strategy);
987
988 let writable_high_water_mark = extract_high_water_mark(writable_strategy, 1.0)?;
990
991 let writable_size_algorithm = extract_size_algorithm(cx, writable_strategy);
993
994 let start_promise = Promise::new(cx, global);
996
997 let stream = TransformStream::new_with_proto(cx, global, proto);
1000 stream.initialize(
1001 cx,
1002 global,
1003 start_promise.clone(),
1004 writable_high_water_mark,
1005 writable_size_algorithm,
1006 readable_high_water_mark,
1007 readable_size_algorithm,
1008 )?;
1009
1010 stream.set_up_transform_stream_default_controller_from_transformer(
1012 cx,
1013 global,
1014 transformer_obj.handle(),
1015 &transformer_dict,
1016 );
1017
1018 if let Some(start) = &transformer_dict.start {
1022 rooted!(&in(cx) let mut result_object = ptr::null_mut::<JSObject>());
1023 rooted!(&in(cx) let mut result: JSVal);
1024 rooted!(&in(cx) let this_object = transformer_obj.get());
1025 start.Call_(
1026 cx,
1027 &this_object.handle(),
1028 &stream.get_controller(),
1029 result.handle_mut(),
1030 ExceptionHandling::Rethrow,
1031 )?;
1032 let is_promise = unsafe {
1033 if result.is_object() {
1034 result_object.set(result.to_object());
1035 IsPromiseObject(result_object.handle().into_handle())
1036 } else {
1037 false
1038 }
1039 };
1040 let promise = if is_promise {
1041 Promise::new_with_js_promise(cx, result_object.handle())
1042 } else {
1043 Promise::new_resolved(cx, global, result.get())
1044 };
1045 start_promise.resolve_native(cx, &promise);
1046 } else {
1047 start_promise.resolve_native(cx, &());
1049 };
1050
1051 Ok(stream)
1052 }
1053
1054 fn Readable(&self) -> DomRoot<ReadableStream> {
1056 self.readable.get().expect("readable stream is not set")
1058 }
1059
1060 fn Writable(&self) -> DomRoot<WritableStream> {
1062 self.writable.get().expect("writable stream is not set")
1064 }
1065}
1066
1067impl Transferable for TransformStream {
1069 type Index = MessagePortIndex;
1070 type Data = TransformStreamData;
1071
1072 fn transfer(&self, cx: &mut JSContext) -> Fallible<(MessagePortId, TransformStreamData)> {
1074 let global = self.global();
1075 let mut realm = enter_auto_realm(cx, &*global);
1076 let mut realm = realm.current_realm();
1077 let cx = &mut realm;
1078
1079 let readable = self.get_readable();
1081
1082 let writable = self.get_writable();
1084
1085 if readable.is_locked() || writable.is_locked() {
1090 return Err(Error::DataClone(None));
1091 }
1092
1093 let port1 = MessagePort::new(cx, &global);
1095 global.track_message_port(&port1, None);
1096 let port1_peer = MessagePort::new(cx, &global);
1097 global.track_message_port(&port1_peer, None);
1098 global.entangle_ports(*port1.message_port_id(), *port1_peer.message_port_id());
1099
1100 let proxy_readable = ReadableStream::new_with_proto(cx, &global, None);
1101 proxy_readable.setup_cross_realm_transform_readable(cx, &port1);
1102 proxy_readable
1103 .pipe_to(cx, &global, &writable, false, false, false, None)
1104 .set_promise_is_handled(cx);
1105
1106 let port2 = MessagePort::new(cx, &global);
1108 global.track_message_port(&port2, None);
1109 let port2_peer = MessagePort::new(cx, &global);
1110 global.track_message_port(&port2_peer, None);
1111 global.entangle_ports(*port2.message_port_id(), *port2_peer.message_port_id());
1112
1113 let proxy_writable = WritableStream::new_with_proto(cx, &global, None);
1114 proxy_writable.setup_cross_realm_transform_writable(cx, &port2);
1115
1116 readable
1118 .pipe_to(cx, &global, &proxy_writable, false, false, false, None)
1119 .set_promise_is_handled(cx);
1120
1121 Ok((
1126 *port1_peer.message_port_id(),
1127 TransformStreamData {
1128 readable: port1_peer.transfer(cx)?,
1129 writable: port2_peer.transfer(cx)?,
1130 },
1131 ))
1132 }
1133
1134 fn transfer_receive(
1136 cx: &mut JSContext,
1137 owner: &GlobalScope,
1138 _id: MessagePortId,
1139 data: TransformStreamData,
1140 ) -> Result<DomRoot<Self>, ()> {
1141 let port1 = MessagePort::transfer_receive(cx, owner, data.readable.0, data.readable.1)?;
1142 let port2 = MessagePort::transfer_receive(cx, owner, data.writable.0, data.writable.1)?;
1143
1144 let proxy_readable = ReadableStream::new_with_proto(cx, owner, None);
1148 proxy_readable.setup_cross_realm_transform_readable(cx, &port2);
1149
1150 let proxy_writable = WritableStream::new_with_proto(cx, owner, None);
1154 proxy_writable.setup_cross_realm_transform_writable(cx, &port1);
1155
1156 let stream = TransformStream::new_with_proto(cx, owner, None);
1162 stream.readable.set(Some(&proxy_readable));
1163 stream.writable.set(Some(&proxy_writable));
1164
1165 Ok(stream)
1166 }
1167
1168 fn serialized_storage<'a>(
1169 data: StructuredData<'a, '_>,
1170 ) -> &'a mut Option<FxHashMap<MessagePortId, Self::Data>> {
1171 match data {
1172 StructuredData::Reader(r) => &mut r.transform_streams_port_impls,
1173 StructuredData::Writer(w) => &mut w.transform_streams_port,
1174 }
1175 }
1176}