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, JSObject};
12use js::jsval::{JSVal, ObjectValue, UndefinedValue};
13use js::realm::CurrentRealm;
14use js::rust::{HandleObject as SafeHandleObject, HandleValue as SafeHandleValue};
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, RootedPromise, TracedPromise};
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 result_promise: TracedPromise,
57
58 #[ignore_malloc_size_of = "mozjs"]
59 chunk: Box<Heap<JSVal>>,
60
61 writable: Dom<WritableStream>,
63
64 controller: Dom<TransformStreamDefaultController>,
65}
66
67impl Callback for TransformBackPressureChangePromiseFulfillment {
68 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
70 if self.writable.is_erroring() {
74 rooted!(&in(cx) let mut error = UndefinedValue());
75 self.writable.get_stored_error(error.handle_mut());
76 self.result_promise.reject(cx, error.handle());
77 return;
78 }
79
80 assert!(self.writable.is_writable());
82
83 rooted!(&in(cx) let mut chunk = UndefinedValue());
85 chunk.set(self.chunk.get());
86 let transform_result = self
87 .controller
88 .transform_stream_default_controller_perform_transform(
89 cx,
90 &self.writable.global(),
91 chunk.handle(),
92 )
93 .expect("perform transform failed");
94
95 let handler = PromiseNativeHandler::new(
98 cx,
99 &self.writable.global(),
100 Some(Box::new(PerformTransformFulfillment {
101 result_promise: self.result_promise.clone(),
102 })),
103 Some(Box::new(PerformTransformRejection {
104 result_promise: self.result_promise.clone(),
105 })),
106 );
107
108 let mut realm = enter_auto_realm(cx, &*self.writable.global());
109 let realm = &mut realm.current_realm();
110 transform_result.append_native_handler(realm, &handler);
111 }
112}
113
114#[derive(JSTraceable, MallocSizeOf)]
115#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
116struct PerformTransformFulfillment {
119 result_promise: TracedPromise,
120}
121
122impl Callback for PerformTransformFulfillment {
123 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
124 self.result_promise.resolve_native(cx, &());
126 }
127}
128
129#[derive(JSTraceable, MallocSizeOf)]
130#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
131struct PerformTransformRejection {
134 result_promise: TracedPromise,
135}
136
137impl Callback for PerformTransformRejection {
138 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
139 self.result_promise.reject(cx, v);
141 }
142}
143
144#[derive(JSTraceable, MallocSizeOf)]
145#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
146struct BackpressureChangeRejection {
149 result_promise: TracedPromise,
150}
151
152impl Callback for BackpressureChangeRejection {
153 fn callback(&self, cx: &mut CurrentRealm, reason: SafeHandleValue) {
154 self.result_promise.reject(cx, reason);
155 }
156}
157
158impl js::gc::Rootable for CancelPromiseFulfillment {}
159
160#[derive(JSTraceable, MallocSizeOf)]
163#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
164struct CancelPromiseFulfillment {
165 readable: Dom<ReadableStream>,
166 controller: Dom<TransformStreamDefaultController>,
167 #[ignore_malloc_size_of = "mozjs"]
168 reason: Box<Heap<JSVal>>,
169}
170
171impl Callback for CancelPromiseFulfillment {
172 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
174 if self.readable.is_errored() {
176 rooted!(&in(cx) let mut error = UndefinedValue());
177 self.readable.get_stored_error(error.handle_mut());
178 self.controller
179 .get_finish_promise(cx)
180 .expect("finish promise is not set")
181 .reject_native(cx, &error.handle());
182 } else {
183 rooted!(&in(cx) let mut reason = UndefinedValue());
186 reason.set(self.reason.get());
187 self.readable
188 .get_default_controller()
189 .error(cx, reason.handle());
190
191 self.controller
193 .get_finish_promise(cx)
194 .expect("finish promise is not set")
195 .resolve_native(cx, &());
196 }
197 }
198}
199
200impl js::gc::Rootable for CancelPromiseRejection {}
201
202#[derive(JSTraceable, MallocSizeOf)]
205#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
206struct CancelPromiseRejection {
207 readable: Dom<ReadableStream>,
208 controller: Dom<TransformStreamDefaultController>,
209}
210
211impl Callback for CancelPromiseRejection {
212 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
214 self.readable.get_default_controller().error(cx, v);
216
217 self.controller
219 .get_finish_promise(cx)
220 .expect("finish promise is not set")
221 .reject(cx, v);
222 }
223}
224
225impl js::gc::Rootable for SourceCancelPromiseFulfillment {}
226
227#[derive(JSTraceable, MallocSizeOf)]
230#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
231struct SourceCancelPromiseFulfillment {
232 writeable: Dom<WritableStream>,
233 controller: Dom<TransformStreamDefaultController>,
234 stream: Dom<TransformStream>,
235 #[ignore_malloc_size_of = "mozjs"]
236 reason: Box<Heap<JSVal>>,
237}
238
239impl Callback for SourceCancelPromiseFulfillment {
240 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
242 let finish_promise = self
244 .controller
245 .get_finish_promise(cx)
246 .expect("finish promise is not set");
247
248 let global = &self.writeable.global();
249 if self.writeable.is_errored() {
251 rooted!(&in(cx) let mut error = UndefinedValue());
252 self.writeable.get_stored_error(error.handle_mut());
253 finish_promise.reject(cx, error.handle());
254 } else {
255 rooted!(&in(cx) let mut reason = UndefinedValue());
258 reason.set(self.reason.get());
259 self.writeable
260 .get_default_controller()
261 .error_if_needed(cx, reason.handle(), global);
262
263 self.stream.unblock_write(cx, global);
265
266 finish_promise.resolve_native(cx, &());
268 }
269 }
270}
271
272impl js::gc::Rootable for SourceCancelPromiseRejection {}
273
274#[derive(JSTraceable, MallocSizeOf)]
277#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
278struct SourceCancelPromiseRejection {
279 writeable: Dom<WritableStream>,
280 controller: Dom<TransformStreamDefaultController>,
281 stream: Dom<TransformStream>,
282}
283
284impl Callback for SourceCancelPromiseRejection {
285 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
287 let global = &self.writeable.global();
289
290 self.writeable
291 .get_default_controller()
292 .error_if_needed(cx, v, global);
293
294 self.stream.unblock_write(cx, global);
296
297 self.controller
299 .get_finish_promise(cx)
300 .expect("finish promise is not set")
301 .reject(cx, v);
302 }
303}
304
305impl js::gc::Rootable for FlushPromiseFulfillment {}
306
307#[derive(JSTraceable, MallocSizeOf)]
310#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
311struct FlushPromiseFulfillment {
312 readable: Dom<ReadableStream>,
313 controller: Dom<TransformStreamDefaultController>,
314}
315
316impl Callback for FlushPromiseFulfillment {
317 fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
319 let finish_promise = self
321 .controller
322 .get_finish_promise(cx)
323 .expect("finish promise is not set");
324
325 if self.readable.is_errored() {
327 rooted!(&in(cx) let mut error = UndefinedValue());
328 self.readable.get_stored_error(error.handle_mut());
329 finish_promise.reject(cx, error.handle());
330 } else {
331 self.readable.get_default_controller().close(cx);
334
335 finish_promise.resolve_native(cx, &());
337 }
338 }
339}
340
341impl js::gc::Rootable for FlushPromiseRejection {}
342#[derive(JSTraceable, MallocSizeOf)]
346#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
347struct FlushPromiseRejection {
348 readable: Dom<ReadableStream>,
349 controller: Dom<TransformStreamDefaultController>,
350}
351
352impl Callback for FlushPromiseRejection {
353 fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
355 self.readable.get_default_controller().error(cx, v);
358
359 self.controller
361 .get_finish_promise(cx)
362 .expect("finish promise is not set")
363 .reject(cx, v);
364 }
365}
366
367impl js::gc::Rootable for CrossRealmTransform {}
368
369#[derive(Clone, JSTraceable, MallocSizeOf)]
372#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
373pub(crate) enum CrossRealmTransform {
374 Readable(CrossRealmTransformReadable),
376 Writable(CrossRealmTransformWritable),
378}
379
380#[dom_struct]
382pub struct TransformStream {
383 reflector_: Reflector,
384
385 backpressure: Cell<bool>,
387
388 backpressure_change_promise: DomRefCell<Option<TracedPromise>>,
390
391 controller: MutNullableDom<TransformStreamDefaultController>,
393
394 detached: Cell<bool>,
396
397 readable: MutNullableDom<ReadableStream>,
399
400 writable: MutNullableDom<WritableStream>,
402}
403
404impl TransformStream {
405 fn new_inherited() -> TransformStream {
407 TransformStream {
408 reflector_: Reflector::new(),
409 backpressure: Default::default(),
410 backpressure_change_promise: DomRefCell::new(None),
411 controller: MutNullableDom::new(None),
412 detached: Cell::new(false),
413 readable: MutNullableDom::new(None),
414 writable: MutNullableDom::new(None),
415 }
416 }
417
418 pub(crate) fn new_with_proto(
419 cx: &mut JSContext,
420 global: &GlobalScope,
421 proto: Option<SafeHandleObject>,
422 ) -> DomRoot<TransformStream> {
423 reflect_dom_object_with_proto(
424 cx,
425 Box::new(TransformStream::new_inherited()),
426 global,
427 proto,
428 )
429 }
430
431 #[cfg_attr(crown, expect(crown::unrooted_must_root))]
434 pub(crate) fn set_up(
435 &self,
436 cx: &mut JSContext,
437 global: &GlobalScope,
438 transformer_type: TransformerType,
439 ) -> Fallible<()> {
440 let writable_high_water_mark = 1.0;
442
443 let writable_size_algorithm = extract_size_algorithm(cx, &Default::default());
445
446 let readable_high_water_mark = 0.0;
448
449 let readable_size_algorithm = extract_size_algorithm(cx, &Default::default());
451
452 let start_promise = Promise::new_resolved_rooted(cx, global, ());
459
460 self.initialize(
464 cx,
465 global,
466 &start_promise,
467 writable_high_water_mark,
468 writable_size_algorithm,
469 readable_high_water_mark,
470 readable_size_algorithm,
471 )?;
472
473 let controller = TransformStreamDefaultController::new(cx, global, transformer_type);
475
476 self.set_up_transform_stream_default_controller(&controller);
480
481 Ok(())
482 }
483
484 pub(crate) fn get_controller(&self) -> DomRoot<TransformStreamDefaultController> {
485 self.controller.get().expect("controller is not set")
486 }
487
488 pub(crate) fn get_writable(&self) -> DomRoot<WritableStream> {
489 self.writable.get().expect("writable stream is not set")
490 }
491
492 pub(crate) fn get_readable(&self) -> DomRoot<ReadableStream> {
493 self.readable.get().expect("readable stream is not set")
494 }
495
496 pub(crate) fn get_backpressure(&self) -> bool {
497 self.backpressure.get()
498 }
499
500 #[expect(clippy::too_many_arguments)]
502 fn initialize(
503 &self,
504 cx: &mut JSContext,
505 global: &GlobalScope,
506 start_promise: &RootedPromise,
507 writable_high_water_mark: f64,
508 writable_size_algorithm: Rc<QueuingStrategySize>,
509 readable_high_water_mark: f64,
510 readable_size_algorithm: Rc<QueuingStrategySize>,
511 ) -> Fallible<()> {
512 let writable = create_writable_stream(
524 cx,
525 global,
526 writable_high_water_mark,
527 writable_size_algorithm,
528 UnderlyingSinkType::Transform(Dom::from_ref(self), start_promise.to_traced()),
529 )?;
530 self.writable.set(Some(&writable));
531
532 let readable = create_readable_stream(
544 cx,
545 global,
546 UnderlyingSourceType::Transform(self, start_promise),
547 Some(readable_size_algorithm),
548 Some(readable_high_water_mark),
549 );
550 self.readable.set(Some(&readable));
551
552 self.set_backpressure(cx, global, true);
557
558 self.controller.set(None);
560
561 Ok(())
562 }
563
564 pub(crate) fn set_backpressure(
566 &self,
567 cx: &mut JSContext,
568 global: &GlobalScope,
569 backpressure: bool,
570 ) {
571 assert!(self.backpressure.get() != backpressure);
573
574 rooted!(&in(cx) let promise = self.backpressure_change_promise.borrow_mut().take());
577 if let Some(ref promise) = *promise {
578 promise.resolve_native(cx, &());
579 }
580
581 *self.backpressure_change_promise.borrow_mut() =
583 Some(Promise::new_rooted(cx, global).to_traced());
584
585 self.backpressure.set(backpressure);
587 }
588
589 fn set_up_transform_stream_default_controller(
591 &self,
592 controller: &TransformStreamDefaultController,
593 ) {
594 assert!(self.controller.get().is_none());
599
600 controller.set_stream(self);
602
603 self.controller.set(Some(controller));
605
606 }
611
612 fn set_up_transform_stream_default_controller_from_transformer(
614 &self,
615 cx: &mut JSContext,
616 global: &GlobalScope,
617 transformer_obj: SafeHandleObject,
618 transformer: &Transformer,
619 ) {
620 let controller = TransformStreamDefaultController::new(
622 cx,
623 global,
624 TransformerType::new_from_js_transformer(transformer),
625 );
626
627 controller.set_transform_obj(transformer_obj);
650
651 self.set_up_transform_stream_default_controller(&controller);
654 }
655
656 pub(crate) fn transform_stream_default_sink_write_algorithm(
658 &self,
659 cx: &mut JSContext,
660 global: &GlobalScope,
661 chunk: SafeHandleValue,
662 ) -> Fallible<RootedPromise> {
663 assert!(self.writable.get().is_some());
665
666 let controller = self.controller.get().expect("controller is not set");
668
669 if self.backpressure.get() {
671 let backpressure_change_promise = self.backpressure_change_promise.borrow();
673
674 assert!(backpressure_change_promise.is_some());
676
677 let result_promise = Promise::new_rooted(cx, global);
679 rooted!(&in(cx) let mut fulfillment_handler = Some(TransformBackPressureChangePromiseFulfillment {
680 controller: Dom::from_ref(&controller),
681 writable: Dom::from_ref(&self.writable.get().expect("writable stream")),
682 chunk: Heap::boxed(chunk.get()),
683 result_promise: result_promise.to_traced(),
684 }));
685
686 let handler = PromiseNativeHandler::new(
687 cx,
688 global,
689 fulfillment_handler.take().map(|h| Box::new(h) as Box<_>),
690 Some(Box::new(BackpressureChangeRejection {
691 result_promise: result_promise.to_traced(),
692 })),
693 );
694 let mut realm = enter_auto_realm(cx, global);
695 let realm = &mut realm.current_realm();
696 backpressure_change_promise
697 .as_ref()
698 .expect("Promise must be some by now.")
699 .append_native_handler(realm, &handler);
700
701 return Ok(result_promise);
702 }
703
704 controller.transform_stream_default_controller_perform_transform(cx, global, chunk)
706 }
707
708 pub(crate) fn transform_stream_default_sink_abort_algorithm(
710 &self,
711 cx: &mut JSContext,
712 global: &GlobalScope,
713 reason: SafeHandleValue,
714 ) -> Fallible<RootedPromise> {
715 let controller = self.controller.get().expect("controller is not set");
717
718 if let Some(finish_promise) = controller.get_finish_promise(cx) {
720 return Ok(finish_promise);
721 }
722
723 let readable = self.readable.get().expect("readable stream is not set");
725
726 controller.set_finish_promise(&Promise::new_rooted(cx, global));
728
729 let cancel_promise = controller.perform_cancel(cx, global, reason)?;
731
732 controller.clear_algorithms();
734
735 let handler = PromiseNativeHandler::new(
737 cx,
738 global,
739 Some(Box::new(CancelPromiseFulfillment {
740 readable: Dom::from_ref(&readable),
741 controller: Dom::from_ref(&controller),
742 reason: Heap::boxed(reason.get()),
743 })),
744 Some(Box::new(CancelPromiseRejection {
745 readable: Dom::from_ref(&readable),
746 controller: Dom::from_ref(&controller),
747 })),
748 );
749 let mut realm = enter_auto_realm(cx, global);
750 let cx = &mut realm.current_realm();
751 cancel_promise.append_native_handler(cx, &handler);
752
753 let finish_promise = controller
755 .get_finish_promise(cx)
756 .expect("finish promise is not set");
757 Ok(finish_promise)
758 }
759
760 pub(crate) fn transform_stream_default_sink_close_algorithm(
762 &self,
763 cx: &mut JSContext,
764 global: &GlobalScope,
765 ) -> Fallible<RootedPromise> {
766 let controller = self
768 .controller
769 .get()
770 .ok_or(Error::Type(c"controller is not set".to_owned()))?;
771
772 if let Some(finish_promise) = controller.get_finish_promise(cx) {
774 return Ok(finish_promise);
775 }
776
777 let readable = self
779 .readable
780 .get()
781 .ok_or(Error::Type(c"readable stream is not set".to_owned()))?;
782
783 controller.set_finish_promise(&Promise::new_rooted(cx, global));
785
786 let flush_promise = controller.perform_flush(cx, global)?;
788
789 controller.clear_algorithms();
791
792 let handler = PromiseNativeHandler::new(
794 cx,
795 global,
796 Some(Box::new(FlushPromiseFulfillment {
797 readable: Dom::from_ref(&readable),
798 controller: Dom::from_ref(&controller),
799 })),
800 Some(Box::new(FlushPromiseRejection {
801 readable: Dom::from_ref(&readable),
802 controller: Dom::from_ref(&controller),
803 })),
804 );
805
806 {
807 let mut realm = enter_auto_realm(cx, global);
808 let realm = &mut realm.current_realm();
809 flush_promise.append_native_handler(realm, &handler);
810 }
811
812 let finish_promise = controller
814 .get_finish_promise(cx)
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<RootedPromise> {
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(cx) {
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_rooted(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 mut realm = enter_auto_realm(cx, global);
870 let cx = &mut realm.current_realm();
871 cancel_promise.append_native_handler(cx, &handler);
872
873 let finish_promise = controller
875 .get_finish_promise(cx)
876 .expect("finish promise is not set");
877 Ok(finish_promise)
878 }
879
880 pub(crate) fn transform_stream_default_source_pull(
882 &self,
883 cx: &mut JSContext,
884 global: &GlobalScope,
885 ) -> Fallible<RootedPromise> {
886 assert!(self.backpressure.get());
888
889 assert!(self.backpressure_change_promise.borrow().is_some());
891
892 self.set_backpressure(cx, global, false);
894
895 Ok(self
897 .backpressure_change_promise
898 .borrow()
899 .as_ref()
900 .expect("Promise must be some by now.")
901 .root(cx))
902 }
903
904 pub(crate) fn error_writable_and_unblock_write(
906 &self,
907 cx: &mut JSContext,
908 global: &GlobalScope,
909 error: SafeHandleValue,
910 ) {
911 self.get_controller().clear_algorithms();
913
914 self.get_writable()
916 .get_default_controller()
917 .error_if_needed(cx, error, global);
918
919 self.unblock_write(cx, global)
921 }
922
923 pub(crate) fn unblock_write(&self, cx: &mut JSContext, global: &GlobalScope) {
925 if self.backpressure.get() {
927 self.set_backpressure(cx, global, false);
928 }
929 }
930
931 pub(crate) fn error(&self, cx: &mut JSContext, global: &GlobalScope, error: SafeHandleValue) {
933 self.get_readable()
935 .get_default_controller()
936 .error(cx, error);
937
938 self.error_writable_and_unblock_write(cx, global, error);
940 }
941}
942
943impl TransformStreamMethods<crate::DomTypeHolder> for TransformStream {
944 fn Constructor(
946 cx: &mut JSContext,
947 global: &GlobalScope,
948 proto: Option<SafeHandleObject>,
949 transformer: Option<*mut JSObject>,
950 writable_strategy: &QueuingStrategy,
951 readable_strategy: &QueuingStrategy,
952 ) -> Fallible<DomRoot<TransformStream>> {
953 rooted!(&in(cx) let transformer_obj = transformer.unwrap_or(ptr::null_mut()));
955
956 let transformer_dict = if !transformer_obj.is_null() {
959 rooted!(&in(cx) let obj_val = ObjectValue(transformer_obj.get()));
960 match Transformer::new(cx, obj_val.handle()) {
961 Ok(ConversionResult::Success(val)) => val,
962 Ok(ConversionResult::Failure(error)) => {
963 return Err(Error::Type(error.into_owned()));
964 },
965 _ => {
966 return Err(Error::JSFailed);
967 },
968 }
969 } else {
970 Transformer::empty()
971 };
972
973 if !transformer_dict.readableType.handle().is_undefined() {
975 return Err(Error::Range(c"readableType is set".to_owned()));
976 }
977
978 if !transformer_dict.writableType.handle().is_undefined() {
980 return Err(Error::Range(c"writableType is set".to_owned()));
981 }
982
983 let readable_high_water_mark = extract_high_water_mark(readable_strategy, 0.0)?;
985
986 let readable_size_algorithm = extract_size_algorithm(cx, readable_strategy);
988
989 let writable_high_water_mark = extract_high_water_mark(writable_strategy, 1.0)?;
991
992 let writable_size_algorithm = extract_size_algorithm(cx, writable_strategy);
994
995 let start_promise = Promise::new_rooted(cx, global);
997
998 let stream = TransformStream::new_with_proto(cx, global, proto);
1001 stream.initialize(
1002 cx,
1003 global,
1004 &start_promise,
1005 writable_high_water_mark,
1006 writable_size_algorithm,
1007 readable_high_water_mark,
1008 readable_size_algorithm,
1009 )?;
1010
1011 stream.set_up_transform_stream_default_controller_from_transformer(
1013 cx,
1014 global,
1015 transformer_obj.handle(),
1016 &transformer_dict,
1017 );
1018
1019 if let Some(start) = &transformer_dict.start {
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 promise = Promise::resolve_or_wrap_promise(cx, result.handle(), global);
1033 start_promise.resolve_native(cx, &promise);
1034 } else {
1035 start_promise.resolve_native(cx, &());
1037 };
1038
1039 Ok(stream)
1040 }
1041
1042 fn Readable(&self) -> DomRoot<ReadableStream> {
1044 self.readable.get().expect("readable stream is not set")
1046 }
1047
1048 fn Writable(&self) -> DomRoot<WritableStream> {
1050 self.writable.get().expect("writable stream is not set")
1052 }
1053}
1054
1055impl Transferable for TransformStream {
1057 type Index = MessagePortIndex;
1058 type Data = TransformStreamData;
1059
1060 fn transfer(&self, cx: &mut JSContext) -> Fallible<(MessagePortId, TransformStreamData)> {
1062 let global = self.global();
1063 let mut realm = enter_auto_realm(cx, &*global);
1064 let mut realm = realm.current_realm();
1065 let cx = &mut realm;
1066
1067 let readable = self.get_readable();
1069
1070 let writable = self.get_writable();
1072
1073 if readable.is_locked() || writable.is_locked() {
1078 return Err(Error::DataClone(None));
1079 }
1080
1081 let port1 = MessagePort::new(cx, &global);
1083 global.track_message_port(&port1, None);
1084 let port1_peer = MessagePort::new(cx, &global);
1085 global.track_message_port(&port1_peer, None);
1086 global.entangle_ports(*port1.message_port_id(), *port1_peer.message_port_id());
1087
1088 let proxy_readable = ReadableStream::new_with_proto(cx, &global, None);
1089 proxy_readable.setup_cross_realm_transform_readable(cx, &port1);
1090 proxy_readable
1091 .pipe_to(cx, &global, &writable, false, false, false, None)
1092 .set_promise_is_handled(cx);
1093
1094 let port2 = MessagePort::new(cx, &global);
1096 global.track_message_port(&port2, None);
1097 let port2_peer = MessagePort::new(cx, &global);
1098 global.track_message_port(&port2_peer, None);
1099 global.entangle_ports(*port2.message_port_id(), *port2_peer.message_port_id());
1100
1101 let proxy_writable = WritableStream::new_with_proto(cx, &global, None);
1102 proxy_writable.setup_cross_realm_transform_writable(cx, &port2);
1103
1104 readable
1106 .pipe_to(cx, &global, &proxy_writable, false, false, false, None)
1107 .set_promise_is_handled(cx);
1108
1109 Ok((
1114 *port1_peer.message_port_id(),
1115 TransformStreamData {
1116 readable: port1_peer.transfer(cx)?,
1117 writable: port2_peer.transfer(cx)?,
1118 },
1119 ))
1120 }
1121
1122 fn transfer_receive(
1124 cx: &mut JSContext,
1125 owner: &GlobalScope,
1126 _id: MessagePortId,
1127 data: TransformStreamData,
1128 ) -> Result<DomRoot<Self>, ()> {
1129 let port1 = MessagePort::transfer_receive(cx, owner, data.readable.0, data.readable.1)?;
1130 let port2 = MessagePort::transfer_receive(cx, owner, data.writable.0, data.writable.1)?;
1131
1132 let proxy_readable = ReadableStream::new_with_proto(cx, owner, None);
1136 proxy_readable.setup_cross_realm_transform_readable(cx, &port2);
1137
1138 let proxy_writable = WritableStream::new_with_proto(cx, owner, None);
1142 proxy_writable.setup_cross_realm_transform_writable(cx, &port1);
1143
1144 let stream = TransformStream::new_with_proto(cx, owner, None);
1150 stream.readable.set(Some(&proxy_readable));
1151 stream.writable.set(Some(&proxy_writable));
1152
1153 Ok(stream)
1154 }
1155
1156 fn serialized_storage<'a>(
1157 data: StructuredData<'a, '_>,
1158 ) -> &'a mut Option<FxHashMap<MessagePortId, Self::Data>> {
1159 match data {
1160 StructuredData::Reader(r) => &mut r.transform_streams_port_impls,
1161 StructuredData::Writer(w) => &mut w.transform_streams_port,
1162 }
1163 }
1164}