Skip to main content

script/dom/stream/
transformstream.rs

1/* This Source Code Form is subject to the terms of the Mozilla Public
2 * License, v. 2.0. If a copy of the MPL was not distributed with this
3 * file, You can obtain one at https://mozilla.org/MPL/2.0/. */
4
5use 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/// Reacting to backpressureChangePromise as part of
51/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
52#[derive(JSTraceable, MallocSizeOf)]
53#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
54struct TransformBackPressureChangePromiseFulfillment {
55    /// The result of reacting to backpressureChangePromise.
56    #[conditional_malloc_size_of]
57    result_promise: Rc<Promise>,
58
59    #[ignore_malloc_size_of = "mozjs"]
60    chunk: Box<Heap<JSVal>>,
61
62    /// The writable used in the fulfillment steps
63    writable: Dom<WritableStream>,
64
65    controller: Dom<TransformStreamDefaultController>,
66}
67
68impl Callback for TransformBackPressureChangePromiseFulfillment {
69    /// Reacting to backpressureChangePromise with the following fulfillment steps:
70    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
71        // Let writable be stream.[[writable]].
72        // Let state be writable.[[state]].
73        // If state is "erroring", throw writable.[[storedError]].
74        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: state is "writable".
82        assert!(self.writable.is_writable());
83
84        // Return ! TransformStreamDefaultControllerPerformTransform(controller, chunk).
85        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        // PerformTransformFulfillment and PerformTransformRejection do not need
97        // to be rooted because they only contain an Rc.
98        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)]
117/// Reacting to fulfillment of performTransform as part of
118/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
119struct 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        // Fulfilled: resolve the outer promise
127        self.result_promise.resolve_native(cx, &());
128    }
129}
130
131#[derive(JSTraceable, MallocSizeOf)]
132#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
133/// Reacting to rejection of performTransform as part of
134/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
135struct 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        // Stream already errored in perform_transform, just reject result_promise
143        self.result_promise.reject(cx, v);
144    }
145}
146
147#[derive(JSTraceable, MallocSizeOf)]
148#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
149/// Reacting to rejection of backpressureChangePromise as part of
150/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
151struct 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/// Reacting to fulfillment of the cancelpromise as part of
165/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-abort-algorithm>
166#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
177    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
178        // If readable.[[state]] is "errored", reject controller.[[finishPromise]] with readable.[[storedError]].
179        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            // Otherwise:
188            // Perform ! ReadableStreamDefaultControllerError(readable.[[controller]], reason).
189            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            // Resolve controller.[[finishPromise]] with undefined.
196            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/// Reacting to rejection of cancelpromise as part of
207/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-abort-algorithm>
208#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
217    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
218        // Perform ! ReadableStreamDefaultControllerError(readable.[[controller]], r).
219        self.readable.get_default_controller().error(cx, v);
220
221        // Reject controller.[[finishPromise]] with r.
222        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/// Reacting to fulfillment of the cancelpromise as part of
232/// <https://streams.spec.whatwg.org/#transform-stream-default-source-cancel>
233#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
245    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
246        // If cancelPromise was fulfilled, then:
247        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 writable.[[state]] is "errored", reject controller.[[finishPromise]] with writable.[[storedError]].
254        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            // Otherwise:
260            // Perform ! WritableStreamDefaultControllerErrorIfNeeded(writable.[[controller]], reason).
261            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            // Perform ! TransformStreamUnblockWrite(stream).
268            self.stream.unblock_write(cx, global);
269
270            // Resolve controller.[[finishPromise]] with undefined.
271            finish_promise.resolve_native(cx, &());
272        }
273    }
274}
275
276impl js::gc::Rootable for SourceCancelPromiseRejection {}
277
278/// Reacting to rejection of cancelpromise as part of
279/// <https://streams.spec.whatwg.org/#transform-stream-default-source-cancel>
280#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
290    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
291        // Perform ! WritableStreamDefaultControllerErrorIfNeeded(writable.[[controller]], r).
292        let global = &self.writeable.global();
293
294        self.writeable
295            .get_default_controller()
296            .error_if_needed(cx, v, global);
297
298        // Perform ! TransformStreamUnblockWrite(stream).
299        self.stream.unblock_write(cx, global);
300
301        // Reject controller.[[finishPromise]] with r.
302        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/// Reacting to fulfillment of the flushpromise as part of
312/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-close-algorithm>
313#[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    /// Reacting to flushpromise with the following fulfillment steps:
322    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
323        // If flushPromise was fulfilled, then:
324        let finish_promise = self
325            .controller
326            .get_finish_promise()
327            .expect("finish promise is not set");
328
329        // If readable.[[state]] is "errored", reject controller.[[finishPromise]] with readable.[[storedError]].
330        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            // Otherwise:
336            // Perform ! ReadableStreamDefaultControllerClose(readable.[[controller]]).
337            self.readable.get_default_controller().close(cx);
338
339            // Resolve controller.[[finishPromise]] with undefined.
340            finish_promise.resolve_native(cx, &());
341        }
342    }
343}
344
345impl js::gc::Rootable for FlushPromiseRejection {}
346/// Reacting to rejection of flushpromise as part of
347/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-close-algorithm>
348
349#[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    /// Reacting to flushpromise with the following fulfillment steps:
358    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
359        // If flushPromise was rejected with reason r, then:
360        // Perform ! ReadableStreamDefaultControllerError(readable.[[controller]], r).
361        self.readable.get_default_controller().error(cx, v);
362
363        // Reject controller.[[finishPromise]] with r.
364        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/// A wrapper to handle `message` and `messageerror` events
374/// for the message port used by the transfered stream.
375#[derive(Clone, JSTraceable, MallocSizeOf)]
376#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
377pub(crate) enum CrossRealmTransform {
378    /// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformreadable>
379    Readable(CrossRealmTransformReadable),
380    /// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformwritable>
381    Writable(CrossRealmTransformWritable),
382}
383
384/// <https://streams.spec.whatwg.org/#ts-class>
385#[dom_struct]
386pub struct TransformStream {
387    reflector_: Reflector,
388
389    /// <https://streams.spec.whatwg.org/#transformstream-backpressure>
390    backpressure: Cell<bool>,
391
392    /// <https://streams.spec.whatwg.org/#transformstream-backpressurechangepromise>
393    #[conditional_malloc_size_of]
394    backpressure_change_promise: DomRefCell<Option<Rc<Promise>>>,
395
396    /// <https://streams.spec.whatwg.org/#transformstream-controller>
397    controller: MutNullableDom<TransformStreamDefaultController>,
398
399    /// <https://streams.spec.whatwg.org/#transformstream-detached>
400    detached: Cell<bool>,
401
402    /// <https://streams.spec.whatwg.org/#transformstream-readable>
403    readable: MutNullableDom<ReadableStream>,
404
405    /// <https://streams.spec.whatwg.org/#transformstream-writable>
406    writable: MutNullableDom<WritableStream>,
407}
408
409impl TransformStream {
410    /// <https://streams.spec.whatwg.org/#initialize-transform-stream>
411    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    /// Creates and set up the newly created transform stream following
437    /// <https://streams.spec.whatwg.org/#transformstream-set-up>
438    #[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        // Step1. Let writableHighWaterMark be 1.
446        let writable_high_water_mark = 1.0;
447
448        // Step 2. Let writableSizeAlgorithm be an algorithm that returns 1.
449        let writable_size_algorithm = extract_size_algorithm(cx, &Default::default());
450
451        // Step 3. Let readableHighWaterMark be 0.
452        let readable_high_water_mark = 0.0;
453
454        // Step 4. Let readableSizeAlgorithm be an algorithm that returns 1.
455        let readable_size_algorithm = extract_size_algorithm(cx, &Default::default());
456
457        // Step 5. Let transformAlgorithmWrapper be an algorithm that runs these steps given a value chunk:
458        // Step 6. Let flushAlgorithmWrapper be an algorithm that runs these steps:
459        // Step 7. Let cancelAlgorithmWrapper be an algorithm that runs these steps given a value reason:
460        // NOTE: These steps are implemented in `TransformStreamDefaultController::new`
461
462        // Step 8. Let startPromise be a promise resolved with undefined.
463        let start_promise = Promise::new_resolved(cx, global, ());
464
465        // Step 9. Perform ! InitializeTransformStream(stream, startPromise,
466        // writableHighWaterMark, writableSizeAlgorithm, readableHighWaterMark,
467        // readableSizeAlgorithm).
468        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        // Step 10. Let controller be a new TransformStreamDefaultController.
479        let controller = TransformStreamDefaultController::new(cx, global, transformer_type);
480
481        // Step 11. Perform ! SetUpTransformStreamDefaultController(stream,
482        // controller, transformAlgorithmWrapper, flushAlgorithmWrapper,
483        // cancelAlgorithmWrapper).
484        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    /// <https://streams.spec.whatwg.org/#initialize-transform-stream>
506    #[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 startAlgorithm be an algorithm that returns startPromise.
518        // Let writeAlgorithm be the following steps, taking a chunk argument:
519        //  Return ! TransformStreamDefaultSinkWriteAlgorithm(stream, chunk).
520        // Let abortAlgorithm be the following steps, taking a reason argument:
521        //  Return ! TransformStreamDefaultSinkAbortAlgorithm(stream, reason).
522        // Let closeAlgorithm be the following steps:
523        //  Return ! TransformStreamDefaultSinkCloseAlgorithm(stream).
524        // Set stream.[[writable]] to ! CreateWritableStream(startAlgorithm, writeAlgorithm,
525        // closeAlgorithm, abortAlgorithm, writableHighWaterMark, writableSizeAlgorithm).
526        // Note: Those steps are implemented using UnderlyingSinkType::Transform.
527
528        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 pullAlgorithm be the following steps:
538
539        // Return ! TransformStreamDefaultSourcePullAlgorithm(stream).
540
541        // Let cancelAlgorithm be the following steps, taking a reason argument:
542
543        // Return ! TransformStreamDefaultSourceCancelAlgorithm(stream, reason).
544
545        // Set stream.[[readable]] to ! CreateReadableStream(startAlgorithm, pullAlgorithm,
546        // cancelAlgorithm, readableHighWaterMark, readableSizeAlgorithm).
547
548        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        // Set stream.[[backpressure]] and stream.[[backpressureChangePromise]] to undefined.
558        // Note: This is done in the constructor.
559
560        // Perform ! TransformStreamSetBackpressure(stream, true).
561        self.set_backpressure(cx, global, true);
562
563        // Set stream.[[controller]] to undefined.
564        self.controller.set(None);
565
566        Ok(())
567    }
568
569    /// <https://streams.spec.whatwg.org/#transform-stream-set-backpressure>
570    pub(crate) fn set_backpressure(
571        &self,
572        cx: &mut JSContext,
573        global: &GlobalScope,
574        backpressure: bool,
575    ) {
576        // Assert: stream.[[backpressure]] is not backpressure.
577        assert!(self.backpressure.get() != backpressure);
578
579        // If stream.[[backpressureChangePromise]] is not undefined, resolve
580        // stream.[[backpressureChangePromise]] with undefined.
581        if let Some(promise) = self.backpressure_change_promise.borrow_mut().take() {
582            promise.resolve_native(cx, &());
583        }
584
585        // Set stream.[[backpressureChangePromise]] to a new promise.;
586        *self.backpressure_change_promise.borrow_mut() = Some(Promise::new(cx, global));
587
588        // Set stream.[[backpressure]] to backpressure.
589        self.backpressure.set(backpressure);
590    }
591
592    /// <https://streams.spec.whatwg.org/#set-up-transform-stream-default-controller>
593    fn set_up_transform_stream_default_controller(
594        &self,
595        controller: &TransformStreamDefaultController,
596    ) {
597        // Assert: stream implements TransformStream.
598        // Note: this is checked with type.
599
600        // Assert: stream.[[controller]] is undefined.
601        assert!(self.controller.get().is_none());
602
603        // Set controller.[[stream]] to stream.
604        controller.set_stream(self);
605
606        // Set stream.[[controller]] to controller.
607        self.controller.set(Some(controller));
608
609        // Set controller.[[transformAlgorithm]] to transformAlgorithm.
610        // Set controller.[[flushAlgorithm]] to flushAlgorithm.
611        // Set controller.[[cancelAlgorithm]] to cancelAlgorithm.
612        // Note: These are set in the constructor.
613    }
614
615    /// <https://streams.spec.whatwg.org/#set-up-transform-stream-default-controller-from-transformer>
616    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 be a new TransformStreamDefaultController.
624        let controller = TransformStreamDefaultController::new(
625            cx,
626            global,
627            TransformerType::new_from_js_transformer(transformer),
628        );
629
630        // Let transformAlgorithm be the following steps, taking a chunk argument:
631        // Let result be TransformStreamDefaultControllerEnqueue(controller, chunk).
632        // If result is an abrupt completion, return a promise rejected with result.[[Value]].
633        // Otherwise, return a promise resolved with undefined.
634
635        // Let flushAlgorithm be an algorithm which returns a promise resolved with undefined.
636        // Let cancelAlgorithm be an algorithm which returns a promise resolved with undefined.
637
638        // If transformerDict["transform"] exists, set transformAlgorithm to an algorithm which
639        // takes an argument
640        // chunk and returns the result of invoking transformerDict["transform"] with argument
641        // list « chunk, controller »
642        // and callback this value transformer.
643
644        // If transformerDict["flush"] exists, set flushAlgorithm to an algorithm which returns
645        // the result
646        // of invoking transformerDict["flush"] with argument list « controller » and callback
647        // this value transformer.
648
649        // If transformerDict["cancel"] exists, set cancelAlgorithm to an algorithm which takes an argument
650        // reason and returns the result of invoking transformerDict["cancel"] with argument list « reason »
651        // and callback this value transformer.
652        controller.set_transform_obj(transformer_obj);
653
654        // Perform ! SetUpTransformStreamDefaultController(stream, controller,
655        // transformAlgorithm, flushAlgorithm, cancelAlgorithm).
656        self.set_up_transform_stream_default_controller(&controller);
657    }
658
659    /// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
660    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: stream.[[writable]].[[state]] is "writable".
667        assert!(self.writable.get().is_some());
668
669        // Let controller be stream.[[controller]].
670        let controller = self.controller.get().expect("controller is not set");
671
672        // If stream.[[backpressure]] is true,
673        if self.backpressure.get() {
674            // Let backpressureChangePromise be stream.[[backpressureChangePromise]].
675            let backpressure_change_promise = self.backpressure_change_promise.borrow();
676
677            // Assert: backpressureChangePromise is not undefined.
678            assert!(backpressure_change_promise.is_some());
679
680            // Return the result of reacting to backpressureChangePromise with the following fulfillment steps:
681            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        // Return ! TransformStreamDefaultControllerPerformTransform(controller, chunk).
708        controller.transform_stream_default_controller_perform_transform(cx, global, chunk)
709    }
710
711    /// <https://streams.spec.whatwg.org/#transform-stream-default-sink-abort-algorithm>
712    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 be stream.[[controller]].
719        let controller = self.controller.get().expect("controller is not set");
720
721        // If controller.[[finishPromise]] is not undefined, return controller.[[finishPromise]].
722        if let Some(finish_promise) = controller.get_finish_promise() {
723            return Ok(finish_promise);
724        }
725
726        // Let readable be stream.[[readable]].
727        let readable = self.readable.get().expect("readable stream is not set");
728
729        // Let controller.[[finishPromise]] be a new promise.
730        controller.set_finish_promise(Promise::new(cx, global));
731
732        // Let cancelPromise be the result of performing controller.[[cancelAlgorithm]], passing reason.
733        let cancel_promise = controller.perform_cancel(cx, global, reason)?;
734
735        // Perform ! TransformStreamDefaultControllerClearAlgorithms(controller).
736        controller.clear_algorithms();
737
738        // React to cancelPromise:
739        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        // Return controller.[[finishPromise]].
757        let finish_promise = controller
758            .get_finish_promise()
759            .expect("finish promise is not set");
760        Ok(finish_promise)
761    }
762
763    /// <https://streams.spec.whatwg.org/#transform-stream-default-sink-close-algorithm>
764    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 be stream.[[controller]].
770        let controller = self
771            .controller
772            .get()
773            .ok_or(Error::Type(c"controller is not set".to_owned()))?;
774
775        // If controller.[[finishPromise]] is not undefined, return controller.[[finishPromise]].
776        if let Some(finish_promise) = controller.get_finish_promise() {
777            return Ok(finish_promise);
778        }
779
780        // Let readable be stream.[[readable]].
781        let readable = self
782            .readable
783            .get()
784            .ok_or(Error::Type(c"readable stream is not set".to_owned()))?;
785
786        // Let controller.[[finishPromise]] be a new promise.
787        controller.set_finish_promise(Promise::new(cx, global));
788
789        // Let flushPromise be the result of performing controller.[[flushAlgorithm]].
790        let flush_promise = controller.perform_flush(cx, global)?;
791
792        // Perform ! TransformStreamDefaultControllerClearAlgorithms(controller).
793        controller.clear_algorithms();
794
795        // React to flushPromise:
796        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        // Return controller.[[finishPromise]].
813        let finish_promise = controller
814            .get_finish_promise()
815            .expect("finish promise is not set");
816        Ok(finish_promise)
817    }
818
819    /// <https://streams.spec.whatwg.org/#transform-stream-default-source-cancel>
820    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 be stream.[[controller]].
827        let controller = self
828            .controller
829            .get()
830            .ok_or(Error::Type(c"controller is not set".to_owned()))?;
831
832        // If controller.[[finishPromise]] is not undefined, return controller.[[finishPromise]].
833        if let Some(finish_promise) = controller.get_finish_promise() {
834            return Ok(finish_promise);
835        }
836
837        // Let writable be stream.[[writable]].
838        let writable = self
839            .writable
840            .get()
841            .ok_or(Error::Type(c"writable stream is not set".to_owned()))?;
842
843        // Let controller.[[finishPromise]] be a new promise.
844        controller.set_finish_promise(Promise::new(cx, global));
845
846        // Let cancelPromise be the result of performing controller.[[cancelAlgorithm]], passing reason.
847        let cancel_promise = controller.perform_cancel(cx, global, reason)?;
848
849        // Perform ! TransformStreamDefaultControllerClearAlgorithms(controller).
850        controller.clear_algorithms();
851
852        // React to cancelPromise:
853        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        // Return controller.[[finishPromise]].
870        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    /// <https://streams.spec.whatwg.org/#transform-stream-default-source-pull>
880    pub(crate) fn transform_stream_default_source_pull(
881        &self,
882        cx: &mut JSContext,
883        global: &GlobalScope,
884    ) -> Fallible<Rc<Promise>> {
885        // Assert: stream.[[backpressure]] is true.
886        assert!(self.backpressure.get());
887
888        // Assert: stream.[[backpressureChangePromise]] is not undefined.
889        assert!(self.backpressure_change_promise.borrow().is_some());
890
891        // Perform ! TransformStreamSetBackpressure(stream, false).
892        self.set_backpressure(cx, global, false);
893
894        // Return stream.[[backpressureChangePromise]].
895        Ok(self
896            .backpressure_change_promise
897            .borrow()
898            .clone()
899            .expect("Promise must be some by now."))
900    }
901
902    /// <https://streams.spec.whatwg.org/#transform-stream-error-writable-and-unblock-write>
903    pub(crate) fn error_writable_and_unblock_write(
904        &self,
905        cx: &mut JSContext,
906        global: &GlobalScope,
907        error: SafeHandleValue,
908    ) {
909        // Perform ! TransformStreamDefaultControllerClearAlgorithms(stream.[[controller]]).
910        self.get_controller().clear_algorithms();
911
912        // Perform ! WritableStreamDefaultControllerErrorIfNeeded(stream.[[writable]].[[controller]], e).
913        self.get_writable()
914            .get_default_controller()
915            .error_if_needed(cx, error, global);
916
917        // Perform ! TransformStreamUnblockWrite(stream).
918        self.unblock_write(cx, global)
919    }
920
921    /// <https://streams.spec.whatwg.org/#transform-stream-unblock-write>
922    pub(crate) fn unblock_write(&self, cx: &mut JSContext, global: &GlobalScope) {
923        // If stream.[[backpressure]] is true, perform ! TransformStreamSetBackpressure(stream, false).
924        if self.backpressure.get() {
925            self.set_backpressure(cx, global, false);
926        }
927    }
928
929    /// <https://streams.spec.whatwg.org/#transform-stream-error>
930    pub(crate) fn error(&self, cx: &mut JSContext, global: &GlobalScope, error: SafeHandleValue) {
931        // Perform ! ReadableStreamDefaultControllerError(stream.[[readable]].[[controller]], e).
932        self.get_readable()
933            .get_default_controller()
934            .error(cx, error);
935
936        // Perform ! TransformStreamErrorWritableAndUnblockWrite(stream, e).
937        self.error_writable_and_unblock_write(cx, global, error);
938    }
939}
940
941impl TransformStreamMethods<crate::DomTypeHolder> for TransformStream {
942    /// <https://streams.spec.whatwg.org/#ts-constructor>
943    #[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        // If transformer is missing, set it to null.
953        rooted!(&in(cx) let transformer_obj = transformer.unwrap_or(ptr::null_mut()));
954
955        // Let underlyingSinkDict be underlyingSink,
956        // converted to an IDL value of type UnderlyingSink.
957        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 transformerDict["readableType"] exists, throw a RangeError exception.
973        if !transformer_dict.readableType.handle().is_undefined() {
974            return Err(Error::Range(c"readableType is set".to_owned()));
975        }
976
977        // If transformerDict["writableType"] exists, throw a RangeError exception.
978        if !transformer_dict.writableType.handle().is_undefined() {
979            return Err(Error::Range(c"writableType is set".to_owned()));
980        }
981
982        // Let readableHighWaterMark be ? ExtractHighWaterMark(readableStrategy, 0).
983        let readable_high_water_mark = extract_high_water_mark(readable_strategy, 0.0)?;
984
985        // Let readableSizeAlgorithm be ! ExtractSizeAlgorithm(readableStrategy).
986        let readable_size_algorithm = extract_size_algorithm(cx, readable_strategy);
987
988        // Let writableHighWaterMark be ? ExtractHighWaterMark(writableStrategy, 1).
989        let writable_high_water_mark = extract_high_water_mark(writable_strategy, 1.0)?;
990
991        // Let writableSizeAlgorithm be ! ExtractSizeAlgorithm(writableStrategy).
992        let writable_size_algorithm = extract_size_algorithm(cx, writable_strategy);
993
994        // Let startPromise be a new promise.
995        let start_promise = Promise::new(cx, global);
996
997        // Perform ! InitializeTransformStream(this, startPromise, writableHighWaterMark,
998        // writableSizeAlgorithm, readableHighWaterMark, readableSizeAlgorithm).
999        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        // Perform ? SetUpTransformStreamDefaultControllerFromTransformer(this, transformer, transformerDict).
1011        stream.set_up_transform_stream_default_controller_from_transformer(
1012            cx,
1013            global,
1014            transformer_obj.handle(),
1015            &transformer_dict,
1016        );
1017
1018        // If transformerDict["start"] exists, then resolve startPromise with the
1019        // result of invoking transformerDict["start"]
1020        // with argument list « this.[[controller]] » and callback this value transformer.
1021        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            // Otherwise, resolve startPromise with undefined.
1048            start_promise.resolve_native(cx, &());
1049        };
1050
1051        Ok(stream)
1052    }
1053
1054    /// <https://streams.spec.whatwg.org/#ts-readable>
1055    fn Readable(&self) -> DomRoot<ReadableStream> {
1056        // Return this.[[readable]].
1057        self.readable.get().expect("readable stream is not set")
1058    }
1059
1060    /// <https://streams.spec.whatwg.org/#ts-writable>
1061    fn Writable(&self) -> DomRoot<WritableStream> {
1062        // Return this.[[writable]].
1063        self.writable.get().expect("writable stream is not set")
1064    }
1065}
1066
1067/// <https://streams.spec.whatwg.org/#ts-transfer>
1068impl Transferable for TransformStream {
1069    type Index = MessagePortIndex;
1070    type Data = TransformStreamData;
1071
1072    /// <https://streams.spec.whatwg.org/#ref-for-transfer-steps②>
1073    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        // Step 1. Let readable be value.[[readable]].
1080        let readable = self.get_readable();
1081
1082        // Step 2. Let writable be value.[[writable]].
1083        let writable = self.get_writable();
1084
1085        // Step 3. If ! IsReadableStreamLocked(readable) is true, throw a
1086        // "DataCloneError" DOMException.
1087        // Step 4. If ! IsWritableStreamLocked(writable) is true, throw a
1088        // "DataCloneError" DOMException.
1089        if readable.is_locked() || writable.is_locked() {
1090            return Err(Error::DataClone(None));
1091        }
1092
1093        // First port pair (readable → proxy writable)
1094        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        // Second port pair (proxy readable → writable)
1107        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        // Pipe readable into the proxy writable (→ port_1)
1117        readable
1118            .pipe_to(cx, &global, &proxy_writable, false, false, false, None)
1119            .set_promise_is_handled(cx);
1120
1121        // Step 5. Set dataHolder.[[readable]] to !
1122        // StructuredSerializeWithTransfer(readable, « readable »).
1123        // Step 6. Set dataHolder.[[writable]] to !
1124        // StructuredSerializeWithTransfer(writable, « writable »).
1125        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    /// <https://streams.spec.whatwg.org/#ref-for-transfer-receiving-steps②>
1135    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        // Step 1. Let readableRecord be !
1145        // StructuredDeserializeWithTransfer(dataHolder.[[readable]], the
1146        // current Realm).
1147        let proxy_readable = ReadableStream::new_with_proto(cx, owner, None);
1148        proxy_readable.setup_cross_realm_transform_readable(cx, &port2);
1149
1150        // Step 2. Let writableRecord be !
1151        // StructuredDeserializeWithTransfer(dataHolder.[[writable]], the
1152        // current Realm).
1153        let proxy_writable = WritableStream::new_with_proto(cx, owner, None);
1154        proxy_writable.setup_cross_realm_transform_writable(cx, &port1);
1155
1156        // Step 3. Set value.[[readable]] to readableRecord.[[Deserialized]].
1157        // Step 4. Set value.[[writable]] to writableRecord.[[Deserialized]].
1158        // Step 5. Set value.[[backpressure]],
1159        // value.[[backpressureChangePromise]], and value.[[controller]] to
1160        // undefined.
1161        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}