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, 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/// 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    result_promise: TracedPromise,
57
58    #[ignore_malloc_size_of = "mozjs"]
59    chunk: Box<Heap<JSVal>>,
60
61    /// The writable used in the fulfillment steps
62    writable: Dom<WritableStream>,
63
64    controller: Dom<TransformStreamDefaultController>,
65}
66
67impl Callback for TransformBackPressureChangePromiseFulfillment {
68    /// Reacting to backpressureChangePromise with the following fulfillment steps:
69    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
70        // Let writable be stream.[[writable]].
71        // Let state be writable.[[state]].
72        // If state is "erroring", throw writable.[[storedError]].
73        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: state is "writable".
81        assert!(self.writable.is_writable());
82
83        // Return ! TransformStreamDefaultControllerPerformTransform(controller, chunk).
84        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        // PerformTransformFulfillment and PerformTransformRejection do not need
96        // to be rooted because they only contain an Rc.
97        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)]
116/// Reacting to fulfillment of performTransform as part of
117/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
118struct PerformTransformFulfillment {
119    result_promise: TracedPromise,
120}
121
122impl Callback for PerformTransformFulfillment {
123    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
124        // Fulfilled: resolve the outer promise
125        self.result_promise.resolve_native(cx, &());
126    }
127}
128
129#[derive(JSTraceable, MallocSizeOf)]
130#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
131/// Reacting to rejection of performTransform as part of
132/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
133struct PerformTransformRejection {
134    result_promise: TracedPromise,
135}
136
137impl Callback for PerformTransformRejection {
138    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
139        // Stream already errored in perform_transform, just reject result_promise
140        self.result_promise.reject(cx, v);
141    }
142}
143
144#[derive(JSTraceable, MallocSizeOf)]
145#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
146/// Reacting to rejection of backpressureChangePromise as part of
147/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
148struct 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/// Reacting to fulfillment of the cancelpromise as part of
161/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-abort-algorithm>
162#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
173    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
174        // If readable.[[state]] is "errored", reject controller.[[finishPromise]] with readable.[[storedError]].
175        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            // Otherwise:
184            // Perform ! ReadableStreamDefaultControllerError(readable.[[controller]], reason).
185            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            // Resolve controller.[[finishPromise]] with undefined.
192            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/// Reacting to rejection of cancelpromise as part of
203/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-abort-algorithm>
204#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
213    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
214        // Perform ! ReadableStreamDefaultControllerError(readable.[[controller]], r).
215        self.readable.get_default_controller().error(cx, v);
216
217        // Reject controller.[[finishPromise]] with r.
218        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/// Reacting to fulfillment of the cancelpromise as part of
228/// <https://streams.spec.whatwg.org/#transform-stream-default-source-cancel>
229#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
241    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
242        // If cancelPromise was fulfilled, then:
243        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 writable.[[state]] is "errored", reject controller.[[finishPromise]] with writable.[[storedError]].
250        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            // Otherwise:
256            // Perform ! WritableStreamDefaultControllerErrorIfNeeded(writable.[[controller]], reason).
257            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            // Perform ! TransformStreamUnblockWrite(stream).
264            self.stream.unblock_write(cx, global);
265
266            // Resolve controller.[[finishPromise]] with undefined.
267            finish_promise.resolve_native(cx, &());
268        }
269    }
270}
271
272impl js::gc::Rootable for SourceCancelPromiseRejection {}
273
274/// Reacting to rejection of cancelpromise as part of
275/// <https://streams.spec.whatwg.org/#transform-stream-default-source-cancel>
276#[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    /// Reacting to backpressureChangePromise with the following fulfillment steps:
286    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
287        // Perform ! WritableStreamDefaultControllerErrorIfNeeded(writable.[[controller]], r).
288        let global = &self.writeable.global();
289
290        self.writeable
291            .get_default_controller()
292            .error_if_needed(cx, v, global);
293
294        // Perform ! TransformStreamUnblockWrite(stream).
295        self.stream.unblock_write(cx, global);
296
297        // Reject controller.[[finishPromise]] with r.
298        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/// Reacting to fulfillment of the flushpromise as part of
308/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-close-algorithm>
309#[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    /// Reacting to flushpromise with the following fulfillment steps:
318    fn callback(&self, cx: &mut CurrentRealm, _v: SafeHandleValue) {
319        // If flushPromise was fulfilled, then:
320        let finish_promise = self
321            .controller
322            .get_finish_promise(cx)
323            .expect("finish promise is not set");
324
325        // If readable.[[state]] is "errored", reject controller.[[finishPromise]] with readable.[[storedError]].
326        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            // Otherwise:
332            // Perform ! ReadableStreamDefaultControllerClose(readable.[[controller]]).
333            self.readable.get_default_controller().close(cx);
334
335            // Resolve controller.[[finishPromise]] with undefined.
336            finish_promise.resolve_native(cx, &());
337        }
338    }
339}
340
341impl js::gc::Rootable for FlushPromiseRejection {}
342/// Reacting to rejection of flushpromise as part of
343/// <https://streams.spec.whatwg.org/#transform-stream-default-sink-close-algorithm>
344
345#[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    /// Reacting to flushpromise with the following fulfillment steps:
354    fn callback(&self, cx: &mut CurrentRealm, v: SafeHandleValue) {
355        // If flushPromise was rejected with reason r, then:
356        // Perform ! ReadableStreamDefaultControllerError(readable.[[controller]], r).
357        self.readable.get_default_controller().error(cx, v);
358
359        // Reject controller.[[finishPromise]] with r.
360        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/// A wrapper to handle `message` and `messageerror` events
370/// for the message port used by the transfered stream.
371#[derive(Clone, JSTraceable, MallocSizeOf)]
372#[cfg_attr(crown, crown::unrooted_must_root_lint::must_root)]
373pub(crate) enum CrossRealmTransform {
374    /// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformreadable>
375    Readable(CrossRealmTransformReadable),
376    /// <https://streams.spec.whatwg.org/#abstract-opdef-setupcrossrealmtransformwritable>
377    Writable(CrossRealmTransformWritable),
378}
379
380/// <https://streams.spec.whatwg.org/#ts-class>
381#[dom_struct]
382pub struct TransformStream {
383    reflector_: Reflector,
384
385    /// <https://streams.spec.whatwg.org/#transformstream-backpressure>
386    backpressure: Cell<bool>,
387
388    /// <https://streams.spec.whatwg.org/#transformstream-backpressurechangepromise>
389    backpressure_change_promise: DomRefCell<Option<TracedPromise>>,
390
391    /// <https://streams.spec.whatwg.org/#transformstream-controller>
392    controller: MutNullableDom<TransformStreamDefaultController>,
393
394    /// <https://streams.spec.whatwg.org/#transformstream-detached>
395    detached: Cell<bool>,
396
397    /// <https://streams.spec.whatwg.org/#transformstream-readable>
398    readable: MutNullableDom<ReadableStream>,
399
400    /// <https://streams.spec.whatwg.org/#transformstream-writable>
401    writable: MutNullableDom<WritableStream>,
402}
403
404impl TransformStream {
405    /// <https://streams.spec.whatwg.org/#initialize-transform-stream>
406    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    /// Creates and set up the newly created transform stream following
432    /// <https://streams.spec.whatwg.org/#transformstream-set-up>
433    #[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        // Step1. Let writableHighWaterMark be 1.
441        let writable_high_water_mark = 1.0;
442
443        // Step 2. Let writableSizeAlgorithm be an algorithm that returns 1.
444        let writable_size_algorithm = extract_size_algorithm(cx, &Default::default());
445
446        // Step 3. Let readableHighWaterMark be 0.
447        let readable_high_water_mark = 0.0;
448
449        // Step 4. Let readableSizeAlgorithm be an algorithm that returns 1.
450        let readable_size_algorithm = extract_size_algorithm(cx, &Default::default());
451
452        // Step 5. Let transformAlgorithmWrapper be an algorithm that runs these steps given a value chunk:
453        // Step 6. Let flushAlgorithmWrapper be an algorithm that runs these steps:
454        // Step 7. Let cancelAlgorithmWrapper be an algorithm that runs these steps given a value reason:
455        // NOTE: These steps are implemented in `TransformStreamDefaultController::new`
456
457        // Step 8. Let startPromise be a promise resolved with undefined.
458        let start_promise = Promise::new_resolved_rooted(cx, global, ());
459
460        // Step 9. Perform ! InitializeTransformStream(stream, startPromise,
461        // writableHighWaterMark, writableSizeAlgorithm, readableHighWaterMark,
462        // readableSizeAlgorithm).
463        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        // Step 10. Let controller be a new TransformStreamDefaultController.
474        let controller = TransformStreamDefaultController::new(cx, global, transformer_type);
475
476        // Step 11. Perform ! SetUpTransformStreamDefaultController(stream,
477        // controller, transformAlgorithmWrapper, flushAlgorithmWrapper,
478        // cancelAlgorithmWrapper).
479        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    /// <https://streams.spec.whatwg.org/#initialize-transform-stream>
501    #[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 startAlgorithm be an algorithm that returns startPromise.
513        // Let writeAlgorithm be the following steps, taking a chunk argument:
514        //  Return ! TransformStreamDefaultSinkWriteAlgorithm(stream, chunk).
515        // Let abortAlgorithm be the following steps, taking a reason argument:
516        //  Return ! TransformStreamDefaultSinkAbortAlgorithm(stream, reason).
517        // Let closeAlgorithm be the following steps:
518        //  Return ! TransformStreamDefaultSinkCloseAlgorithm(stream).
519        // Set stream.[[writable]] to ! CreateWritableStream(startAlgorithm, writeAlgorithm,
520        // closeAlgorithm, abortAlgorithm, writableHighWaterMark, writableSizeAlgorithm).
521        // Note: Those steps are implemented using UnderlyingSinkType::Transform.
522
523        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 pullAlgorithm be the following steps:
533
534        // Return ! TransformStreamDefaultSourcePullAlgorithm(stream).
535
536        // Let cancelAlgorithm be the following steps, taking a reason argument:
537
538        // Return ! TransformStreamDefaultSourceCancelAlgorithm(stream, reason).
539
540        // Set stream.[[readable]] to ! CreateReadableStream(startAlgorithm, pullAlgorithm,
541        // cancelAlgorithm, readableHighWaterMark, readableSizeAlgorithm).
542
543        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        // Set stream.[[backpressure]] and stream.[[backpressureChangePromise]] to undefined.
553        // Note: This is done in the constructor.
554
555        // Perform ! TransformStreamSetBackpressure(stream, true).
556        self.set_backpressure(cx, global, true);
557
558        // Set stream.[[controller]] to undefined.
559        self.controller.set(None);
560
561        Ok(())
562    }
563
564    /// <https://streams.spec.whatwg.org/#transform-stream-set-backpressure>
565    pub(crate) fn set_backpressure(
566        &self,
567        cx: &mut JSContext,
568        global: &GlobalScope,
569        backpressure: bool,
570    ) {
571        // Assert: stream.[[backpressure]] is not backpressure.
572        assert!(self.backpressure.get() != backpressure);
573
574        // If stream.[[backpressureChangePromise]] is not undefined, resolve
575        // stream.[[backpressureChangePromise]] with undefined.
576        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        // Set stream.[[backpressureChangePromise]] to a new promise.;
582        *self.backpressure_change_promise.borrow_mut() =
583            Some(Promise::new_rooted(cx, global).to_traced());
584
585        // Set stream.[[backpressure]] to backpressure.
586        self.backpressure.set(backpressure);
587    }
588
589    /// <https://streams.spec.whatwg.org/#set-up-transform-stream-default-controller>
590    fn set_up_transform_stream_default_controller(
591        &self,
592        controller: &TransformStreamDefaultController,
593    ) {
594        // Assert: stream implements TransformStream.
595        // Note: this is checked with type.
596
597        // Assert: stream.[[controller]] is undefined.
598        assert!(self.controller.get().is_none());
599
600        // Set controller.[[stream]] to stream.
601        controller.set_stream(self);
602
603        // Set stream.[[controller]] to controller.
604        self.controller.set(Some(controller));
605
606        // Set controller.[[transformAlgorithm]] to transformAlgorithm.
607        // Set controller.[[flushAlgorithm]] to flushAlgorithm.
608        // Set controller.[[cancelAlgorithm]] to cancelAlgorithm.
609        // Note: These are set in the constructor.
610    }
611
612    /// <https://streams.spec.whatwg.org/#set-up-transform-stream-default-controller-from-transformer>
613    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 be a new TransformStreamDefaultController.
621        let controller = TransformStreamDefaultController::new(
622            cx,
623            global,
624            TransformerType::new_from_js_transformer(transformer),
625        );
626
627        // Let transformAlgorithm be the following steps, taking a chunk argument:
628        // Let result be TransformStreamDefaultControllerEnqueue(controller, chunk).
629        // If result is an abrupt completion, return a promise rejected with result.[[Value]].
630        // Otherwise, return a promise resolved with undefined.
631
632        // Let flushAlgorithm be an algorithm which returns a promise resolved with undefined.
633        // Let cancelAlgorithm be an algorithm which returns a promise resolved with undefined.
634
635        // If transformerDict["transform"] exists, set transformAlgorithm to an algorithm which
636        // takes an argument
637        // chunk and returns the result of invoking transformerDict["transform"] with argument
638        // list « chunk, controller »
639        // and callback this value transformer.
640
641        // If transformerDict["flush"] exists, set flushAlgorithm to an algorithm which returns
642        // the result
643        // of invoking transformerDict["flush"] with argument list « controller » and callback
644        // this value transformer.
645
646        // If transformerDict["cancel"] exists, set cancelAlgorithm to an algorithm which takes an argument
647        // reason and returns the result of invoking transformerDict["cancel"] with argument list « reason »
648        // and callback this value transformer.
649        controller.set_transform_obj(transformer_obj);
650
651        // Perform ! SetUpTransformStreamDefaultController(stream, controller,
652        // transformAlgorithm, flushAlgorithm, cancelAlgorithm).
653        self.set_up_transform_stream_default_controller(&controller);
654    }
655
656    /// <https://streams.spec.whatwg.org/#transform-stream-default-sink-write-algorithm>
657    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: stream.[[writable]].[[state]] is "writable".
664        assert!(self.writable.get().is_some());
665
666        // Let controller be stream.[[controller]].
667        let controller = self.controller.get().expect("controller is not set");
668
669        // If stream.[[backpressure]] is true,
670        if self.backpressure.get() {
671            // Let backpressureChangePromise be stream.[[backpressureChangePromise]].
672            let backpressure_change_promise = self.backpressure_change_promise.borrow();
673
674            // Assert: backpressureChangePromise is not undefined.
675            assert!(backpressure_change_promise.is_some());
676
677            // Return the result of reacting to backpressureChangePromise with the following fulfillment steps:
678            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        // Return ! TransformStreamDefaultControllerPerformTransform(controller, chunk).
705        controller.transform_stream_default_controller_perform_transform(cx, global, chunk)
706    }
707
708    /// <https://streams.spec.whatwg.org/#transform-stream-default-sink-abort-algorithm>
709    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 be stream.[[controller]].
716        let controller = self.controller.get().expect("controller is not set");
717
718        // If controller.[[finishPromise]] is not undefined, return controller.[[finishPromise]].
719        if let Some(finish_promise) = controller.get_finish_promise(cx) {
720            return Ok(finish_promise);
721        }
722
723        // Let readable be stream.[[readable]].
724        let readable = self.readable.get().expect("readable stream is not set");
725
726        // Let controller.[[finishPromise]] be a new promise.
727        controller.set_finish_promise(&Promise::new_rooted(cx, global));
728
729        // Let cancelPromise be the result of performing controller.[[cancelAlgorithm]], passing reason.
730        let cancel_promise = controller.perform_cancel(cx, global, reason)?;
731
732        // Perform ! TransformStreamDefaultControllerClearAlgorithms(controller).
733        controller.clear_algorithms();
734
735        // React to cancelPromise:
736        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        // Return controller.[[finishPromise]].
754        let finish_promise = controller
755            .get_finish_promise(cx)
756            .expect("finish promise is not set");
757        Ok(finish_promise)
758    }
759
760    /// <https://streams.spec.whatwg.org/#transform-stream-default-sink-close-algorithm>
761    pub(crate) fn transform_stream_default_sink_close_algorithm(
762        &self,
763        cx: &mut JSContext,
764        global: &GlobalScope,
765    ) -> Fallible<RootedPromise> {
766        // Let controller be stream.[[controller]].
767        let controller = self
768            .controller
769            .get()
770            .ok_or(Error::Type(c"controller is not set".to_owned()))?;
771
772        // If controller.[[finishPromise]] is not undefined, return controller.[[finishPromise]].
773        if let Some(finish_promise) = controller.get_finish_promise(cx) {
774            return Ok(finish_promise);
775        }
776
777        // Let readable be stream.[[readable]].
778        let readable = self
779            .readable
780            .get()
781            .ok_or(Error::Type(c"readable stream is not set".to_owned()))?;
782
783        // Let controller.[[finishPromise]] be a new promise.
784        controller.set_finish_promise(&Promise::new_rooted(cx, global));
785
786        // Let flushPromise be the result of performing controller.[[flushAlgorithm]].
787        let flush_promise = controller.perform_flush(cx, global)?;
788
789        // Perform ! TransformStreamDefaultControllerClearAlgorithms(controller).
790        controller.clear_algorithms();
791
792        // React to flushPromise:
793        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        // Return controller.[[finishPromise]].
813        let finish_promise = controller
814            .get_finish_promise(cx)
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<RootedPromise> {
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(cx) {
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_rooted(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        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        // Return controller.[[finishPromise]].
874        let finish_promise = controller
875            .get_finish_promise(cx)
876            .expect("finish promise is not set");
877        Ok(finish_promise)
878    }
879
880    /// <https://streams.spec.whatwg.org/#transform-stream-default-source-pull>
881    pub(crate) fn transform_stream_default_source_pull(
882        &self,
883        cx: &mut JSContext,
884        global: &GlobalScope,
885    ) -> Fallible<RootedPromise> {
886        // Assert: stream.[[backpressure]] is true.
887        assert!(self.backpressure.get());
888
889        // Assert: stream.[[backpressureChangePromise]] is not undefined.
890        assert!(self.backpressure_change_promise.borrow().is_some());
891
892        // Perform ! TransformStreamSetBackpressure(stream, false).
893        self.set_backpressure(cx, global, false);
894
895        // Return stream.[[backpressureChangePromise]].
896        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    /// <https://streams.spec.whatwg.org/#transform-stream-error-writable-and-unblock-write>
905    pub(crate) fn error_writable_and_unblock_write(
906        &self,
907        cx: &mut JSContext,
908        global: &GlobalScope,
909        error: SafeHandleValue,
910    ) {
911        // Perform ! TransformStreamDefaultControllerClearAlgorithms(stream.[[controller]]).
912        self.get_controller().clear_algorithms();
913
914        // Perform ! WritableStreamDefaultControllerErrorIfNeeded(stream.[[writable]].[[controller]], e).
915        self.get_writable()
916            .get_default_controller()
917            .error_if_needed(cx, error, global);
918
919        // Perform ! TransformStreamUnblockWrite(stream).
920        self.unblock_write(cx, global)
921    }
922
923    /// <https://streams.spec.whatwg.org/#transform-stream-unblock-write>
924    pub(crate) fn unblock_write(&self, cx: &mut JSContext, global: &GlobalScope) {
925        // If stream.[[backpressure]] is true, perform ! TransformStreamSetBackpressure(stream, false).
926        if self.backpressure.get() {
927            self.set_backpressure(cx, global, false);
928        }
929    }
930
931    /// <https://streams.spec.whatwg.org/#transform-stream-error>
932    pub(crate) fn error(&self, cx: &mut JSContext, global: &GlobalScope, error: SafeHandleValue) {
933        // Perform ! ReadableStreamDefaultControllerError(stream.[[readable]].[[controller]], e).
934        self.get_readable()
935            .get_default_controller()
936            .error(cx, error);
937
938        // Perform ! TransformStreamErrorWritableAndUnblockWrite(stream, e).
939        self.error_writable_and_unblock_write(cx, global, error);
940    }
941}
942
943impl TransformStreamMethods<crate::DomTypeHolder> for TransformStream {
944    /// <https://streams.spec.whatwg.org/#ts-constructor>
945    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        // If transformer is missing, set it to null.
954        rooted!(&in(cx) let transformer_obj = transformer.unwrap_or(ptr::null_mut()));
955
956        // Let underlyingSinkDict be underlyingSink,
957        // converted to an IDL value of type UnderlyingSink.
958        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 transformerDict["readableType"] exists, throw a RangeError exception.
974        if !transformer_dict.readableType.handle().is_undefined() {
975            return Err(Error::Range(c"readableType is set".to_owned()));
976        }
977
978        // If transformerDict["writableType"] exists, throw a RangeError exception.
979        if !transformer_dict.writableType.handle().is_undefined() {
980            return Err(Error::Range(c"writableType is set".to_owned()));
981        }
982
983        // Let readableHighWaterMark be ? ExtractHighWaterMark(readableStrategy, 0).
984        let readable_high_water_mark = extract_high_water_mark(readable_strategy, 0.0)?;
985
986        // Let readableSizeAlgorithm be ! ExtractSizeAlgorithm(readableStrategy).
987        let readable_size_algorithm = extract_size_algorithm(cx, readable_strategy);
988
989        // Let writableHighWaterMark be ? ExtractHighWaterMark(writableStrategy, 1).
990        let writable_high_water_mark = extract_high_water_mark(writable_strategy, 1.0)?;
991
992        // Let writableSizeAlgorithm be ! ExtractSizeAlgorithm(writableStrategy).
993        let writable_size_algorithm = extract_size_algorithm(cx, writable_strategy);
994
995        // Let startPromise be a new promise.
996        let start_promise = Promise::new_rooted(cx, global);
997
998        // Perform ! InitializeTransformStream(this, startPromise, writableHighWaterMark,
999        // writableSizeAlgorithm, readableHighWaterMark, readableSizeAlgorithm).
1000        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        // Perform ? SetUpTransformStreamDefaultControllerFromTransformer(this, transformer, transformerDict).
1012        stream.set_up_transform_stream_default_controller_from_transformer(
1013            cx,
1014            global,
1015            transformer_obj.handle(),
1016            &transformer_dict,
1017        );
1018
1019        // If transformerDict["start"] exists, then resolve startPromise with the
1020        // result of invoking transformerDict["start"]
1021        // with argument list « this.[[controller]] » and callback this value transformer.
1022        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            // Otherwise, resolve startPromise with undefined.
1036            start_promise.resolve_native(cx, &());
1037        };
1038
1039        Ok(stream)
1040    }
1041
1042    /// <https://streams.spec.whatwg.org/#ts-readable>
1043    fn Readable(&self) -> DomRoot<ReadableStream> {
1044        // Return this.[[readable]].
1045        self.readable.get().expect("readable stream is not set")
1046    }
1047
1048    /// <https://streams.spec.whatwg.org/#ts-writable>
1049    fn Writable(&self) -> DomRoot<WritableStream> {
1050        // Return this.[[writable]].
1051        self.writable.get().expect("writable stream is not set")
1052    }
1053}
1054
1055/// <https://streams.spec.whatwg.org/#ts-transfer>
1056impl Transferable for TransformStream {
1057    type Index = MessagePortIndex;
1058    type Data = TransformStreamData;
1059
1060    /// <https://streams.spec.whatwg.org/#ref-for-transfer-steps②>
1061    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        // Step 1. Let readable be value.[[readable]].
1068        let readable = self.get_readable();
1069
1070        // Step 2. Let writable be value.[[writable]].
1071        let writable = self.get_writable();
1072
1073        // Step 3. If ! IsReadableStreamLocked(readable) is true, throw a
1074        // "DataCloneError" DOMException.
1075        // Step 4. If ! IsWritableStreamLocked(writable) is true, throw a
1076        // "DataCloneError" DOMException.
1077        if readable.is_locked() || writable.is_locked() {
1078            return Err(Error::DataClone(None));
1079        }
1080
1081        // First port pair (readable → proxy writable)
1082        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        // Second port pair (proxy readable → writable)
1095        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        // Pipe readable into the proxy writable (→ port_1)
1105        readable
1106            .pipe_to(cx, &global, &proxy_writable, false, false, false, None)
1107            .set_promise_is_handled(cx);
1108
1109        // Step 5. Set dataHolder.[[readable]] to !
1110        // StructuredSerializeWithTransfer(readable, « readable »).
1111        // Step 6. Set dataHolder.[[writable]] to !
1112        // StructuredSerializeWithTransfer(writable, « writable »).
1113        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    /// <https://streams.spec.whatwg.org/#ref-for-transfer-receiving-steps②>
1123    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        // Step 1. Let readableRecord be !
1133        // StructuredDeserializeWithTransfer(dataHolder.[[readable]], the
1134        // current Realm).
1135        let proxy_readable = ReadableStream::new_with_proto(cx, owner, None);
1136        proxy_readable.setup_cross_realm_transform_readable(cx, &port2);
1137
1138        // Step 2. Let writableRecord be !
1139        // StructuredDeserializeWithTransfer(dataHolder.[[writable]], the
1140        // current Realm).
1141        let proxy_writable = WritableStream::new_with_proto(cx, owner, None);
1142        proxy_writable.setup_cross_realm_transform_writable(cx, &port1);
1143
1144        // Step 3. Set value.[[readable]] to readableRecord.[[Deserialized]].
1145        // Step 4. Set value.[[writable]] to writableRecord.[[Deserialized]].
1146        // Step 5. Set value.[[backpressure]],
1147        // value.[[backpressureChangePromise]], and value.[[controller]] to
1148        // undefined.
1149        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}